Тема 08

Брокеры сообщений и асинхронность

Секция, где проверяют не знание кнопок в Kafka, а умение рассуждать про распределённую систему с ненадёжной сетью. Почти все вопросы здесь сводятся к одному: «что произойдёт, если сейчас упасть — и кто это заметит». Поэтому каждая глава идёт от того, где именно теряются и дублируются сообщения, а уже потом к настройкам.

Что реально проверяют этой секцией

Три вещи. Первая — понимаешь ли ты, что «доставлено» и «обработано» — это разные события, и что между ними всегда есть окно, в котором можно упасть. Вторая — умеешь ли ты проектировать под дубликаты: кандидат, который на вопрос про exactly-once отвечает «включаем транзакции в Kafka», и кандидат, который отвечает «в общем случае невозможно, поэтому делаем at-least-once плюс идемпотентный обработчик, а транзакции Kafka закрывают только read-process-write внутри самой Kafka» — это два разных грейда. Третья — видишь ли ты границу порядка: где он гарантирован (партиция, очередь с одним консьюмером) и что его ломает (ретраи, параллелизм, DLQ).

8.1Общая теория

Базис, на котором стоят и Kafka, и RabbitMQ, и любой самописный «воркер поверх таблицы». Пока не видно, где именно рвётся, разговор про acks и prefetch сводится к перечню настроек без смысла.

Сначала — словарь: кто здесь кто

Дальше на каждом шагу будут «продюсер», «консьюмер», «ack» и «гарантия доставки». Все четыре понятия простые, но договориться о них надо до разговора про надёжность, иначе остальной текст рассыпается на непонятные слова.

1. Сообщение — это данные плюс конверт

Сообщением называют порцию байтов, которую одна программа передаёт другой не напрямую, а через посредника. У сообщения есть тело (payload: сами данные, обычно JSON или бинарный формат) и заголовки, служебные поля рядом с телом: идентификатор, тип события, время, данные для трассировки. По заголовкам получатель решает, что делать с сообщением, не разбирая тело.

2. Продюсер пишет, консьюмер читает

Продюсер (producer, издатель, отправитель) отправляет сообщение. Консьюмер (consumer, потребитель, подписчик) его получает и обрабатывает. Один и тот же сервис обычно и то и другое сразу: читает из одной очереди, пишет в другую. Русские кальки «производитель / потребитель» звучат криво, поэтому и в разговоре, и на собесе почти всегда говорят по-английски.

3. Брокер — это посредник с диском

Брокер сообщений живёт на отдельном сервере (обычно на кластере серверов): принимает сообщения от продюсеров, кладёт их на диск и отдаёт консьюмерам. Диск тут главное: из-за него продюсер и консьюмер не обязаны работать одновременно. Продюсер отдал сообщение и ушёл, консьюмер придёт через час и заберёт. Kafka, RabbitMQ, NATS, SQS — всё это брокеры; устроены они по-разному, но роль одна.

4. Подтверждение (ack) — момент, когда сообщение считается обработанным

Консьюмер шлёт брокеру ack (от acknowledgement, «подтверждение»), и смысл у него такой: «это сообщение можно мне больше не показывать». Пока ack не пришёл, брокер считает сообщение невыполненной работой и рано или поздно отдаст его снова — этому же консьюмеру или другому. Обратный сигнал, nack (negative ack, «не смог»), говорит «верни в очередь» или «убери в свалку битых сообщений».

Дальше всё упирается в одно: подтверждение и обработка происходят в разные моменты, и между ними всегда есть щель, в которую можно упасть. Если подтвердить раньше, чем сделать, падение потеряет работу. Если позже — после падения работу сделают дважды. Третьего варианта нет, и ниже станет ясно, почему его не бывает в принципе.

Зачем вообще асинхронность

При синхронном вызове мало дождаться ответа: вызываемый сервис обязан быть жив прямо сейчас. Такая связанность называется temporal coupling (связанность во времени), и именно она превращает набор микросервисов в распределённый монолит: падает почта — не создаётся заказ.

Брокер даёт три вещи, и стоит проговаривать их именно так:

  • Развязка (decoupling). Продюсер знает только имя топика или exchange. Он не знает, сколько у события подписчиков, живы ли они и что они с ним сделают. Чтобы добавить четвёртого потребителя, код меняют только у него.
  • Сглаживание пиков (load levelling). Брокер работает буфером. Продюсер может выдавать 20 000 сообщений в секунду в пиковую минуту, а консьюмер спокойно переваривать 3 000/с следующие полчаса. Без буфера пришлось бы держать железо под пик.
  • Надёжность. Сообщение лежит на диске и переживает падение потребителя. Синхронный HTTP-запрос при 500-й ошибке теряется целиком — если только вызывающий сам не построил себе очередь ретраев, то есть тот же брокер, только плохой.
Синхронно — жёсткая связка во времени API заказов Платежи Почта Аналитика 40 мс 120 мс 300 мс 80 мс Клиент ждёт 540 мс. Легла почта — не создался заказ. Доступность = произведение доступностей. Через брокер — развязка во времени API заказов ответ за 40 мс Брокер очередь / лог на диске Платежи Почта Аналитика лежит 10 мин — ок упала — догонит добавлена вчера Платим за это: eventual consistency, дубликаты, потеря глобального порядка, сложная сквозная трассировка и ещё одна stateful-система в проде.
Развязка. Сверху доступность всей операции равна произведению доступностей всех участников. Снизу продюсеру достаточно, чтобы был жив брокер, а потребители догоняют в своём темпе, даже те, кого добавили задним числом.
Половина ответа про цену

Ответ «асинхронность нужна для развязки и надёжности» звучит как заученный. В сильном ответе всегда есть и цена: клиенту теперь нельзя сказать «готово» (нужен статус/поллинг/вебхук), появились дубликаты (нужна идемпотентность), порядок больше не глобальный, отладка стала многошаговой (нужны correlation id и трассировка через заголовки сообщений), и сам брокер стал единой точкой отказа, которую надо кластеризовать и мониторить.

Push vs pull

Вопрос в том, кто начинает передачу: брокер или потребитель. В модели push («толкать») брокер сам шлёт сообщение, как только оно появилось. В pull («тянуть») консьюмер сам приходит и спрашивает: «есть что-нибудь?». Это не деталь реализации: от неё зависит, где именно копится работа, когда потребитель не успевает.

Через backpressure (обратное давление) медленный получатель заставляет быстрого отправителя сбавить темп. Бытовой пример: раздача на кухне ресторана. Если повар ставит тарелки на стойку быстрее, чем официанты их разносят, стойка переполнится, и дальше возможны только два исхода: либо повар видит, что места нет, и притормаживает (обратное давление сработало), либо тарелки летят на пол (потери). В push тормозить отправителя приходится специально, в pull это выходит само собой: не пришёл за сообщением — оно просто осталось лежать у брокера.

PUSH — брокер толкает RabbitMQ basic.consume, core NATS, вебхуки PULL — консьюмер тянет Kafka fetch, SQS long polling, JetStream pull Брокер решает, когда слать Консьюмер обрабатывает шлёт, не спрашивая буфер в памяти растёт + минимальная задержка доставки + брокер сам балансирует между консьюмерами − без prefetch / кредитов консьюмер захлёбывается Обратное давление приходится встраивать вручную Брокер хранит лог Консьюмер сам задаёт темп fetch(offset) батч записей темп задаёт читатель растёт лаг, не память + backpressure получается сам собой + батчинг и перечитывание с любого offset − задержка на опрос: лечится long poll fetch.min.bytes + fetch.max.wait.ms = компромисс
Кто инициирует передачу. В push перегрузку потребителя приходится решать отдельно (prefetch, кредиты). В pull она решается сама: потребитель просто реже приходит, а очередь накапливается в брокере, где для неё есть диск.

Чистого push и чистого pull почти не бывает. RabbitMQ толкает, но с кредитным окном (basic.qos(prefetch)): брокер не пошлёт больше N неподтверждённых сообщений. Kafka тянет, но с long polling: запрос висит на брокере до fetch.max.wait.ms или пока не накопится fetch.min.bytes, поэтому пустых опросов почти нет и задержка остаётся единицами миллисекунд.

Гарантии доставки

Три уровня, и все три называют по-английски, потому что русские варианты не прижились. Читаются они буквально: at-most-once значит «не более одного раза», at-least-once — «не менее одного раза», exactly-once — «ровно один раз». Формально речь не о «доставке», а о том, сколько раз может проявиться эффект обработки: не сколько раз сообщение прилетело по сети, а сколько раз списались деньги или ушло письмо.

Что физически значит «ровно один раз», видно на бытовом примере. Курьер привёз посылку и должен отметить в системе, что доставил. Он позвонил в дверь, отдал коробку — и в этот момент у него сел телефон, отметка не ушла. Диспетчер видит: заявка не закрыта. Что произошло, он не знает: курьер вообще не доехал или доехал, отдал и не отчитался. Вариантов у диспетчера ровно два: послать вторую посылку (получатель получит две) или не посылать (может остаться ни с чем). Способа узнать правду у него нет, и дело не в плохой системе: оборвался единственный канал связи. Между брокером и консьюмером такое случается десятки раз в сутки.

ГарантияМеханикаПотериДублиГде встречается
at-most-onceПодтвердить (ack / commit offset) до обработки, либо вообще не подтверждатьданетметрики, телеметрия, UDP-подобная логика, Redis Pub/Sub
at-least-onceПодтвердить после успешной обработки, при сбое — переотправканетдадефолт в Kafka, RabbitMQ, SQS, JetStream — 95 % прода
exactly-onceНедостижимо в общем случае. На практике — at-least-once плюс дедупликациянетнет эффектаKafka EOS внутри Kafka; везде ещё — «effectively once»
at-most-once ack до обработки брокер консьюмер 1. доставка 2. ack сразу обработка процесс упал в середине брокер уже всё забыл потерь: да дубликатов: нет at-least-once ack после обработки брокер консьюмер 1. доставка обработка 2. ack после упал между этими шагами переотправка = дубликат потерь: нет дубликатов: да exactly-once проблема двух генералов брокер консьюмер 1. доставка обработка 2. ack потерян сетью брокер не знает: было или нет повторить = дубль, не повторить = потеря решения нет ни при каком числе сообщений
Где именно рвётся. Гарантии отличаются только порядком двух действий: подтверждения и обработки. Третья панель показывает, почему «правильного» порядка не придумать: подтверждение само может потеряться, а ещё одно подтверждение лишь сдвигает проблему на шаг.

Почему exactly-once невозможен: проблема двух генералов

Классическая формулировка: две армии на холмах по обе стороны от долины, в долине противник. Атака удастся только при одновременном ударе. Гонцы ходят через долину, и их могут перехватить. Доказано, что никакого протокола с конечным числом сообщений, гарантирующего согласованное решение, не существует. Доказательство от противного: возьмём кратчайший корректный протокол; его последнее сообщение может потеряться, но раз протокол корректен и без него, это сообщение лишнее, и мы получили протокол короче. Повторяем, пока сообщений не останется вовсе, и приходим к противоречию.

Прикладной перевод: отправитель никогда не может достоверно узнать, был ли получен и обработан его последний запрос. Не получив ответа, он вынужден выбирать между «повторить» (риск дубликата) и «не повторять» (риск потери). Брокер этот выбор не отменяет, а лишь берёт на себя. Сюда же относится FLP-теорема (Fischer, Lynch, Paterson, 1985): в асинхронной системе, где может отказать даже один узел, детерминированный консенсус за конечное время невозможен.

Формулировка, которую стоит выучить

«Exactly-once доставка невозможна. Достижима exactly-once обработка, и только как сочетание at-least-once доставки с идемпотентным потребителем или транзакцией, которая атомарно фиксирует и результат работы, и факт потребления. То есть проблема решается не на транспорте, а end-to-end, на прикладном уровне.» Именно это и хотят услышать.

Что тогда такое «exactly-once» у Kafka

Две отдельные вещи, которые маркетинг слепил в одну. Во-первых, идемпотентный продюсер (enable.idempotence=true, включён по умолчанию с Kafka 3.0): продюсер получает PID и нумерует записи sequence number по партиции, брокер отбрасывает повтор с уже виденным номером. Это убирает дубли от ретраев самого продюсера, и только в пределах одной сессии продюсера. Во-вторых, транзакции: атомарная запись в несколько партиций плюс атомарный коммит offset входного топика через sendOffsetsToTransaction. Вместе они дают exactly-once для сценария read-process-write внутри Kafka. Как только обработчик делает HTTP-вызов, шлёт письмо или пишет в стороннюю БД без общей транзакции, — гарантия заканчивается: откатить письмо Kafka не умеет.

Идемпотентность консьюмера

Операция идемпотентна, если даёт один и тот же результат, сколько бы раз её ни повторили: один раз, три раза, сто. Так устроена кнопка вызова лифта. Нажал её один раз или продолбил пятнадцать — лифт приедет один раз, состояние «лифт вызван» от повторов не меняется. А вот кнопка «взять талончик» в очереди неидемпотентна: каждое нажатие выдаёт новый номер. Слово латинское: idem — «тот же», potens — «имеющий силу»; формально f(f(x)) = f(x).

Раз дубликаты неизбежны, обработчик устраивают так, чтобы повторная обработка того же сообщения не меняла результат. Путей два, и они принципиально разные.

Путь первый: сделать саму операцию идемпотентной по природе. Он лучше, потому что ничего не надо хранить. UPDATE users SET status='active' WHERE id=$1 идемпотентен, а UPDATE accounts SET balance = balance - 100 — нет. Приёмы: UPSERT по естественному ключу, установка абсолютного значения вместо инкремента, оптимистическая блокировка по версии (WHERE version < $2), условная вставка ON CONFLICT DO NOTHING.

Путь второй: явная дедупликация по ключу. Дедупликация означает «вести список уже сделанного и сверяться с ним»: у каждого сообщения есть ключ, обработчик перед работой смотрит, нет ли этого ключа в списке, и если есть, молча пропускает сообщение. Как журнал на вахте: фамилия уже записана — второй раз не пускаем. Этот путь нужен, когда операция принципиально неидемпотентна (списание денег, отправка письма, вызов внешнего API). Ключ дедупликации выбирают в таком порядке:

  • Бизнес-ключ: payment_id, order_id + event_type, aggregate_id + version. Самый надёжный вариант: такой ключ переживает и переотправку продюсером, и повторную генерацию события.
  • Message ID от продюсера: UUID, который кладут в заголовок при первой отправке. Генерируют его один раз, а не на каждом ретрае.
  • Координаты в брокере (topic:partition:offset) работают только внутри одной партиции и ломаются при перекладывании через retry-топик. Это крайний вариант.
-- Таблица дедупликации. Хватает одного PRIMARY KEY.
CREATE TABLE processed_messages (
    consumer      text        NOT NULL,   -- имя группы: один и тот же msg могут есть разные консьюмеры
    message_key   text        NOT NULL,   -- бизнес-ключ или message_id
    processed_at  timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer, message_key)
);

-- Чистка: без неё таблица растёт вечно. Старое удаляют пачками по этому индексу:
-- секционировать по processed_at нельзя, ключ таблицы обязан её включать.
CREATE INDEX ON processed_messages (processed_at);
// Дедупликация и бизнес-логика обязаны быть в одной транзакции.
// Иначе получаем ту же проблему двойной записи, только в миниатюре.
func (h *Handler) Handle(ctx context.Context, m Message) error {
    tx, err := h.db.BeginTx(ctx, nil)
    if err != nil {
        return err
    }
    defer tx.Rollback() // no-op после успешного Commit

    // ON CONFLICT DO NOTHING: если строка уже есть, RowsAffected == 0
    res, err := tx.ExecContext(ctx,
        `INSERT INTO processed_messages (consumer, message_key)
         VALUES ($1, $2) ON CONFLICT DO NOTHING`,
        h.group, m.Key)
    if err != nil {
        return err
    }
    if n, _ := res.RowsAffected(); n == 0 {
        // Уже обрабатывали. Коммитить нечего, просто подтверждаем брокеру.
        return nil
    }

    if err := h.apply(ctx, tx, m); err != nil { // вся бизнес-логика через tx
        return err
    }
    return tx.Commit()
}
Три способа сломать себе дедупликацию
  • Ключ генерируется на каждой попытке. Продюсер делает uuid.New() внутри цикла ретраев — дубликаты получают разные ключи, и дедупликация не срабатывает. Ключ должен рождаться один раз, вместе с событием.
  • Дедуп в Redis, бизнес-данные в Postgres. Между SETNX и COMMIT можно упасть: ключ отметили, работу не сделали, повтор пропустят. Redis годится как быстрый фильтр перед транзакцией, но источником истины должна быть та же БД, что и данные.
  • Дедуп-таблица без ретеншена. Через полгода это самая большая таблица в базе, и INSERT в неё обходится дороже самой бизнес-операции.

Персистентность

Персистентность (от persistent, «стойкий, сохраняющийся») гарантирует, что сообщение переживёт перезапуск брокера. Для этого сообщение мало принять: его записывают туда, откуда оно достанется после включения. Обратный режим называется transient: сообщение живёт только в оперативной памяти брокера и исчезает вместе с процессом.

Тонкость, на которой ловят: «записано на диск» не равно «не потеряется». Между write() и физическим носителем стоит page cache, область оперативной памяти, где ядро ОС держит данные, ещё не сброшенные на диск. Вызов write() возвращает успех сразу, как только байты легли в этот кэш, а на пластину или флеш они попадут когда-нибудь потом. Заставить сбросить их прямо сейчас умеет только fsync; без него отключение питания стирает всё, что ещё лежало в кэше.

  • Kafka сознательно не делает fsync на каждую запись: пишет в page cache, а сбрасывает по log.flush.interval.messages / .ms (по умолчанию практически «никогда, отдано ОС»). Надёжность держится не на диске одного узла, а на репликации: acks=all + min.insync.replicas=2 при replication.factor=3. Данные потеряются, только если одновременно упадут несколько брокеров.
  • RabbitMQ идёт от сообщения: очередь должна быть durable, сообщение помечено delivery_mode=2 (persistent), и обязательно нужны publisher confirms, иначе basic.publish работает как fire-and-forget, и брокер вправе молча потерять сообщение. Quorum queues добавляют репликацию через Raft.

За персистентность платят задержкой и пропускной способностью. Поэтому в проде встречается смешанный режим: критичные события (платежи) идут persistent с подтверждениями, а метрики и телеметрия — transient, потому что терять их не страшно, а платить за них полной надёжностью дорого.

Ordering: где порядок есть, а где его нет

Под ordering понимают порядок сообщений: гарантию, что консьюмер увидит события в той же последовательности, в какой их отправлял продюсер. Звучит как мелочь, но без неё «заказ отменён» иногда приходит раньше «заказ создан», и обработчик не понимает, что ему отменять.

Общее правило: порядок гарантируется только внутри одной единицы упорядочивания. Так называют минимальный кусок системы, где сообщения гарантированно выстроены в одну линию. Таких единиц в системе много, и между собой они никак не синхронизированы. В Kafka эту роль играет партиция (один физический файл-лог внутри топика, подробно в 8.2), в RabbitMQ — очередь при единственном консьюмере и без requeue, в JetStream — стрим (или subject-фильтр внутри него). Глобального порядка «по топику» нет ни у кого, иначе не было бы горизонтального масштабирования.

Что ломает порядокПочемуЧто делать
Несколько партиций / очередейМежду ними нет общих часовКлюч партиционирования = идентификатор сущности
Несколько консьюмеров на одной очередиОбрабатывают параллельно, финишируют вразнобойОдин консьюмер, либо шардировать по ключу
Параллелизм внутри консьюмераГорутина на сообщение = гонкаПул воркеров, где воркер выбирается как hash(key) % N
Ретраи с requeue / retry-топикиСообщение уезжает в конец, следующие обгоняютОстанавливать обработку ключа, либо принять расхождение
DLQ и ручной перезапускСообщение вернётся часы спустяВерсионирование: отбрасывать событие, если версия устарела
Продюсер с max.in.flight.requests.per.connection > 1 без идемпотентностиПовтор первого батча встанет после второгоenable.idempotence=true — тогда безопасно до 5
Практическое правило

Глобальный порядок почти никогда не нужен, нужен порядок на сущность: события одного заказа не должны перемешаться, а события разных заказов друг другу безразличны. Отсюда рецепт: ключ сообщения = order_id, и всё остальное решается само. Если же порядок нужен по-настоящему сквозной, самый честный ответ на собесе звучит так: «делаем одну партицию и теряем масштабирование, либо кладём в событие версию/timestamp и делаем обработчик устойчивым к перестановке (last-write-wins по версии)».

Split brain у кластеров брокеров

Сеть делится на две части (network partition). Обе половины видят, что «остальные умерли», обе выбирают себе лидера и обе начинают принимать записи. Когда сеть починят, останутся две расходящиеся истории, и корректно слить их нельзя: сообщения либо теряются, либо дублируются, либо оказываются в разном порядке у разных потребителей.

Защищаются все одинаково: через кворум большинства. Кворумом называют минимальное число узлов, которые должны согласиться, чтобы решение считалось принятым; «большинство» значит N/2+1. Смысл в том, что два любых большинства из одного набора узлов обязательно пересекаются хотя бы в одном узле, а он не проголосует за два разных решения сразу. Поэтому две половины разорванной сети не могут обе набрать кворум, и второй лидер не появляется. Реализовано это у всех по-разному:

  • Kafka. Раньше контроллер и метаданные жили в ZooKeeper (кворум ZAB), с KIP-500 переехали в саму Kafka, в KRaft (Raft-кворум контроллеров). Production-ready с 3.3, миграция с ZooKeeper готова для продакшена с 3.6, а в Kafka 4.0 (2025) ZooKeeper удалён совсем. Данные страхует unclean.leader.election.enable=false (так по умолчанию): если ни одна реплика из ISR (in-sync replicas, копии партиции, которые не отстали от лидера; подробно в 8.2) не доступна, партиция становится недоступной, но не расходящейся. Это осознанный выбор CP в терминах CAP.
  • RabbitMQ. Классические зеркалируемые очереди (mirrored) split brain переживали плохо: при автоматическом слиянии одна половина просто откатывалась, теряя сообщения. Поэтому они объявлены устаревшими в 3.9 и удалены в 4.0, а штатным вариантом стали quorum queues на Raft: пишет только большинство, меньшинство отказывает в записи. Для кластера настраивают cluster_partition_handling = pause_minority или autoheal.
  • NATS JetStream ведёт отдельную Raft-группу на каждый стрим, R=3 или R=5.
Подвох: «настроил acks=all — значит, не потеряю»

Не значит. acks=all означает «все реплики из текущего ISR подтвердили». Если ISR схлопнулся до одной реплики (остальные отстали или упали), acks=all превращается в acks=1, и падение этого брокера означает потерю. Ровно для этого есть min.insync.replicas=2: при меньшем ISR продюсер получит NotEnoughReplicasException и запись не пройдёт. Правильная пара: RF=3 и min.insync.replicas=2. С ней кластер переживает падение одного брокера без потерь и без остановки записи.

Очередь vs топик

Базовых моделей обмена две, и различаются они тем, сколько потребителей увидит одно сообщение.

Очередь — point-to-point Топик — publish / subscribe Producer очередь tasks m4 m3 m2 m1 C1 C2 C3 одно сообщение получает ровно один потребитель масштабирование = добавить консьюмера порядок теряется, как только консьюмеров больше одного Producer лог / fanout topic: orders S1 S2 S3 копия события уходит каждому подписчику Kafka: одна группа = очередь, разные группы = pub/sub RabbitMQ: fanout exchange + своя очередь на подписчика
Две модели. Очередь распределяет работу, топик распространяет факт. Kafka реализует обе поверх одного лога: внутри consumer group партиции делятся между участниками (очередь), а разные группы читают всё независимо (pub/sub).
Дело не в терминах

От выбора зависит, кто знает о потребителях. Через очередь идёт команда («сделай задачу»), и отправитель обычно знает, что её выполнят. В топик уходит факт («заказ создан»), и отправителю всё равно, кто на него отреагирует. На собесе по архитектуре из этого вырастает вопрос про commands vs events: команды маршрутизируют в конкретную очередь и валидируют, события публикуют в топик и не валидируют, потому что событие уже случилось и отменять нечего.

Вопросы

8
Суть: брокер убирает связанность во времени. Продюсеру больше не нужно, чтобы потребитель был жив прямо сейчас, — и из этого автоматически следуют буферизация пиков и переживание падений.

У синхронной цепочки HTTP-вызовов два свойства, которые под нагрузкой становятся смертельными. Доступность перемножается: четыре сервиса по 99,9 % дают 99,6 %, то есть примерно три часа простоя в месяц вместо сорока минут. А задержка суммируется, и клиент ждёт худшую ветку, даже если ему интересна только первая.

Три эффекта, которые даёт брокер

  • Развязка (decoupling). Продюсер публикует факт в топик и не знает ни числа подписчиков, ни их состояния. Появился пятый потребитель события «заказ создан» — команда заказов не выкатывает ни строчки кода. Это разница между «изменение в N местах» и «изменение в одном».
  • Сглаживание пиков (load levelling). Брокер держит буфер на диске. В чёрную пятницу продюсер выдаёт 30 000 сообщений/с, а консьюмеры продолжают жевать по 4 000/с; растёт лаг, но ничего не падает. Без буфера пришлось бы держать железо под пик, который бывает два раза в год.
  • Надёжность. Сообщение лежит в персистентном логе. Потребитель может лежать двадцать минут, подняться и доесть накопившееся. Синхронный запрос при недоступности адресата просто теряется, если только вызывающий сам не построил себе очередь ретраев, то есть плохой брокер.

Чем добить ответ

Обязательно назвать цену, иначе ответ звучит как реклама. За асинхронность платят конечной согласованностью вместо мгновенной (клиент нажал кнопку, а эффект наступит через секунду или через минуту), дубликатами, которые надо обрабатывать, потерей сквозного порядка, отсутствием простого стектрейса (нужен трейсинг с пробросом traceparent в заголовках сообщения) и ещё одной инфраструктурной системой, которую надо мониторить и обновлять. Хорошая формулировка: «асинхронность — это не оптимизация, а смена контракта: вместо ответа мы отдаём обещание».

Когда брокер не нужен

Когда вызывающему нужен результат прямо сейчас, чтобы показать его пользователю, это работа для синхронного вызова, а «запрос-ответ через две очереди» превращается в RPC с ручной корреляцией, таймаутами и мусорными очередями. Асинхронность оправдана там, где ответ не нужен немедленно: побочные эффекты (письма, индексация, аналитика), долгие задачи, интеграция с внешними системами и веерная рассылка событий.

Суть: вопрос в том, кто задаёт темп. В push темп задаёт брокер, и обратное давление приходится встраивать руками (prefetch, кредиты). В pull темп задаёт потребитель, и backpressure получается бесплатно — растёт лаг, а не память консьюмера.

Push

Брокер сам отправляет сообщение, как только оно появилось: RabbitMQ basic.consume, core NATS, вебхуки. Задержка минимальная, брокер видит всех потребителей и может балансировать между ними. Но если потребитель медленнее потока, у него растёт очередь в памяти, пока процесс не убьёт OOM. Лечится кредитным окном: basic.qos(prefetch=N) означает «не шли больше N неподтверждённых». По сути так вручную делают то, что pull даёт даром.

Pull

Потребитель сам запрашивает следующую порцию: Kafka fetch, SQS long polling, JetStream pull consumers. Плюсы: естественное обратное давление, крупные батчи (эффективнее по CPU и сети), возможность перечитать с любого offset, простая семантика при добавлении потребителей. Расплачиваться приходится задержкой на опрос и пустыми запросами; их убирает long polling: Kafka держит запрос на брокере до fetch.max.wait.ms (500 мс по умолчанию) или пока не наберётся fetch.min.bytes, поэтому при живом трафике данные приходят практически мгновенно.

Чистых моделей нет

RabbitMQ толкает с кредитами и под нагрузкой ведёт себя «pull-подобно». Kafka тянет с long poll и по задержке похожа на push. JetStream умеет обе модели явно: push-консьюмер с MaxAckPending (те же кредиты) и pull-консьюмер с Fetch(batch). Правильный вывод для собеса: модель доставки определяет, где копится очередь при перегрузке — в памяти потребителя или на диске брокера. Второе почти всегда лучше.

Суть: разница между гарантиями — это порядок двух действий, подтверждения и обработки. Exactly-once доставка недостижима (проблема двух генералов); достижима exactly-once обработка — только end-to-end, через at-least-once плюс идемпотентность или общую транзакцию.

Три уровня

  • При at-most-once подтверждаем до обработки. Упали в середине — сообщение потеряно навсегда. Осмысленно для метрик и телеметрии, где потеря дешевле дубликата.
  • С at-least-once подтверждение идёт после успешной обработки. Если упасть между обработкой и подтверждением, брокер переотправит сообщение, и будет дубль. Это дефолт Kafka, RabbitMQ, SQS, JetStream и 95 % прода.
  • exactly-once обещает ровно один эффект, но как свойство транспорта в общем случае недостижим.

Почему невозможно

Проблема двух генералов. Две армии по краям долины, атака удастся только одновременная, гонцов перехватывают. Доказано: протокола с конечным числом сообщений, гарантирующего согласованное решение, не существует. Доказательство от противного: возьмём кратчайший корректный протокол; его последнее сообщение может потеряться, но раз протокол корректен и без него, оно лишнее, значит, есть протокол короче; повторяем до нуля сообщений и получаем противоречие. Прикладной перевод: отправитель принципиально не может узнать, был ли обработан его последний запрос, и вынужден выбирать между риском дубликата и риском потери. Рядом стоит FLP (1985): в асинхронной системе с возможным отказом одного узла детерминированный консенсус за конечное время невозможен.

Что делать на практике

Признать at-least-once и сделать обработчик идемпотентным, либо атомарно фиксировать результат работы и факт потребления в одной транзакции. Это и есть «effectively once». Kafka EOS относится ко второму варианту: транзакция покрывает запись в выходные партиции и коммит offset входных, то есть работает только для read-process-write внутри Kafka. Как только обработчик отправляет письмо или дёргает внешний API, гарантия испаряется.

На что ловят

Кандидат говорит «в Kafka есть exactly-once, включаем и всё». Follow-up: «а если консьюмер пишет результат в Postgres — что даст транзакция Kafka?» Правильный ответ: ничего, потому что коммит offset в Kafka и коммит в Postgres идут двумя разными транзакциями без общего координатора; нужен либо идемпотентный апсерт в Postgres, либо inbox-таблица, либо outbox на стороне записи.

Суть: идемпотентность — свойство «повторная обработка не меняет результат». Достигается либо природой операции (UPSERT, установка абсолютного значения), либо явной дедупликацией по ключу в той же транзакции, что и бизнес-логика.

Путь первый: операция идемпотентна сама по себе

Он лучше, потому что не требует хранить состояние. UPDATE users SET status='active' WHERE id=$1 идемпотентен, UPDATE accounts SET balance = balance - 100 — нет. Приёмы: INSERT ... ON CONFLICT DO NOTHING или DO UPDATE по естественному ключу, запись абсолютного значения вместо дельты, оптимистическая блокировка по версии (WHERE version < $2), «перевод в состояние» вместо «переключения состояния».

Путь второй: явная дедупликация

Нужен, когда операция неидемпотентна по природе: списание, письмо, вызов внешнего API. Ключ выбираем в таком порядке. Лучше всех бизнес-ключ (payment_id, order_id + event_type, aggregate_id + version): переживает и переотправку продюсером, и повторную генерацию события. Дальше message_id, созданный один раз вместе с событием и положенный в заголовок. Последними идут координаты в брокере (topic:partition:offset): работают только внутри одной партиции и ломаются при перекладывании через retry-топик.

CREATE TABLE processed_messages (
    consumer      text        NOT NULL,   -- имя группы: одно сообщение едят несколько консьюмеров
    message_key   text        NOT NULL,
    processed_at  timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer, message_key)
);
-- PARTITION BY RANGE (processed_at) здесь не выйдет: первичный ключ секционированной
-- таблицы обязан включать processed_at. Старое чистят DELETE пачками по индексу.
CREATE INDEX ON processed_messages (processed_at);

Вставку в дедуп-таблицу и бизнес-изменения делают одной транзакцией. Если их разнести, между ними можно упасть, и повтор отбросят, хотя работа не сделана. Дедуп в Redis допустим только как быстрый фильтр перед транзакцией, но источником истины он быть не может.

Окно дедупликации

Хранить ключи вечно нельзя. Разумный горизонт равен максимальному времени жизни сообщения в системе с запасом: retention топика плюс максимальная задержка retry-цепочки плюс время, за которое реально разбирают DLQ. Обычно от суток до недели. Об этом почти никто не говорит, поэтому, если назвать окно ретеншена, ответ сразу выделится.

Суть: персистентность — гарантия, что сообщение переживёт перезапуск брокера. Тонкость: «записано на диск» не равно «не потеряется» — между write() и пластиной стоит page cache ядра.

Без персистентности брокер превращается в буфер в оперативной памяти: перезапуск пода, OOM или падение узла стирают всё, что не успели прочитать. Персистентность нужна везде, где потеря сообщения означает потерю денег или расхождение данных между сервисами: платежи, изменения статусов, любые события, по которым другой сервис строит своё состояние.

Две разные реализации

  • Kafka. Сознательно не делает fsync на каждую запись: пишет в page cache, сброс отдан ОС (log.flush.interval.messages по умолчанию практически «никогда»). Надёжность даёт репликация, а не диск одного узла: acks=all при replication.factor=3 и min.insync.replicas=2. Запись пропадёт, только если питание отключится сразу у нескольких брокеров. Это осознанный размен: fsync на запись стоил бы порядка производительности.
  • RabbitMQ. Идёт от сообщения: очередь объявлена durable, сообщение помечено delivery_mode=2 (persistent), и обязательно включены publisher confirms. Без confirms basic.publish работает как fire-and-forget: TCP-запись прошла, а брокер мог не принять, и ты об этом не узнаешь. Quorum queues добавляют репликацию через Raft и пишут на диск всегда.
Классический подвох

«Очередь durable — значит, сообщения не потеряются?» Нет. durable у классической очереди означает, что переживёт перезапуск сама очередь как объект, а не её содержимое. Нужны обе галочки: durable queue плюс persistent message. И даже этого мало без publisher confirms — потерять можно на участке «продюсер отправил, брокер ещё не записал».

Суть: порядок есть только внутри одной единицы упорядочивания: партиция в Kafka, очередь с единственным консьюмером в RabbitMQ, стрим в JetStream. Глобального порядка по топику не бывает ни у кого — иначе не было бы масштабирования.

Что ломает порядок

  • Несколько партиций. Между ними нет общих часов и нет координации. Лечится ключом партиционирования: hash(order_id) % partitions, тогда все события одного заказа лягут в одну партицию и порядок сохранится.
  • Несколько консьюмеров на одной очереди. Сообщения раздаются по кругу и финишируют вразнобой: второе может обработаться раньше первого просто потому, что оно легче.
  • Параллелизм внутри консьюмера. «Горутина на сообщение» уничтожает порядок даже при одной партиции. Если нужны и параллелизм, и порядок, берут пул воркеров с выбором hash(key) % N: один ключ всегда идёт в одного воркера.
  • Ретраи. Сообщение уходит в конец очереди или в retry-топик, а следующие его обгоняют. Это самый недооценённый источник расхождений.
  • Продюсер с max.in.flight.requests.per.connection > 1 без идемпотентности. Ретрай первого батча приходит после второго. С enable.idempotence=true брокер по sequence number отклоняет батч с разрывом в номерах (OUT_OF_ORDER_SEQUENCE_NUMBER), продюсер отправляет заново по порядку, и до 5 in-flight безопасно.

Как отвечать

Сначала спросить, какой именно порядок нужен. Почти всегда нужен порядок на сущность, а не глобальный: события одного заказа не должны перемешаться, события разных заказов друг другу безразличны. Отсюда рецепт в одну строку: ключ сообщения равен идентификатору сущности. Если же порядок нужен по-настоящему сквозной, честный ответ: «одна партиция и один консьюмер, масштабирование теряем», либо «кладём в событие версию и делаем обработчик устойчивым к перестановке (last-write-wins по версии, устаревшее событие отбрасываем)».

Суть: сеть распалась, обе половины считают себя живыми, обе выбрали лидера, обе принимают записи — и слить эти истории потом нельзя. Защита у всех одна: кворум большинства, писать разрешено только той половине, где больше половины узлов.

Механика: network partition разрезает кластер на группы, которые не видят друг друга. Молчащий узел не отличить от мёртвого, поэтому каждая группа решает, что остальные умерли, и продолжает работу. Когда связность восстанавливается, в двух половинах лежат разные сообщения под одними и теми же offset. Корректного автоматического слияния не существует: любое решение либо теряет данные, либо создаёт дубликаты, либо ломает порядок.

Как защищаются конкретные брокеры

  • Kafka. Метаданные и выборы лидера идут через кворум: раньше ZooKeeper (ZAB), с KIP-500 — KRaft, Raft-кворум контроллеров внутри самой Kafka. Production-ready с 3.3, миграция с 3.6, в Kafka 4.0 ZooKeeper удалён полностью. Для данных стоит предохранитель unclean.leader.election.enable=false (дефолт): если ни одной реплики из ISR нет, партиция становится недоступной, но не расходящейся. Это выбор CP в терминах CAP.
  • RabbitMQ. Классические зеркалируемые очереди split brain переживали плохо: при автослиянии одна половина откатывалась и теряла сообщения. Они deprecated в 3.9 и удалены в 4.0; теперь штатно используют quorum queues на Raft. Плюс настройка кластера cluster_partition_handling: pause_minority (меньшинство само останавливается), autoheal или ignore — последний режим и есть разрешённый split brain.
  • NATS JetStream — Raft на уровне стрима, R=3 или R=5; при потере кворума стрим недоступен целиком — ни записи, ни чтения.
Чем добить

Упомянуть, что кворум требует нечётного числа узлов: кластер из двух брокеров не переживает потерю одного (нет большинства), а из четырёх переживает столько же отказов, сколько из трёх, но стоит дороже. И что acks=all не спасает сам по себе: он означает «все реплики из текущего ISR», а если ISR схлопнулся до одной, это фактически acks=1. Поэтому пара RF=3 + min.insync.replicas=2 считается стандартным прод-минимумом.

Суть: в очереди сообщение получает ровно один потребитель из конкурирующих — это распределение работы. В топике копию получает каждый подписчик — это распространение факта.

Очередь (point-to-point)

Несколько консьюмеров подписаны на одну очередь, брокер отдаёт каждое сообщение одному из них. Это паттерн competing consumers: масштабирование = добавить процесс, брокер сам перераспределит нагрузку. Сообщение после подтверждения удаляется. Порядок сохраняется, только пока консьюмер один. Так обычно раздают задачи: сгенерировать PDF, отправить письмо, пересчитать отчёт.

Топик (publish/subscribe)

Продюсер публикует событие, каждый подписчик получает свою копию и обрабатывает независимо. Продюсер не знает о подписчиках вообще. Сообщение не «съедается», а живёт по retention-политике или пока его не заберут все. Сюда ложатся события домена: «заказ создан», на который реагируют склад, биллинг, уведомления и аналитика.

Как это выглядит в реальных брокерах

В RabbitMQ это разные объекты: очередь — это очередь, а pub/sub делается через fanout-exchange и по отдельной очереди на каждого подписчика. В Kafka обе модели выражены через один лог и понятие consumer group: внутри группы партиции поделены между участниками (получается очередь), разные группы читают весь топик независимо (получается pub/sub). Поэтому в Kafka, чтобы добавить нового потребителя, достаточно завести новый group.id, и старые потребители об этом даже не узнают.

Команды и события

Из этой пары вырастает архитектурный вопрос. Команда («сними деньги») адресована конкретному обработчику, у неё один исполнитель, её можно отклонить, и живёт она в очереди. Событие («деньги сняты») сообщает о факте в прошедшем времени, у него ноль или много подписчиков, отменять нечего, и место ему в топике. Путаница между ними («отправляем событие SendEmail») обычно означает, что связанность на самом деле никуда не делась.

8.2Kafka

Kafka устроена не как очередь, а как распределённый упорядоченный лог с указателями чтения. Из этой фразы выводятся почти все её особенности (и почти все грабли): сообщение не удаляется после прочтения, порядок живёт внутри партиции, а «кто что прочитал» Kafka записывает отдельно, в служебный топик.

Сначала — что такое лог и чем он не очередь

Большинство недоразумений вокруг Kafka растёт из одного места: человек читает «Kafka — брокер сообщений» и достраивает привычную картинку «положил, забрал, исчезло». В Kafka сообщение никуда не исчезает. Разберём это на трёх шагах, без них дальше будет каша.

1. Лог — это файл, в который только дописывают

Слово лог здесь означает не «логи приложения», а структуру данных: последовательность записей, к которой можно только дописать в конец. Изменить или удалить запись в середине нельзя — такой операции просто нет. По-английски это называется append-only log, «журнал только на дозапись».

Так устроен бумажный кассовый журнал. Каждая операция пишется новой строкой, строки нумеруются подряд, вымарывать написанное запрещено. Ошибся — пишешь новую строку «отмена операции №17». Такой журнал можно перечитать с начала и восстановить всю историю: из него ничего не пропало.

2. Читатель — это закладка, а не потребитель

В очереди чтение сообщение расходует: забрал — его больше нет. В логе чтение не меняет ничего. Каждый читатель держит собственный номер строки, на которой он остановился; этот номер и называется offset («смещение»). Прочитать запись означает взять строку под своим номером и сдвинуть номер на единицу вперёд.

Отсюда сразу три свойства, которых у очереди нет. Читателей может быть сколько угодно, и друг другу они не мешают: у каждого своя закладка. Историю можно перечитать: отмотать закладку назад и обработать вчерашние события заново, например после исправления бага. И отставание выражается числом: конец лога минус твоя закладка и есть всё, что ты ещё не разгрёб.

3. Чем это отличается от очереди — практически

Очередь (RabbitMQ)Лог (Kafka)
Что делает чтениеЗабирает сообщение из очередиДвигает закладку; лог не меняется
Когда данные удаляютсяПосле подтверждения консьюмеромПо сроку или размеру — независимо от того, прочитал их кто-то или нет
Кто помнит прогрессБрокер, по каждому сообщению отдельноБрокер помнит одно число на группу читателей
Перечитать историюНельзя: сообщений уже нетМожно, пока не истёк срок хранения
Что значит «сообщение потеряно»Брокер его выбросил или отдал не томуКонсьюмер не успел прочитать до истечения срока хранения

В эксплуатации это значит, что Kafka не знает, «сколько сообщений в очереди». Она знает, где конец лога и где чья закладка. Разница между ними называется lag, и это единственная метрика, которая тут вообще имеет смысл.

Архитектура: из чего это собрано

Кластер состоит из нескольких брокеров (процессов JVM), у каждого свой набор дисков. Топик остаётся логическим именем, физического воплощения у него нет. Хранятся данные в партициях: партиция занимает каталог на диске конкретного брокера, внутри лежит append-only лог. Топик с 12 партициями даёт 12 независимых логов, размазанных по кластеру.

Партиция нарезана на сегменты, каждый из трёх файлов:

  • 00000000000000012345.log хранит сами записи, батчами, ровно в том виде, в каком их прислал продюсер;
  • .index — разреженный индекс «относительный offset → байтовое смещение в .log» (запись примерно раз в index.interval.bytes = 4 КБ);
  • .timeindex переводит «timestamp → offset» и нужен для retention.ms и для offsetsForTimes() при перемотке по времени.

Число в имени файла, base offset, совпадает с offset первой записи сегмента. Пишут всегда только в активный сегмент; он закрывается и создаётся новый по segment.bytes (1 ГБ по умолчанию) или segment.ms (7 суток). Удаление и компакция работают сегментами, а не записями, и активный сегмент не трогают никогда. При низком трафике «протухшее» сообщение может пролежать сильно дольше retention.ms.

Offset служит номером записи внутри партиции: номера монотонно растут, назначает их лидер. Он не глобальный: offset 100 в P0 и offset 100 в P1 не связаны ничем. Запись по offset брокер ищет так: бинарным поиском по именам сегментов, затем по .index, затем линейным досканированием внутри 4 КБ. Выходит, что чтение с произвольного offset стоит примерно столько же, сколько чтение с конца.

У каждой партиции есть лидер и фолловеры (replication.factor). Весь трафик записи и (по умолчанию) чтения идёт через лидера; фолловеры непрерывно тянут данные тем же fetch-протоколом, что и обычные консьюмеры. В ISR (in-sync replicas) входят реплики, которые не отстали больше чем на replica.lag.time.max.ms (30 с). High watermark отмечает максимальный offset, реплицированный на все реплики из ISR; консьюмеры видят записи только до HW, иначе после смены лидера они могли бы прочитать то, чего в новом лидере нет.

Координацию (кто лидер, где партиции, кто в ISR) раньше вёл ZooKeeper. После KIP-500 метаданные живут в самой Kafka, в режиме KRaft: выделенные controller-узлы держат Raft-кворум и ведут собственный лог метаданных. Production-ready с 3.3, миграция с 3.6, в Kafka 4.0 (2025) ZooKeeper удалён полностью. На практике failover контроллера стал на порядок быстрее, а кластер держит миллионы партиций вместо пары сотен тысяч.

topic: orders · partitions = 3 · ключ сообщения = order_id Буква в ячейке — ключ записи. Одинаковый ключ всегда ложится в одну партицию, поэтому порядок событий одного заказа сохраняется. offset 0 1 2 3 4 5 6 7 P0 лидер b1 A D A A D D A A committed offset = 5 log end offset = 8 P1 лидер b2 B E B E B E B E committed = 3 P2 лидер b3 C F C C F C F C committed = 6 lag = 2 consumer group: billing C1 C2 владеет P0, P1 владеет P2 Внутри группы партиция принадлежит ровно одному консьюмеру. Три партиции на двух консьюмеров = 2 + 1: идеального баланса не бывает. Четвёртый консьюмер в этой же группе будет простаивать — параллелизм внутри группы ограничен числом партиций. Группа analytics читает тот же лог со своими offset и про billing ничего не знает: это pub/sub поверх одного лога.
Топик, партиции, offset, группа. После чтения записи не удаляются, двигается только указатель. Committed offset хранится в служебном топике __consumer_offsets, а лаг считается как разница между концом лога и этим указателем.

Почему Kafka быстрая

Вопрос любят задавать, потому что ответ проверяет понимание работы с диском и сетью, а не знание Kafka. Четыре причины, по убыванию значимости.

  1. Последовательная запись append-only. Нет обновлений на месте, нет B-дерева, нет случайных seek. Последовательная запись на обычный HDD даёт сотни МБ/с, что сопоставимо со случайным доступом к памяти. Kafka нарочно спроектирована так, чтобы обращаться к диску только последовательно.
  2. Page cache вместо своего кэша. Kafka не держит данные в JVM heap, а пишет и читает через страничный кэш ядра. GC от этого не страдает, кэш переживает перезапуск брокера, а «горячие» читатели (те, кто у хвоста лога) вообще не касаются диска. Отсюда практическое правило: брокеру дают маленький heap (6–8 ГБ), а всю остальную память отдают ОС.
  3. Zero-copy. Консьюмеру данные отдаются через sendfile(2): page cache → сокет, без копирования в user space и обратно. Это экономит два копирования и два переключения контекста. Оговорка: zero-copy отваливается, если включён TLS (данные надо шифровать в user space) или если брокеру приходится конвертировать формат записи под старого клиента. Отсюда и ответ на вопрос «почему включение TLS уронило пропускную способность вдвое».
  4. Батчинг и единый формат. Продюсер собирает record batch, сжимает его целиком (а не каждое сообщение), брокер кладёт батч на диск как есть и как есть отдаёт консьюмеру. Распаковки и перепаковки на брокере нет, и стоимость сетевого и дискового вызова размазывается на тысячи сообщений.
Что ещё можно упомянуть

Бинарный протокол без текстового парсинга. Индексов по содержимому нет: искать сообщение Kafka умеет только по offset и времени, за счёт этого она и быстрая. Масштабирование горизонтальное, через партиции, и на пути данных нет общего координатора. И offset потребителя хранится отдельно от сообщений, в __consumer_offsets, так что брокеру не нужно менять уже записанные данные, чтобы отметить «прочитано». В RabbitMQ, наоборот, состояние сообщения меняется, и потому его модель принципиально дороже.

Producer: путь одной записи

Send() в Kafka-клиенте не отправляет ничего. Он сериализует ключ и значение, спрашивает у partitioner номер партиции, кладёт запись в аккумулятор (буфер батчей, по батчу на партицию) и возвращает управление. Отдельный поток-отправитель забирает готовые батчи, группирует их по брокерам-лидерам и шлёт один запрос на брокер.

acks и min.insync.replicas

acksКогда считается успехомЧто теряемСтоимость
0Как только запись ушла в сокетВсё: брокер мог быть недоступен, ошибку не увидимМаксимальная скорость, fire-and-forget
1Лидер записал в свой лог (в page cache)Лидер упал до репликации — запись пропалаОдна сетевая задержка
all (-1)Все реплики из текущего ISR подтвердилиТолько при одновременной потере всего ISR+ время репликации на самую медленную реплику ISR
Главный подвох

acks=all сам по себе ничего не гарантирует. Если реплики отстали и ISR схлопнулся до одной (самого лидера), acks=all становится ровно acks=1. Страхует от этого min.insync.replicas=2 на топике: при меньшем размере ISR продюсер получает NotEnoughReplicasException и запись не проходит. Рабочая тройка для прода: replication.factor=3, min.insync.replicas=2, acks=all. С ней кластер переживает падение одного брокера без потерь и без остановки записи. Частая ошибка: поставить min.insync.replicas=3 при RF=3. Тогда любой рестарт брокера останавливает запись в топик.

Идемпотентный продюсер

Продюсер отправил батч, брокер записал, ответ потерялся в сети, продюсер повторил, и в логе две одинаковые записи. Идемпотентность (enable.idempotence=true, включена по умолчанию с Kafka 3.0) это лечит: продюсер при старте получает PID (producer id) и ведёт sequence number отдельно на каждую партицию. Брокер помнит последний принятый номер по паре (PID, партиция): пришёл тот же номер — DuplicateSequenceNumber, батч молча отбрасывается с успешным ответом; пришёл номер с разрывом — OutOfOrderSequenceException.

Побочный эффект, за который её и любят: идемпотентность чинит порядок при ретраях. Без неё при max.in.flight.requests.per.connection > 1 повторно отправленный первый батч приедет после второго. С ней брокер по sequence number отклоняет батч с разрывом в номерах, продюсер отправляет заново по порядку, и до 5 in-flight безопасно. Границы гарантии: одна сессия продюсера, одна партиция. Перезапустили процесс — новый PID, дедупликация не работает. На другую партицию она тоже не распространяется.

Ключ, партиция и порядок

// Ключ есть   -> partition = murmur2(key) % numPartitions  (детерминированно)
// Ключа нет   -> Java-клиент: sticky partitioner, льёт в одну партицию, пока не закроется
//                батч, и переключается на следующую. kafka-go Murmur2Balancer берёт
//                для записи без ключа случайную партицию.
w := &kafka.Writer{
    Addr:     kafka.TCP("kafka-1:9092", "kafka-2:9092"),
    Topic:    "orders",
    Balancer: &kafka.Murmur2Balancer{}, // как в Java-клиенте. kafka.Hash{} считает FNV-1a
                                      // и совпадает с Sarama, но не с Java
    RequiredAcks: kafka.RequireAll,   // acks=all
    Async:    false,                  // Async=true = потеря при падении процесса
    BatchTimeout: 10 * time.Millisecond, // аналог linger.ms
    BatchSize:    100,
    Compression:  kafka.Lz4,
}
err := w.WriteMessages(ctx, kafka.Message{
    Key:   []byte(orderID),           // всё про один заказ -> одна партиция -> порядок
    Value: payload,
    Headers: []kafka.Header{
        {Key: "event_id",    Value: []byte(eventID)},   // для дедупликации у потребителя
        {Key: "traceparent", Value: []byte(traceID)},   // трейсинг сквозь брокер
    },
})
Три следствия ключа
  • Порядок. Все записи с одним ключом попадают в одну партицию, и порядок для этой сущности гарантирован.
  • Перекос (hot partition). Если 40 % трафика идёт от одного крупного клиента, его партиция станет узким местом, и добавление консьюмеров не поможет. Лечится составным ключом (tenant_id + bucket) ценой потери порядка внутри тенанта.
  • Изменение числа партиций ломает раскладку. murmur2(key) % N при другом N даёт другую партицию, и события одного заказа окажутся в двух логах. Уменьшить число партиций нельзя вообще. Увеличить можно, но это разовая осознанная операция, а не «докинем на пиковую нагрузку». Обычно партиции закладывают с запасом сразу.

Батчинг, linger.ms и компрессия

batch.size (16 КБ) ограничивает батч на партицию по байтам. linger.ms (с Kafka 4.0 по умолчанию 5 мс, раньше 0) говорит, сколько ждать доукомплектования батча перед отправкой. Даже ноль не означает «без батчинга»: батч всё равно собирается за время, пока предыдущий запрос в полёте. Дешевле всего поднять пропускную способность в разы, увеличив linger.ms до 5–20 мс: батчи крупнее, компрессия эффективнее, запросов меньше. Платим фиксированной добавкой к задержке.

Сжимается батч целиком, и чем он крупнее, тем выгоднее компрессия. На практике lz4 почти всегда правильный дефолт (быстрый, приличная степень); zstd (с Kafka 2.1) жмёт заметно лучше при чуть большем CPU и хорош для JSON и дорогого трафика между зонами; gzip самый медленный, брать его не стоит; snappy остался легаси-компромиссом. Сжатый батч так и лежит на диске и в том же виде уходит консьюмеру: брокер его не распаковывает (кроме случаев, когда включена валидация или конвертация формата).

Consumer groups и ребалансировка

Consumer group (группа потребителей) объединяет несколько процессов, которые вместе читают один топик и делят его партиции: каждая партиция достаётся ровно одному участнику группы. Группа распараллеливает обработку, не ломая порядок внутри партиции. Ребалансировкой (rebalance) называют перераздачу партиций между участниками: кто-то пришёл, кто-то ушёл или перестал отвечать, и надо заново решить, кто какие партиции читает. Похоже на официантов, которые делят столы в зале: вышел один покурить, его столы срочно раздают остальным, и на время передела обслуживание встаёт.

Консьюмеры объединяются в группу по group.id. Один брокер становится для этой группы group coordinator (выбирается по хешу group.id): ведёт список членов, раздаёт партиции и хранит offset. А вот считает раздачу лидер группы, один из консьюмеров: координатор присылает ему список всех членов и подписок, тот раскладывает партиции по выбранной стратегии и возвращает план. Логика назначения остаётся в клиенте, и её можно менять без обновления кластера.

Ребалансировка запускается, когда меняется состав группы или набор партиций:

  • новый консьюмер вызвал Subscribe и прислал JoinGroup;
  • консьюмер вышел штатно (Close шлёт LeaveGroup);
  • консьюмер перестал слать heartbeat дольше session.timeout.ms (45 с в новых версиях, heartbeat идёт из отдельного потока каждые heartbeat.interval.ms = 3 с);
  • консьюмер не вызвал poll() дольше max.poll.interval.ms (5 минут) — такой считается зависшим, даже если heartbeat идёт;
  • в топик добавили партиции или изменилась подписка по регулярке.

Стратегии назначения

СтратегияКак делитКогда брать
RangeAssignor (дефолт легаси)Партиции каждого топика по порядку режутся диапазонамиКогда нужна ко-локация одинаковых партиций разных топиков (join по ключу). Даёт перекос: при 3 партициях и 2 консьюмерах первый берёт 2
RoundRobinAssignorВсе партиции всех топиков по кругуРовный баланс, но при любом изменении раскладка едет целиком
StickyAssignorРовный баланс + минимум переездов относительно прошлого планаПочти всегда лучше RoundRobin
CooperativeStickyAssignorТо же плюс инкрементальная ребалансировка (KIP-429)Дефолт для новых сервисов: нет глобальной паузы
Eager (Range / RoundRobin / Sticky): группа встаёт целиком работает ребалансировка: revoke ВСЕГО работает по новому плану C1 P0 P1 P2 простой: не читает ничего P0 P1 C2 P3 P4 P5 простой: не читает ничего P2 P3 C3 ещё не в группе JoinGroup, ждёт SyncGroup P4 P5 Пауза длится от JoinGroup до SyncGroup всех участников. Один зависший консьюмер тянет её до rebalance.timeout.ms — то есть до минут. Cooperative sticky (KIP-429): отзываются только переезжающие партиции работает ребалансировка: revoke только лишнего работает по новому плану C1 P0 P1 P2 P0 P1 читает дальше, отдал P2 P0 P1 C2 P3 P4 P5 P3 читает, P4 P5 отданы P2 P3 C3 ещё не в группе ждёт свои партиции P4 P5 Останавливается только то, что реально переезжает. Цена — два раунда ребалансировки вместо одного, зато без общей паузы.
Чем опасна ребалансировка. При eager-стратегии простаивает вся группа, и лаг за это время растёт на весь входящий поток. Кооперативная стратегия убирает глобальную паузу, оставляя её только владельцам переезжающих партиций.
Как выглядит rebalance storm в проде

Обработка одного сообщения замедлилась (тормозит внешний API), консьюмер не успевает вызвать poll() за max.poll.interval.ms, координатор считает его мёртвым и запускает ребалансировку. Она отбирает партиции у всех, лаг растёт, после перераспределения каждый получает больше работы, снова не успевает, и цикл замыкается: группа ребалансируется бесконечно и не обрабатывает ничего. Лечение по порядку: уменьшить max.poll.records (обрабатывать меньше за раз), увеличить max.poll.interval.ms, вынести медленную работу в отдельный пул с pause()/resume() партиции, включить CooperativeStickyAssignor и static membership (group.instance.id, KIP-345). Тогда рестарт пода при деплое вообще не вызывает ребалансировку, пока укладывается в session.timeout.ms.

Правило про параллелизм

Максимум параллельных консьюмеров в одной группе = число партиций топика. При 12 партициях тринадцатый консьюмер будет просто простаивать с пустым назначением. Отсюда правило планирования: партиций закладывают с запасом на будущий рост (например, 24–48 на топик среднего размера), потому что увеличить их можно, но это ломает раскладку по ключу, а уменьшить нельзя вообще. Но и партиции не бесплатны: каждая стоит файловых дескрипторов, памяти на брокере и времени на восстановление при failover.

Коммит offset

Committed offset равен offset следующей записи, которую надо прочитать, то есть «последняя обработанная + 1». Хранится он в служебном топике __consumer_offsets (50 партиций, cleanup.policy=compact) под ключом (group.id, topic, partition). Компакция нужна, чтобы топик не рос бесконечно: по каждому ключу важно только последнее значение.

Авто-коммит (enable.auto.commit=true, интервал 5 с) коммитит не по таймеру в фоне, а внутри вызова poll(): если с прошлого коммита прошло больше интервала, коммитится всё, что было выдано предыдущим poll. Пока цикл синхронный, эти записи к тому моменту уже обработаны, и выходит тот же at-least-once. Опасно, когда обработку отдают в другие горутины: тогда коммит подтверждает записи, обработку которых ты ещё не завершил, и падение в этом окне теряет сообщение навсегда. В librdkafka и в kafka-go с ReadMessage offset запоминается уже при выдаче, так что коммит может обогнать обработку и без всяких горутин.

Упал ДО коммита ручной коммит после обработки время poll() отдал запись с offset 42 обработка: списали 100 руб процесс убит: OOM, деплой, паника рестарт: последний коммит = 42 poll() снова отдаёт offset 42 списали 100 руб ВТОРОЙ раз дубликат — это и есть at-least-once Упал ПОСЛЕ коммита тот же ручной коммит время poll() отдал запись с offset 42 обработка: списали 100 руб commit(43) записан в __consumer_offsets процесс убит рестарт: последний коммит = 43 poll() отдаёт offset 43, идём дальше ни потерь, ни дублей А если enable.auto.commit=true с интервалом 5 секунд? Коммит внутри poll() подтверждает выданное прошлым poll — при синхронном цикле оно обработано. Если обработка ушла в фон, падение в этом окне — ПОТЕРЯ: offset уже сдвинут.
Где рождается дубликат и где потеря. Ручной коммит после обработки даёт at-least-once, и в худшем случае запись повторится. Авто-коммит опасен там, где подтверждение обгоняет обработку: при обработке в фоне, в librdkafka, в kafka-go с ReadMessage. Тогда он превращается в at-most-once с реальной потерей данных.
// segmentio/kafka-go: ручной коммит после успешной обработки.
r := kafka.NewReader(kafka.ReaderConfig{
    Brokers:  []string{"kafka-1:9092"},
    GroupID:  "billing",              // GroupID включает управление группой
    Topic:    "orders",
    MinBytes: 10e3,                   // fetch.min.bytes
    MaxBytes: 10e6,
    MaxWait:  500 * time.Millisecond, // fetch.max.wait.ms
    // CommitInterval = 0: CommitMessages коммитит синхронно
})
defer r.Close()

for {
    // FetchMessage, а не ReadMessage: ReadMessage коммитит сам
    m, err := r.FetchMessage(ctx)
    if err != nil {
        return err                    // ctx отменён или reader закрыт
    }
    if err := handle(ctx, m); err != nil {
        // continue здесь потерял бы сообщение: коммит следующего сдвинет
        // offset и за него. Откладываем в retry-топик, потом коммитим.
        if err := toRetry(ctx, m); err != nil {
            return err                // не отложили — выходим без коммита
        }
    }
    if err := r.CommitMessages(ctx, m); err != nil {
        slog.Error("commit", "err", err)
    }
}
commitSync vs commitAsync

commitSync блокирует и ретраит до успеха. Надёжно, но каждый батч получает лишний round-trip. commitAsync намеренно не ретраит: запоздалый ретрай мог бы записать меньший offset поверх большего и вызвать переобработку. На практике делают так: асинхронный коммит в цикле, а синхронный в defer при выходе и в колбэке onPartitionsRevoked перед ребалансировкой. Коммитить каждое сообщение дорого. Обычно коммитят раз в N записей или раз в секунду и сознательно расширяют окно возможных дублей.

Retention и log compaction

Retention («удержание», срок хранения) определяет, сколько времени или байт Kafka держит записи, прежде чем выбросить. Log compaction («уплотнение лога») убирает иначе: вместо «выбросить старое по времени» брокер выбрасывает устаревшие версии, оставляя по каждому ключу только последнюю запись. Аналогия из магазина: retention говорит «храним чеки за три месяца, дальше в мусор», а compaction — «по каждому товару держим только актуальный ценник, старые выбрасываем, сколько бы им ни было лет».

Общее для обеих политик: Kafka удаляет данные по cleanup.policy, а не по факту прочтения. Никто не прочитал — всё равно удалит по сроку; все прочитали — всё равно будет хранить до срока.

  • delete (дефолт): целые сегменты удаляются, когда все их записи старше retention.ms (7 суток) либо когда суммарный размер партиции превысил retention.bytes. Обе границы работают одновременно: что наступит раньше. Считать надо на партицию, а не на топик: retention.bytes=1GB при 24 партициях дают 24 ГБ на топик, и ещё умножить на replication factor.
  • compact — размер ограничен не временем, а числом уникальных ключей: для каждого ключа хранится последнее значение. Топик превращается из журнала изменений в материализованный снимок состояния, который можно перечитать с нуля и восстановить всю таблицу.
  • compact,delete включает обе: последнее значение по ключу, но не старше retention. Так делают, например, «состояние за последние 30 дней».
До compaction — полная история изменений, cleanup.policy=compact offset 0 1 2 3 4 5 6 7 8 9 key value k1 k2 k1 k3 k2 k1 k4 k3 k2 k1 v1 a1 v2 c1 a2 v3 d1 null a3 v4 Зелёные — последняя запись по своему ключу. Жёлтая — tombstone: запись с value=null, то есть «ключ удалён». log cleaner После compaction — по одному последнему значению на ключ, offsets не перенумеровываются эти offsets теперь просто отсутствуют k4 k3 k2 k1 d1 null a3 v4 6 7 8 9 Tombstone держится ещё delete.retention.ms (сутки) и только потом исчезает — чтобы отставшие консьюмеры успели узнать об удалении ключа.
Compaction. Лог перестаёт быть журналом и становится снимком: «текущее значение каждого ключа». Offsets при этом сохраняются, и в нумерации появляются дыры — консьюмер должен быть к этому готов и не рассчитывать, что offset+1 всегда существует.
Зачем нужна компакция
  • Восстановление состояния с нуля. Новый сервис подписался на compacted-топик user-profiles, прочитал его целиком и получил актуальные профили всех пользователей, не трогая чужую БД. На этом построены CQRS-проекции и KTable в Kafka Streams.
  • Служебные топики самой Kafka. __consumer_offsets и __transaction_state compacted, иначе они росли бы вечно.
  • Change data capture. Debezium публикует изменения таблицы с primary key в качестве ключа сообщения; compacted-топик оказывается зеркалом таблицы.
Что ломает компакцию

Сообщения без ключа. В compacted-топике запись с key=null считается ошибкой: сжимать её не по чему. Ожидание, что дубликатов не останется совсем. Компакция идёт в фоне и «лениво»: активный сегмент не трогается никогда, а остальные чистятся, только когда доля «грязных» байт превысит min.cleanable.dirty.ratio (0.5). Поэтому в compacted-топике вполне может лежать три версии одного ключа, и потребитель обязан корректно обрабатывать их по порядку. И tombstone надо уметь читать: сообщение с value=null означает «удалить», а не «пустое значение».

Consumer lag

Lag («отставание») показывает, сколько записей в партиции уже написано, но этой группой ещё не обработано. Формально lag = log end offset − committed offset для каждой пары (группа, партиция): конец лога минус закладка. Это как стопка непрочитанных писем в ящике: важно не то, что там что-то лежит, а растёт она или тает.

По лагу судят о здоровье асинхронной системы: он показывает не «сколько сообщений в очереди», а насколько потребитель отстал от реальности.

Чем мерить

  • kafka-consumer-groups.sh --describe --group billing — быстро и руками, показывает лаг по каждой партиции.
  • kafka-exporter отдаёт в Prometheus метрику kafka_consumergroup_lag с разбивкой по партициям (у kafka-lag-exporter она называется kafka_consumergroup_group_lag). Алертить надо на производную (лаг растёт N минут подряд), а не на абсолютное значение: 50 000 при потоке 100k/с дают всего полсекунды отставания, и это нормально.
  • Lag в единицах времени лучше всего подходит для SLO: «отстаём на 8 секунд» понятно бизнесу, а «отстаём на 400 000 сообщений» нет. Считается по .timeindex или разностью timestamp последней обработанной записи и текущего времени.
  • Burrow (LinkedIn) смотрит не на сам лаг, а на его тренд, и выдаёт статус OK/WARN/ERR без ручных порогов.

Лаг растёт: что смотреть по порядку

СимптомПричинаЧто делать
Растёт на всех партициях равномерноКонсьюмеров банально не хватает под потокДобавить инстансы до числа партиций; дальше — увеличивать партиции
Растёт на одной-двух партицияхПерекос ключа: hot partition, крупный тенантИзменить ключ (составной, с бакетом), вынести тенанта в отдельный топик
Лаг пилой: падает и снова растётRebalance stormmax.poll.records вниз, cooperative sticky, static membership
Консьюмеров хватает, CPU простаиваетЖдём внешнюю систему: БД, HTTP, блокировкиБатчить запись в БД, поднять пул соединений, параллелить внутри по hash(key)
Лаг скачком уменьшился самRetention съел непрочитанноеПроверить records-lag вместе с OffsetOutOfRange: данные потеряны, это инцидент
Лаг растёт после деплояНовая версия обрабатывает медленнее, либо сбросился offsetСравнить process-latency до и после; проверить auto.offset.reset

Ретраи и DLQ

В Kafka нет ни встроенного requeue, ни встроенного DLQ. Нельзя «вернуть сообщение в очередь»: указатель всего один, offset, и он либо двигается, либо нет. Остаются два примитивных подхода, и оба плохи:

  • Не коммитить и читать заново. Партиция встаёт колом на одном ядовитом сообщении (head-of-line blocking): за ним копятся тысячи корректных, а мы в бесконечном цикле падаем на одном. Вдобавок так быстро превышается max.poll.interval.ms, и начинается ребалансировка.
  • Ретраить в цикле внутри обработчика. То же самое: пока спим между попытками, не вызываем poll(). Работает только при коротком общем бюджете (несколько секунд) и обязательно с pause() партиции и дальнейшими вызовами poll().

В проде делают цепочку retry-топиков с возрастающей задержкой и финальный DLQ.

# Топология
orders                 -> основной консьюмер
orders.retry.5s        -> консьюмер, который перед обработкой спит до (produced_at + 5s)
orders.retry.1m
orders.retry.10m
orders.DLQ             -> никто не читает автоматически; читают глазами и руками

# Заголовки, которые надо переносить вместе с сообщением:
x-original-topic       # откуда пришло изначально
x-retry-count          # сколько попыток уже было
x-first-failed-at      # когда сломалось впервые
x-last-error           # тип и текст последней ошибки
x-original-offset      # координаты оригинала, чтобы найти его в логе
// Консьюмер retry-топика спит до времени готовности сообщения.
// У всех сообщений в orders.retry.1m одинаковая задержка,
// поэтому топик упорядочен по времени готовности и достаточно
// подождать голову: приоритетный разбор не нужен.
func (c *RetryConsumer) run(ctx context.Context) error {
    for {
        m, err := c.reader.FetchMessage(ctx)
        if err != nil {
            return err
        }
        ready := m.Time.Add(c.delay)
        if d := time.Until(ready); d > 0 {
            // kafka-go держит членство в группе фоновым heartbeat, спать можно.
            // Java-клиенту на это время нужны pause() и дальнейшие вызовы poll().
            select {
            case <-time.After(d):
            case <-ctx.Done():
                return ctx.Err()
            }
        }
        var perr error
        switch err := c.handle(ctx, m); {
        case err == nil:
            // успех
        case errors.Is(err, ErrPermanent) || retryCount(m) >= c.maxRetries:
            perr = c.produce(ctx, c.dlqTopic, m) // в DLQ без дальнейших попыток
        default:
            perr = c.produce(ctx, c.nextTopic, m) // на следующую ступень задержки
        }
        if perr != nil {
            return perr // не переложили — не коммитим, прочитаем снова
        }
        if err := c.reader.CommitMessages(ctx, m); err != nil {
            return err
        }
    }
}
Что сказать про DLQ
  • Ошибки бывают двух видов, и путать их нельзя. Транзиентные (таймаут БД, 503 от соседа) ретраят. Постоянные (невалидный JSON, нет обязательного поля, бизнес-запрет) ретраить бессмысленно, их сразу отправляют в DLQ. Без этого разделения retry-цепочка просто откладывает неизбежное и тратит ресурсы.
  • Retry-топики ломают порядок. Сообщение, ушедшее на ретрай, вернётся позже своих соседей по ключу. Если порядок критичен, придётся останавливать обработку всего ключа до разбора или осознанно принять расхождение.
  • DLQ без алерта и без процесса разбора превращается в мусорку. Нужны метрика количества, алерт на «первое сообщение за N минут», ретеншен побольше основного топика и инструмент повторной подачи (reprocess из DLQ обратно в основной топик).

Транзакции и exactly-once

Транзакции Kafka закрывают ровно один сценарий: read-process-write внутри Kafka. Прочитали из топика A, посчитали, записали в топики B и C, и всё это либо видно потребителям целиком, либо не видно вовсе.

Механика по шагам:

  1. Продюсер задаёт transactional.id, стабильный идентификатор, переживающий рестарт (например, billing-worker-p3, привязанный к партиции, а не случайный UUID).
  2. initTransactions() находит transaction coordinator и получает PID с новым epoch. Все старые продюсеры с тем же transactional.id и меньшим epoch получают ProducerFenced. Так работает zombie fencing: подвисший старый инстанс не сможет дописать в транзакцию после того, как его заменили.
  3. Записи внутри транзакции пишутся в целевые партиции сразу, но помечены как транзакционные.
  4. sendOffsetsToTransaction(offsets, consumer.groupMetadata()) кладёт коммит входных offset в ту же транзакцию, и связка «прочитал и записал» становится атомарной.
  5. commitTransaction(): координатор пишет решение в __transaction_state и рассылает во все затронутые партиции control batch (маркер COMMIT или ABORT).

На стороне чтения нужен isolation.level=read_committed. Такой консьюмер видит записи только до LSO (last stable offset, offset первой ещё не завершённой транзакции) и отфильтровывает записи прерванных транзакций по маркерам. А по умолчанию стоит read_uncommitted, отсюда классический подвох «включили транзакции, а дубликаты и грязные данные остались»: консьюмера просто забыли переключить.

Граница exactly-once

Как только обработчик делает что-то за пределами Kafka (пишет в Postgres, дёргает платёжный шлюз, отправляет письмо), транзакция Kafka это действие не покрывает. Откатить письмо нельзя. Поэтому для «Kafka → БД» правильный ответ не «транзакции», а идемпотентная запись: апсерт по бизнес-ключу либо inbox-таблица, где вставка ключа сообщения и бизнес-данные идут в одной транзакции Postgres, а offset коммитится после. Цена EOS тоже не нулевая: дополнительный round-trip к координатору, control-батчи в логе и, что неприятнее всего, read_committed консьюмер не может читать дальше LSO: одна зависшая транзакция блокирует чтение партиции до transaction.timeout.ms.

Библиотеки Go

БиблиотекаПлюсыМинусыКогда брать
segmentio/kafka-goЧистый Go, простое API (Reader/Writer), лаконичный код, хорошая работа с contextНет транзакций (EOS) и cooperative-ребалансировки, часть тонких настроек недоступнаДефолт для типового сервиса: читать, обрабатывать, писать. 80 % задач
IBM/sarama (бывш. Shopify)Самая старая и распространённая, покрывает почти весь протокол, много готовых примеровМногословный API, легко выстрелить в ногу (ConsumerGroupHandler, ручное управление сессией), исторически шероховатая обработка ошибокЛегаси-кодовая база; когда нужна функция, которой нет в kafka-go
twmb/franz-goСамая полная поддержка протокола на чистом Go: транзакции, EOS, cooperative sticky, KIP-848, самые высокие бенчмаркиAPI мощнее и потому объёмнее; меньше готовых рецептов в интернетеКогда нужны транзакции, максимальная производительность или свежие KIP
confluent-kafka-goОбёртка над librdkafka: эталонное поведение, все фичи Confluent, Schema Registry из коробкиcgo: усложняет кросс-компиляцию и статическую сборку, тяжелее контейнер, сложнее профилироватьЕсли в компании стандарт на Confluent-стек и cgo не пугает
Как отвечать про библиотеки

Ответ «работал с segmentio/kafka-go, потому что чистый Go и простой Reader; когда понадобились транзакции, посмотрели на franz-go» звучит на порядок лучше, чем перечисление названий. И полезно назвать конкретную грабельку: в kafka-go ReadMessage коммитит offset сам, поэтому для at-least-once нужен FetchMessage + CommitMessages; а Writer с Async: true молча теряет сообщения при падении процесса и не возвращает ошибку доставки.

Вопросы

9
Суть: Kafka — распределённый append-only лог. Топик логичен, партиция физична: это каталог с сегментами на диске конкретного брокера. Offset — номер записи внутри партиции, и только внутри неё есть порядок.

Снизу вверх

  • Сегмент складывается из файла .log и двух разреженных индексов: .index (offset → байт) и .timeindex (время → offset). В имени файла стоит base offset. Пишут только в активный сегмент, он ротируется по segment.bytes (1 ГБ) или segment.ms (7 суток).
  • Партиция — упорядоченная последовательность сегментов, единица параллелизма, репликации и упорядочивания. Живёт на одном брокере-лидере плюс на фолловерах.
  • Топик просто объединяет партиции под одним именем. Физически его не существует.
  • Брокер держит диски и по одним партициям работает лидером, по другим фолловером. Метаданными ведает кворум контроллеров KRaft, обычно на выделенных узлах.

Offset

Монотонный номер, который назначает лидер партиции. Не глобальный: offset 100 в P0 никак не соотносится с offset 100 в P1. Ищется offset бинарным поиском по именам сегментов, затем по .index, затем досканированием внутри 4 КБ. Поэтому чтение с начала лога стоит примерно столько же, сколько чтение с конца, и «перечитать всё заново» остаётся обычной операцией, а не аварией.

Репликация

replication.factor реплик на партицию, одна из них лидер. Фолловеры тянут данные тем же fetch-протоколом. В ISR входят реплики, отставшие не больше чем на replica.lag.time.max.ms (30 с). High watermark равен максимальному offset, который есть у всех реплик ISR; консьюмер видит только до него, иначе после смены лидера прочитал бы данные, которых в новом лидере нет.

Чем добить

Упомянуть переход с ZooKeeper на KRaft (KIP-500): метаданные и выборы лидера теперь в Raft-логе контроллеров внутри самой Kafka. GA с 3.3, миграция с 3.6, в Kafka 4.0 ZooKeeper удалён. Практический эффект: failover контроллера занимает доли секунды вместо десятков, и кластер держит миллионы партиций. И назвать unclean.leader.election.enable=false как выбор CP: при потере всего ISR партиция становится недоступной, а не расходящейся.

Суть: acks определяет, что считать успешной записью; он работает только в паре с min.insync.replicas. Идемпотентность (PID + sequence number) убирает дубли от ретраев и чинит порядок. Ключ детерминированно выбирает партицию — отсюда и порядок, и перекос.

acks

0 — ушло в сокет, ошибок не увидим вообще. 1 — лидер записал к себе (в page cache), падение лидера до репликации = потеря. all — подтвердили все реплики текущего ISR. Последнее и есть ловушка: если ISR схлопнулся до одного лидера, acks=all эквивалентен acks=1. Спасает min.insync.replicas=2 — при меньшем ISR продюсер получит NotEnoughReplicasException. Прод-конфиг: RF=3, min.insync.replicas=2, acks=all. Ставить min.insync=3 при RF=3 нельзя, потому что любой рестарт брокера остановит запись.

Идемпотентный продюсер

enable.idempotence=true (дефолт с Kafka 3.0). Продюсер получает PID и нумерует записи sequence number по каждой партиции; брокер помнит последний принятый номер и отбрасывает повтор, отвечая успехом. Идемпотентность убирает дубли от собственных ретраев продюсера, в пределах одной сессии и одной партиции. Перезапустили процесс — новый PID, гарантии нет. Полезный побочный эффект: с идемпотентностью безопасно держать max.in.flight.requests.per.connection до 5, потому что брокер восстанавливает порядок по sequence number; без неё ретрай первого батча приедет после второго.

Ключ, партиция, порядок

Есть ключ — murmur2(key) % numPartitions, детерминированно. Нет ключа — sticky partitioner: льём в одну партицию, пока не закроется батч, потом в следующую (батчи крупнее, задержка ниже, раскладка всё равно ровная). Три следствия: порядок гарантирован для одного ключа; неудачный ключ даёт hot partition; изменение числа партиций меняет раскладку и разрывает историю сущности между двумя логами.

Частая ошибка на практике

Разные клиенты считают хеш по-разному, и одного «включить хеширование» мало — важно, какой именно хеш. В segmentio/kafka-go по умолчанию стоит round-robin, а не хеш. Но и явный &kafka.Hash{} не решает задачу: по документации он считает FNV-1a и совпадает с Sarama, а не с Java. Java-клиент и librdkafka используют murmur2, и в kafka-go ему соответствует &kafka.Murmur2Balancer{}. Возьмёшь не тот — ключ вроде бы задан, а сообщения одного заказа расползаются по партициям относительно соседнего сервиса, и «порядок гарантирован» превращается в тыкву. Проверяется это за минуту: прогнать сто сообщений с одним ключом и посмотреть распределение по партициям.

Суть: внутри группы партиция принадлежит ровно одному консьюмеру, поэтому параллелизм упирается в число партиций. Ребалансировка при eager-стратегиях останавливает всю группу, и в патологическом случае группа зацикливается на перераспределении вместо работы.

Кто и как распределяет

Координатор группы (брокер, выбранный по хешу group.id) ведёт членство и хранит offset. Само распределение считает лидер группы, то есть один из консьюмеров: координатор присылает ему список членов и подписок, тот возвращает план, координатор рассылает его в SyncGroup. Логика назначения живёт в клиенте, поэтому её можно менять, не обновляя кластер.

Триггеры ребалансировки

  • консьюмер зашёл (JoinGroup) или вышел (LeaveGroup при Close);
  • пропали heartbeat дольше session.timeout.ms (45 с);
  • не вызывался poll() дольше max.poll.interval.ms (5 мин). В проде это самый частый триггер: обработка затянулась, а heartbeat при этом шёл;
  • добавили партиции в топик.

Чем опасна

При eager-стратегии (Range, RoundRobin, Sticky) все консьюмеры отзывают все партиции и ждут нового плана. Для группы это stop-the-world: лаг растёт на весь входящий поток, незакоммиченная работа переделывается заново другим консьюмером (дубликаты), а прогретые локальные кэши обесцениваются. Хуже всего rebalance storm: медленная обработка → превышение max.poll.interval.ms → ребалансировка → у оставшихся ещё больше партиций → ещё медленнее → снова ребалансировка.

Лечение

  • CooperativeStickyAssignor (KIP-429): отзываются только переезжающие партиции, общей паузы нет. Дефолт для новых сервисов.
  • Static membership (group.instance.id, KIP-345): при рестарте пода консьюмер возвращается со своим прежним назначением, ребалансировки не происходит вовсе, если уложился в session.timeout.ms. Резко упрощает деплой.
  • Уменьшить max.poll.records, вынести медленную работу в отдельный пул с pause()/resume().
  • С Kafka 4.0 доступен новый серверный протокол ребалансировки KIP-848: план считает брокер, переходы инкрементальны.
Формулировка про параллелизм

«Партиция — единица параллелизма. Консьюмеров в группе может быть больше, чем партиций, но лишние будут простаивать с пустым назначением. Значит, потолок горизонтального масштабирования группы задаётся при создании топика.» Дальше уместно добавить: партиции можно только увеличивать, и это ломает раскладку по ключу, поэтому их закладывают с запасом; а внутри одного консьюмера можно параллелить обработку по hash(key) % workers, сохраняя порядок на ключ. Так пробивают потолок, не меняя топик.

Суть: ручной коммит после обработки = at-least-once: падение до коммита даёт дубликат, падение после — ничего страшного. Авто-коммит при синхронном цикле даёт тот же at-least-once, а при обработке в фоне подтверждает ещё не обработанное и даёт потерю.

Где хранится

Committed offset равен offset следующей записи для чтения, то есть «последняя обработанная + 1». Живёт в служебном топике __consumer_offsets (50 партиций, cleanup.policy=compact) под ключом (group.id, topic, partition). Компакция нужна, чтобы топик не рос: значимо только последнее значение.

Падение ДО коммита

Прочитали offset 42, обработали (списали деньги), упали до commit(43). После рестарта координатор отдаёт последний закоммиченный offset 42, консьюмер читает ту же запись и списывает деньги второй раз. Для at-least-once это нормальное поведение, и защищает от него только идемпотентный обработчик, а не настройки Kafka.

Падение ПОСЛЕ коммита

Обработали, закоммитили 43, упали. После рестарта читаем с 43 — ни потерь, ни дублей. При ручном управлении правилен только один порядок: «сначала обработать, потом коммитить».

Авто-коммит

enable.auto.commit=true, auto.commit.interval.ms=5000. Коммит происходит внутри poll(): если интервал истёк, подтверждается всё, что выдал предыдущий poll. При синхронном цикле эти записи уже обработаны, и получается at-least-once. Если же обработка идёт в фоне, подтвердятся и записи, которые ещё в работе: упал в этом окне — offset уже сдвинут, брокер ничего не переотправит, сообщение потеряно. Авто-коммит допустим только там, где потеря дешевле кода (метрики, логи).

Тонкости ручного коммита

  • commitSync блокирует и ретраит; commitAsync намеренно не ретраит, иначе повтор мог бы записать меньший offset поверх большего и вызвать переобработку.
  • Рабочий рецепт: commitAsync в цикле + commitSync в defer при выходе и в колбэке перед отзывом партиций.
  • Коммитить каждое сообщение дорого; обычно коммитят раз в N записей или раз в секунду, сознательно расширяя окно возможных дублей.
  • В segmentio/kafka-go метод ReadMessage коммитит сам; для at-least-once нужен FetchMessage + явный CommitMessages.
Follow-up, который любят

«А если консьюмер работает в несколько горутин и коммитит по завершении каждой?» Тогда можно закоммитить offset 50, пока обработка 47 ещё идёт, и падение потеряет 47–49. При параллельной обработке коммитить надо сам минимальный незавершённый offset: он и есть следующая запись для чтения; держать map незавершённых и двигать «водяной знак» только непрерывно. Такое проговорит только тот, кто реально писал консьюмеры.

Суть: Kafka удаляет данные по политике, а не по факту прочтения. delete режет старые сегменты по времени и размеру; compact оставляет по одному последнему значению на ключ, превращая лог в снимок состояния.

cleanup.policy=delete

Удаляются целые сегменты, все записи которых старше retention.ms (7 суток по умолчанию), либо когда размер партиции превысил retention.bytes. Обе границы действуют одновременно, срабатывает та, что наступит раньше. Считать надо на партицию: retention.bytes=1GB при 24 партициях и RF=3 дают 72 ГБ на кластер. Активный сегмент не трогается никогда, поэтому при слабом трафике данные живут дольше настройки.

cleanup.policy=compact

Фоновый log cleaner проходит по неактивным сегментам и оставляет для каждого ключа только последнюю запись. Offsets при этом не перенумеровываются: в логе появляются дыры, и консьюмер должен быть к этому готов. Чтобы удалить ключ, пишут tombstone, запись с value=null; она хранится ещё delete.retention.ms (сутки), чтобы отставшие потребители успели узнать об удалении, и только потом исчезает.

Зачем компакция

  • Новый сервис перечитывает compacted-топик с нуля и получает актуальный снимок всех сущностей. На этом стоят CQRS-проекции и KTable в Kafka Streams.
  • Сама Kafka использует её для __consumer_offsets и __transaction_state.
  • CDC: Debezium пишет изменения с primary key в ключе, и compacted-топик становится зеркалом таблицы.
Три ловушки
  • Компакция не гарантирует уникальность немедленно. Активный сегмент не чистится, остальные — только когда доля грязных байт превысит min.cleanable.dirty.ratio (0.5). Три версии одного ключа в топике для неё штатная ситуация.
  • Записи без ключа в compacted-топике недопустимы, их не по чему сжимать.
  • Если лаг «сам уменьшился», радоваться рано. Консьюмер отстал сильнее retention, непрочитанные сегменты удалили, он получил OffsetOutOfRange и по auto.offset.reset=latest прыгнул в конец, молча потеряв данные. Это инцидент, а не самоизлечение.
Суть: lag = log end offset − committed offset по каждой партиции. Алертить надо не на абсолютное значение, а на устойчивый рост и на лаг в секундах; лечение начинается с вопроса «лаг на всех партициях или на одной».

Чем мерить

  • kafka-consumer-groups.sh --describe --group billing — руками, по партициям.
  • kafka_consumergroup_lag из kafka-exporter в Prometheus — основной способ.
  • Клиентские JMX-метрики records-lag-max и fetch-latency показывают картину со стороны консьюмера, включая паузы на ребалансировку.
  • Burrow оценивает тренд лага и выдаёт статус без ручных порогов.
  • Лаг во времени («отстаём на 8 секунд») и есть правильная метрика для SLO; его считают по .timeindex или как разницу текущего времени и timestamp последней обработанной записи.

Диагностика по симптому

  • Растёт равномерно на всех партициях — не хватает мощности. Добавить консьюмеров до числа партиций, а дальше увеличивать партиции.
  • Растёт на одной-двух — перекос ключа, hot partition. Менять ключ (составной с бакетом) или выносить тяжёлого тенанта в свой топик.
  • Пилообразный график выдаёт rebalance storm. Смотреть логи координатора, уменьшать max.poll.records, включать cooperative sticky и static membership.
  • Если CPU консьюмера простаивает, упёрлись во внешнюю систему. Батчить запись в БД, увеличить пул соединений, параллелить внутри процесса по hash(key).
  • Лаг упал скачком без деплоя — retention съел непрочитанное. Инцидент.

Экстренные меры

Когда лаг измеряется часами и растёт, а быстро масштабироваться некуда: запустить временную группу-«ускоритель», которая только перекладывает сообщения в несколько параллельных топиков; временно отключить необязательные шаги обработки (обогащение, аналитику); если данные совсем неактуальны, осознанно перемотать offset на конец (kafka-consumer-groups.sh --reset-offsets --to-latest), явно зафиксировав, что часть событий не будет обработана. Последнее уже управленческое решение, а не техническое.

Суть: в Kafka нет ни requeue, ни встроенного DLQ — есть только offset. Стандартное решение — цепочка retry-топиков с растущей задержкой плюс финальный DLQ, с переносом метаданных об ошибке в заголовках.

Почему «просто не коммитить» не работает

Партиция встаёт колом на одном ядовитом сообщении — head-of-line blocking: за ним копятся тысячи корректных. Плюс бесконечный цикл падений быстро превышает max.poll.interval.ms и запускает ребалансировку, которая делает только хуже. Ретраить в цикле прямо в обработчике тоже нельзя дольше пары секунд, ведь во время сна мы не вызываем poll().

Схема с retry-топиками

Топики orders.retry.5s, orders.retry.1m, orders.retry.10m, orders.DLQ. Обработчик при транзиентной ошибке публикует сообщение в следующий retry-топик и коммитит оригинал — партиция освобождается сразу. Консьюмер retry-топика перед обработкой спит до produced_at + delay (в Java-клиенте — поставив партицию на pause() и продолжая вызывать poll()); поскольку у всех сообщений одного retry-топика одинаковая задержка, топик остаётся упорядоченным по времени готовности, и достаточно ждать голову.

Заголовки

Обязательно переносить: x-original-topic, x-original-offset, x-retry-count, x-first-failed-at, x-last-error. Без них DLQ бесполезен — непонятно, что именно сломалось и куда возвращать.

Разделение ошибок

Транзиентные (таймаут БД, 503 соседа, дедлок) имеет смысл ретраить. Постоянные (невалидная схема, отсутствующее обязательное поле, бизнес-запрет) отправляют сразу в DLQ: ретраи их не вылечат и только съедят ресурсы. В коде это разные типы ошибок, а не строковое сравнение.

Что делает ответ сильным
  • Назвать, что retry-топики ломают порядок: сообщение вернётся позже соседей по тому же ключу. Если порядок критичен, остаётся либо блокировать весь ключ до разбора, либо принять расхождение и сделать обработчик устойчивым по версии.
  • Сказать, что DLQ без алерта и процесса разбора становится просто мусоркой. Нужны метрика, алерт на первое сообщение за N минут, увеличенный retention и инструмент повторной подачи в основной топик.
  • Упомянуть альтернативу: Kafka Connect и Spring Kafka имеют DLQ и DefaultErrorHandler из коробки, а в чистом Go это всегда пишут руками, и такой код стоит вынести в общую библиотеку на всю компанию.
Суть: транзакции дают атомарность read-process-write внутри Kafka: запись в несколько партиций и коммит входных offset — одним решением. За пределы Kafka гарантия не распространяется.

Механика

  1. Продюсер задаёт стабильный transactional.id, переживающий рестарт.
  2. initTransactions() находит transaction coordinator и получает PID с новым epoch. Старые продюсеры с тем же id и меньшим epoch получают ProducerFenced, и zombie fencing отсекает подвисший предыдущий инстанс.
  3. Записи пишутся в партиции сразу, но помечены транзакционными.
  4. sendOffsetsToTransaction(offsets, consumer.groupMetadata()) включает коммит входных offset в ту же транзакцию; на этом и держится атомарность связки «прочитал и записал».
  5. commitTransaction(): координатор фиксирует решение в __transaction_state и пишет control batch (маркер COMMIT/ABORT) во все затронутые партиции.

Сторона чтения

Нужен isolation.level=read_committed. Такой консьюмер видит записи только до LSO (last stable offset, offset первой незавершённой транзакции) и отбрасывает записи прерванных транзакций. Дефолт у него read_uncommitted, и из-за забытой настройки случается классическое «включили EOS, а грязные данные всё равно видны».

Границы и цена

Гарантия действует, пока весь эффект обработки сводится к записям в Kafka. Письмо, платёж или запись в чужую БД откатить нельзя, значит exactly-once там достигается только идемпотентностью или inbox-таблицей. Цена: дополнительный round-trip к координатору на каждую транзакцию, control-батчи в логе, повышенная задержка и, главное, read_committed консьюмер не читает дальше LSO, поэтому одна зависшая транзакция блокирует чтение партиции до transaction.timeout.ms.

Как это звучит на собесе

«Exactly-once в Kafka — это не про доставку, а про атомарность конвейера внутри Kafka. Идемпотентный продюсер убирает дубли от ретраев, транзакции добавляют атомарность записи в несколько партиций вместе с коммитом offset, а read_committed прячет незавершённое от читателей. Для интеграции с внешним миром это ничего не даёт, там остаётся at-least-once плюс идемпотентность.»

Суть: segmentio/kafka-go — простой дефолт на чистом Go; sarama — старая и распространённая, но многословная; franz-go — самая полная по протоколу (транзакции, cooperative sticky, KIP-848); confluent-kafka-go — обёртка над librdkafka с cgo.

segmentio/kafka-go

Идиоматичный Go: Reader и Writer, всё через context, минимум церемоний. Нет транзакций, часть тонких настроек недоступна. Две грабли, которые стоит назвать вслух: ReadMessage коммитит offset сам (для at-least-once нужен FetchMessage плюс CommitMessages), а Writer с Async: true теряет сообщения при падении процесса и не возвращает ошибку доставки. Кроме того, по умолчанию балансировщик работает round-robin, а не по хешу, поэтому для совместимости с Java-клиентами нужен явный &kafka.Murmur2Balancer{}: напрашивающийся &kafka.Hash{} считает FNV-1a и совпадает с Sarama, а не с Java.

IBM/sarama

Исторически основная библиотека, покрывает почти весь протокол. Мешает многословный API с ConsumerGroupHandler и ручным управлением сессией: там легко неправильно обработать Setup/Cleanup и получить утечку горутин при ребалансировке. Обычно встречается в легаси-коде.

twmb/franz-go

Самая современная реализация на чистом Go: транзакции и EOS, cooperative sticky, поддержка свежих KIP, лучшие бенчмарки, аккуратная работа с батчами. Берут, когда нужны транзакции или максимальная производительность.

confluent-kafka-go

Обёртка над C-библиотекой librdkafka: эталонное поведение и все фичи Confluent, включая Schema Registry. За это платят cgo: сложнее кросс-компиляция и статическая сборка, тяжелее образ, хуже профилирование и трассировка.

Как отвечать, если опыта немного

Честно назвать одну библиотеку, с которой реально работал, и конкретную проблему, на которую наступил. «Использовали kafka-go; сначала читали через ReadMessage и удивлялись, почему при падении пода часть сообщений теряется — оказалось, он коммитит сам; перешли на FetchMessage с явным коммитом» работает лучше любого сравнительного обзора, потому что показывает практику, а не чтение README.

8.3RabbitMQ и альтернативы

RabbitMQ устроен зеркально Kafka: вместо «лога, который читают указателями» тут умный роутер с очередями, из которых сообщения исчезают после подтверждения. Отсюда гибкая маршрутизация, per-message ack, DLX и TTL, и отсюда же границы: историю не перечитать, одну очередь тяжело масштабировать.

Модель AMQP 0-9-1: кто здесь кто

RabbitMQ построен на открытом протоколе AMQP (Advanced Message Queuing Protocol) версии 0-9-1. Кроме передачи байтов по сети, протокол задаёт модель маршрутизации: как сообщение само находит дорогу к нужной очереди.

Эту модель проще всего понять на обычной почте. Письмо ты не несёшь адресату в квартиру, а сдаёшь в сортировочный центр, написав на конверте адрес. Дальше центр сверяется со своими таблицами и решает, в какие почтовые ящики положить копии. На языке AMQP сортировочный центр называется exchange (обменник), адрес на конверте служит routing key (ключом маршрутизации), а строчка в таблице «письма с таким адресом класть в такой-то ящик» зовётся binding (привязка). Почтовый ящик, где письмо лежит, пока его не заберут, и есть queue (очередь).

Аналогия держится почти до конца, а ломается в одном месте: сортировочный центр ничего не хранит и ничего не возвращает. Если ни одна строчка таблицы не подошла, письмо не уходит обратно отправителю и не ложится в «невостребованные», а просто исчезает. Об этом подробнее ниже.

На собесе первым делом проговори: продюсер в RabbitMQ никогда не публикует в очередь. Он публикует в exchange с каким-то routing key. Буфера в exchange нет, это чистая функция «сообщение + routing key + заголовки → список очередей». Правило, по которому очередь подписана на exchange, называется binding и само содержит binding key (или набор заголовков для headers).

Полный путь сообщения:

  • Connection — одно TCP-соединение (и TLS-сессия) до брокера. Дорогое, живёт долго.
  • Channel — логический подканал: внутри одного соединения их мультиплексируется много. Все операции (publish, consume, ack, qos, declare) идут через канал. Канал между горутинами не делят: в Go правило «один канал — одна горутина». amqp091-go прикрывает publish и ack мьютексом, но ответ синхронного вызова (declare, qos) на общем канале может достаться чужой горутине, а в других клиентах перемешанные фреймы рвут соединение с 505 UNEXPECTED_FRAME. Ошибка протокола закрывает канал (или всё соединение), и клиент должен уметь его пересоздать.
  • Exchange маршрутизирует. Бывает direct, topic, fanout, headers; отдельно стоит системный default exchange с пустым именем "". Он работает как direct, и к нему каждая очередь автоматически привязана по своему имени. Поэтому «наивный» пример publish("", "my-queue", body) и работает без объявления обменника: ты попал в default exchange с routing key, равным имени очереди.
  • Queue — единственное место, где сообщение реально лежит. Это FIFO-структура, и в classic-варианте она живёт на одном узле кластера, на том, где её объявили.
  • Consumer либо подписывается (basic.consume, push-модель), либо сам опрашивает очередь (basic.get, pull). Второе почти всегда антипаттерн: round-trip на каждое сообщение.

Из этой модели следует неприятная вещь: если сообщение не совпало ни с одной привязкой, оно молча исчезает. Не в DLX и не в «unroutable», а в никуда. Чтобы такие потери увидеть, публикуют с флагом mandatory=true и слушают basic.return (в Go через ch.NotifyReturn(...)) либо вешают на exchange alternate-exchange, куда стекается вся «непонятная почта».

Один и тот же publish, четыре типа exchange — куда попадёт сообщение Зелёная стрелка — привязка совпала, копия сообщения кладётся в очередь. Красный крест — не совпала, очередь не получает ничего. direct — binding key должен совпасть с routing key буквально, посимвольно routing key payment direct binding: payment q.pay.main binding: payment q.pay.audit — тоже копия binding: refund q.refunds — ничего не получит topic — ключ режется точками на слова; * = ровно одно слово, # = ноль или больше routing key order.created.eu topic binding: order.*.eu q.eu.orders binding: order.# q.all.orders binding: *.created q.created — два слова против трёх fanout — routing key не смотрится вообще, копия уходит во все привязанные очереди routing key любой fanout binding: ключ игнорируется q.cache.invalidate binding: ключ игнорируется q.search.index binding: ключ игнорируется q.audit.log headers — вместо ключа смотрят заголовки; x-match = all (все) или any (хотя бы один) headers: type=pdf, lang=ru headers all: type=pdf + lang=ru q.render.ru any: type=pdf q.pdf.any all: type=pdf + lang=en q.render.en — lang не тот
Четыре типа exchange на одном примере. Exchange берёт routing key (или заголовки) и по таблице привязок решает, в какие очереди положить копии. Одному консьюмеру две копии не придут: копия кладётся в очередь, и уже очередь отдаёт её одному из своих подписчиков.

Механика routing key на примерах

Чаще всех работает topic: почти в любом проекте хватает одного topic-обменника на домен. Ключ состоит из слов через точку, это иерархия длиной до 255 байт. Спецсимволов всего два: * заменяет ровно одно слово, #ноль или больше слов.

Binding keyorder.createdorder.created.euorder.paid.eu.b2buser.created
order.createdданетнетнет
order.*данетнетнет
order.#дададанет
*.createdданетнетда
#.eu.#нетдаданет
#дададада

Ключ проектируют от общего к частному и кладут в него только то, на что кто-то реально захочет подписаться: домен.сущность.событие.регион. Идентификаторы (order.created.42718) класть в ключ можно, но бессмысленно: на конкретный id всё равно никто не подпишется, а таблица привязок в брокере строится как дерево и растёт.

ТипПо чему маршрутизируетСтоимостьПод какую задачу
directТочное равенство routing key и binding keyХеш-таблица, O(1)Раздача задач по типу: email, sms, push. Приоритеты через отдельные очереди task.high / task.low
topicШаблон по словам ключа (*, #)Дерево префиксов, чуть дороже directШина доменных событий: один обменник, подписчики сами выбирают срез. Дефолтный выбор в 80 % случаев
fanoutНи по чему — всем привязаннымСамый дешёвыйШироковещание: инвалидация кэша на всех подах, «перечитай конфиг», прогрев
headersЗаголовки сообщения, x-match=all|anyСамый дорогой: перебор привязокКогда критериев несколько и они не складываются в иерархию (формат + язык + приоритет). На практике почти всегда заменяется на topic
"" (default)direct, где binding key = имя очередиO(1)Прямая отправка в конкретную очередь. Удобно для RPC-ответов и быстрых прототипов, но жёстко привязывает продюсера к имени очереди
Подвох: «а если два консьюмера на одной очереди — оба получат сообщение?»

Нет. Копии делаются по очередям, а не по консьюмерам. Exchange положит по копии в каждую подходящую очередь, а внутри одной очереди сообщение уйдёт ровно одному подписчику (round-robin с оглядкой на prefetch). Этим pub/sub и отличается от work queue в терминах AMQP: pub/sub = очередь на каждого подписчика, work queue = одна очередь на всех. Классическая ошибка: сделали одну очередь на пять сервисов и удивляются, что событие видит только один.

// Declare идемпотентен, звать его можно при каждом старте.
// Повторный declare с другими параметрами вернёт PRECONDITION_FAILED и закроет канал.
err := ch.ExchangeDeclare(
    "ex.orders", // имя
    "topic",     // тип
    true,        // durable: переживёт рестарт брокера
    false,       // autoDelete: не удалять при отключении последней привязки
    false,       // internal: можно публиковать снаружи
    false,       // noWait: ждём подтверждения от брокера
    nil,
)

q, err := ch.QueueDeclare(
    "q.billing.orders", // именованная очередь: имя = контракт владельца-консьюмера
    true,               // durable
    false,              // autoDelete
    false,              // exclusive: доступна и другим соединениям
    false, nil,
)

// У очереди сколько угодно привязок, и работают они как OR, а не AND.
_ = ch.QueueBind(q.Name, "order.created.#", "ex.orders", false, nil)
_ = ch.QueueBind(q.Name, "order.paid.#",    "ex.orders", false, nil)

Подтверждения на стороне консьюмера: ack, nack, reject

Когда брокер отдаёт сообщение консьюмеру, оно не удаляется из очереди, а переходит в состояние unacked (доставлено, но не подтверждено) и остаётся привязанным к каналу. Каждая доставка получает delivery tag, монотонный счётчик внутри канала. Это не идентификатор сообщения и не сквозной номер: после реконнекта нумерация начинается заново. Закрыть доставку консьюмер может тремя способами:

КомандаЧто делаетКогда применять
basic.ack(tag, multiple)Удаляет сообщение из очереди навсегда. multiple=true подтверждает все доставки с тегом ≤ указанного — батчевое подтверждениеОбработка успешно завершена и её эффект зафиксирован (транзакция БД закоммичена)
basic.nack(tag, multiple, requeue)Расширение RabbitMQ поверх AMQP: то же, что reject, но умеет multipleОшибка. requeue=true — вернуть в очередь, false — выкинуть или в DLX
basic.reject(tag, requeue)Стандарт AMQP, ровно одно сообщениеТо же самое, когда батч не нужен
ничего не делатьСообщение висит в unacked, пока жив каналНикогда. Это утечка: очередь не пустеет, prefetch забит, метрика messages_unacknowledged растёт

Если канал или соединение закрывается (падение пода, сетевой обрыв, ошибка протокола), все unacked-сообщения этого канала автоматически возвращаются в очередь и будут доставлены заново с флагом redelivered=true. Так в RabbitMQ и устроен at-least-once: таймаут видимости, как в SQS, не нужен, гарантию даёт сам разрыв TCP-соединения.

Где теряет сообщения автоack (autoAck=true / noAck)

С autoAck брокер считает сообщение доставленным в момент записи в сокет и сразу удаляет его из очереди. Он не ждёт ни обработки, ни даже того, что клиент реально прочитает байты из TCP-буфера. Значит, пропадает всё, что было:

  • в TCP-буфере ядра, когда под убили;
  • во внутренней очереди клиентской библиотеки (в Go это канал <-chan Delivery, куда библиотека заранее складывает доставки);
  • в обработке, если случилась паника, OOM или SIGKILL при деплое.

Вторая беда: autoAck отключает prefetch. Неподтверждённых сообщений не существует, ограничивать нечего, поэтому брокер льёт в консьюмера всё, что есть, на максимальной скорости, и клиент съедает память. Автоack допустим ровно там, где допустима потеря: метрики, телеметрия, живые логи.

Когда подтверждать

Ack делают последним, а не первым. Сначала транзакция БД закоммичена, файл записан и fsync-нут, письмо ушло, и только потом ack. Сделаешь наоборот и получишь at-most-once с честной потерей данных. У долгой обработки опасность обратная — consumer_timeout (по умолчанию 30 минут, с RabbitMQ 3.8.17): не подтвердил за это время, и брокер принудительно закрывает канал с PRECONDITION_FAILED, а работа уходит на переобработку. Длинные задачи либо режут на шаги, либо поднимают таймаут в конфиге брокера.

Prefetch (basic.qos): сколько сообщений можно держать в руках

Prefetch («предвыборка») задаёт, сколько сообщений брокер может выдать одному консьюмеру авансом, не дожидаясь подтверждений по уже выданным. Это обратное давление из главы 8.1, сжатое в одно число: официант несёт не больше трёх тарелок, и четвёртую ему не дадут, пока он не поставит хотя бы одну.

basic.qos(prefetchSize, prefetchCount, global) задаёт максимум неподтверждённых доставок. Когда их набралось prefetchCount, этот потребитель ничего нового не получит: следующее сообщение уйдёт кому-то другому. prefetchSize (лимит в байтах) на практике не используют. global=false означает «лимит на каждого консьюмера канала», global=true — «лимит на весь канал суммарно».

По умолчанию prefetch не ограничен, и хуже дефолта не придумать: первый подключившийся консьюмер выгребает всю очередь себе в память, остальные простаивают.

basic.qos не вызван — prefetch не ограничен брокер отдаёт всё, что есть, первому же готовому каналу очередь пусто всё уже роздано консьюмер A 10 unacked в памяти консьюмер B простаивает 0 работы Падение A вернёт в очередь все 10 сразу; второй под куплен зря, а память A растёт с размером очереди. basic.qos(prefetchCount = 2) брокер держит не больше 2 неподтверждённых на консьюмера очередь 6 ждут своей очереди консьюмер A 2 в работе консьюмер B 2 в работе Работа делится поровну, память консьюмера ограничена, падение возвращает в очередь максимум 2 сообщения.
Нагрузку на консьюмера в RabbitMQ регулирует только prefetch. Без него распределение «честное» лишь формально: брокер раздаёт по кругу, но мгновенно, и вся очередь оседает в памяти первого подключившегося.

Как подобрать значение? Каждое подтверждение стоит сетевого round-trip, поэтому prefetch=1 даёт идеально ровное распределение и низкий throughput: консьюмер простаивает, пока ack идёт туда, а следующее сообщение обратно. Ориентиры такие:

Профиль задачиprefetchПочему
Долгая обработка (секунды и минуты): конвертация видео, отчёт, вызов внешнего API1Round-trip в единицы миллисекунд на фоне секунд не виден, зато задача не «залипнет» за спиной другой
Быстрая обработка (единицы миллисекунд), много мелких сообщений100–300Иначе сеть станет узким местом раньше, чем CPU
Пул из N воркеров внутри одного процессаN или 2NДержать все воркеры занятыми и один-два в очереди на вход, не больше
Неизвестно10–50Безопасный старт: и память ограничена, и сеть не пилит round-trip на каждое
Head-of-line blocking при большом prefetch

Память тут не единственная цена. Сообщения из буфера консьюмера уже не будут перераспределены: их не увидит новый под, который ты только что поднял под пик. При prefetch = 5000 и 10 секундах на сообщение последнее из буфера подождёт почти 14 часов, а метрика длины очереди всё это время показывает ноль. Поэтому при неоднородном времени обработки prefetch держат маленьким, а не «побольше, чтобы быстрее».

Requeue и бесконечный цикл переобработки

nack(requeue=true) возвращает сообщение в очередь. Вопрос в том, куда: RabbitMQ старается поставить его как можно ближе к исходной позиции, то есть практически в голову. Если консьюмер один, а ошибка постоянная (невалидный JSON, отсутствующая запись, битая схема), выходит цикл: доставка → ошибка → requeue → доставка через микросекунду. Ретраем такое не назовёшь: получается горячий цикл на 100 % CPU, который заодно держит за собой всю очередь.

// Плохо: вечный цикл на постоянной ошибке
msgs, _ := ch.Consume(q, "", false, false, false, false, nil)
for d := range msgs {
    if err := handle(d.Body); err != nil {
        d.Nack(false, true) // requeue = true
        continue            // и через мгновение снова сюда
    }
    d.Ack(false)
}
// Хорошо: разделяем транзиентные и постоянные ошибки
for d := range msgs {
    err := handle(d.Body)
    switch {
    case err == nil:
        d.Ack(false)
    case errors.Is(err, ErrPermanent), attempts(d) >= maxAttempts:
        // DLX у очереди один, и он ведёт в retry. В DLQ кладём сами.
        if publishToDLQ(ctx, ch, d) == nil {
            d.Ack(false)
        } // не вышло — не подтверждаем, сообщение вернётся
    default:
        d.Nack(false, false) // через DLX в retry-очередь с TTL, см. ниже
    }
}

Ещё одна неприятность: в AMQP нет встроенного счётчика попыток. Булево поле redelivered скажет «это не первая доставка», но номер попытки не назовёт. Считать попытки можно тремя способами:

  • Через x-death: при каждом проходе через dead-lettering брокер сам наращивает count в заголовке x-death. Самый честный вариант, руками ничего писать не надо, но работает он, только если ретраи идут через DLX.
  • Своим заголовком: консьюмер читает x-retry-count, увеличивает, публикует сообщение заново и подтверждает исходное. Работает всегда, но публикация + ack не атомарны: упал между ними и потерял сообщение (или продублировал, если порядок обратный).
  • Внешним счётчиком в Redis по message_id. Он нужен, когда попытки считают сквозь несколько разных очередей и сервисов.

DLX — dead letter exchange

Термин dead letter пришёл с почты: так называют «мёртвое письмо», которое не удалось ни доставить, ни вернуть отправителю. Такие письма уходят в отдельный отдел, где их разбирают вручную. Брокеры взяли и слово, и идею: необработанное сообщение не выбрасывают, а откладывают в сторону, чтобы человек потом посмотрел, что с ним не так. Отсюда и DLQ (dead letter queue), «очередь битых сообщений».

DLX (dead letter exchange) не «специальный тип очереди», а обычный exchange: в него очередь пересылает сообщения, которые из неё выбыли не по-хорошему. Задают его двумя аргументами при объявлении очереди, а в проде правильнее через policy, чтобы не передекларировать очередь ради смены маршрута:

args := amqp.Table{
    "x-dead-letter-exchange":    "dlx.orders",   // куда переслать
    "x-dead-letter-routing-key": "orders.dead",  // с каким ключом; без него сохранится исходный
}
ch.QueueDeclare("q.orders", true, false, false, false, args)

Причин уйти в DLX ровно четыре. На собесе назови все: кандидаты обычно помнят только первую.

reason в x-deathЧто произошлоТипичная причина в коде
rejectedbasic.reject или basic.nack с requeue=falseКонсьюмер решил, что ретрай не поможет
expiredИстёк TTL — сообщения (expiration) или очереди (x-message-ttl)Задержка в retry-очереди или протухшие данные
maxlenОчередь упёрлась в x-max-length / x-max-length-bytesПерегрузка; при overflow=drop-head вытесняется самое старое
delivery_limitQuorum-очередь: превышен x-delivery-limit перед доставкойPoison message, брокер сам отправляет его в DLX

При каждом прохождении брокер дописывает запись в заголовок x-death. Там лежит массив структур с полями queue, reason, count, exchange, routing-keys, time. Другого штатного счётчика попыток в AMQP нет: x-death[0].count честно говорит, сколько раз сообщение уже возвращалось через dead-lettering.

Три способа потерять сообщение на DLX
  • DLX не существует или к нему не привязана ни одна очередь. Тогда сообщение молча отбрасывается, как обычный unroutable publish, а консьюмер никакой ошибки не увидит: он уже сделал nack.
  • Цикл: DLX очереди A ведёт обратно в A (или в B, чей DLX ведёт в A). Сообщение крутится по кругу, x-death.count растёт, CPU горит. Сам RabbitMQ обрывает цикл в одном случае: сообщение вернулось в очередь, которую уже проходило, и при этом истёк TTL. На rejected защиты нет.
  • Никто не читает DLQ. Формально сообщение цело, а на деле про него вспомнят через месяц, когда кто-то спросит, почему не прошёл платёж. Нужны метрика rabbitmq_queue_messages{queue=~".*dlq"} (из per-object метрик) и алерт «в DLQ появилось хоть что-то».

Отложенный ретрай: лестница из TTL + DLX

В RabbitMQ нет встроенной «доставки через 30 секунд». Стандартный трюк: завести очередь без консьюмеров, у которой есть x-message-ttl и DLX обратно на рабочий обменник. Сообщение полежит там положенное время, «протухнет» и вернётся в работу.

Retry-лестница: nack(requeue=false) → wait-очередь с TTL → DLX → обратно в работу Ожидание создаётся тем, что у wait-очереди нет консьюмеров: единственный способ покинуть её — истечь по TTL. x.orders topic exchange q.orders рабочая очередь консьюмер ack — сообщение исчезло nack(requeue=false), x-death.count = 0 q.retry.5s ttl 5s · нет консьюмеров истёк TTL reason = expired dlx.retry rk = orders.retry возврат в x.orders, попытка №2 count >= 3 → сразу в DLQ q.retry.30s ttl 30s · попытка №3 q.orders.dlq разбор руками, алерт Число ступеней = числу попыток: 5s → 30s → 5m. Экспоненту делают отдельными очередями, потому что per-message TTL в classic-очереди работает только для головы очереди — сообщение с TTL 5s не обгонит стоящее перед ним с TTL 5m.
Задержку в RabbitMQ покупают через TTL, таймера здесь нет. Ступени делают отдельными очередями, потому что per-message TTL очередь не переупорядочивает: сообщения истекают строго по порядку, а не по величине своего TTL.
На head-of-line у TTL новички спотыкаются чаще всего

Хочется сделать «одну retry-очередь и класть в каждое сообщение свой expiration». Получится система, где сообщение с задержкой 1 секунда, попавшее за сообщение с задержкой 10 минут, выйдет через 10 минут. Classic-очередь проверяет TTL только у головы: она FIFO, и «протухшее в середине» физически не может выйти раньше. Решений два: лестница из очередей с фиксированным TTL (по одной на ступень) или плагин rabbitmq_delayed_message_exchange, который держит отложенные сообщения в Mnesia и публикует по таймеру.

Durable, persistent, confirms: четыре независимых условия

«Сообщение переживёт рестарт брокера» требует не одной галочки, а конъюнкции четырёх условий, и каждое отвечает за свой кусок пути. Провалить достаточно одно.

Что настраиваемГдеЧто будет, если забыть
durable=true у exchangeExchangeDeclareПосле рестарта обменника нет; публикация в него закрывает канал с NOT_FOUND
durable=true у очередиQueueDeclareОпределение очереди не восстановится; очередь исчезает вместе со всем содержимым
DeliveryMode = 2 (persistent)у каждого сообщенияТело живёт только в памяти. Очередь durable — но пустая после рестарта
Publisher confirmsch.Confirm() + NotifyPublishПродюсер не знает, дошло ли. Publish без confirms — это запись в TCP-сокет, а не гарантия

Комбинации ведут себя ровно по таблице, без сюрпризов: transient-сообщение в durable-очереди пропадёт при рестарте (при нехватке памяти его могут выгрузить на диск, но это экономия памяти, а не долговечность), persistent-сообщение в non-durable-очереди пропадёт вместе с очередью.

// Полный «надёжный» publish
if err := ch.Confirm(false); err != nil { return err }          // включаем confirms на канале
confirms := ch.NotifyPublish(make(chan amqp.Confirmation, 256)) // окно на 256 неподтверждённых
returns  := ch.NotifyReturn(make(chan amqp.Return, 8))          // сюда придут unroutable

err := ch.PublishWithContext(ctx, "x.orders", "order.created",
    true,  // mandatory: не смаршрутизировалось — верни, а не выбрось
    false, // immediate: RabbitMQ не поддерживает с 3.0, всегда false
    amqp.Publishing{
        DeliveryMode: amqp.Persistent, // = 2, иначе тело только в RAM
        MessageId:    evt.ID,          // ключ дедупликации для консьюмера
        Timestamp:    time.Now(),
        ContentType:  "application/json",
        Body:         body,
    })

c := <-confirms
if !c.Ack {
    // брокер не принял: публикуем заново, дубликаты тут норма
}
Что именно означает confirm

Для classic-очереди брокер отправляет basic.ack продюсеру, когда сообщение принято всеми очередями, в которые оно смаршрутизировалось, а если оно persistent, то и сброшено на диск. Для quorum-очереди ack приходит, когда запись реплицирована на большинство узлов Raft-группы. Поэтому «включили persistent, значит не потеряем» без confirms не работает: между вызовом Publish и fsync есть окно в сотни миллисекунд. Упал брокер в этом окне, и сообщения нет, а продюсер об этом не узнает.

Если ждать каждый confirm синхронно, платишь round-trip за сообщение, и throughput не поднимется выше единиц тысяч в секунду. Правильно держать окно: публикуем пачкой, ведём map deliveryTag → сообщение, вычёркиваем по приходящим Ack (не забывая про флаг multiple, он подтверждает всё до тега включительно), переотправляем по Nack и по таймауту.

Classic, quorum, streams

  • Classic queue, исходный тип, живёт целиком на одном узле. Реплицировать её умело только устаревшее зеркалирование (ha-mode): его объявили deprecated и выпилили в 4.0. Она дешёвая по ресурсам и годится для некритичных задач и очередей, которые не жалко потерять.
  • Quorum queue — реплицированная очередь на Raft, современный дефолт для всего, что нельзя терять. Всегда durable, подтверждение = запись на большинство узлов, есть штатная обработка poison message через x-delivery-limit. Платить приходится диском и памятью (Raft-лог держится до snapshot), отсутствием transient-режима, exclusive-очередей и global-prefetch.
  • Stream (3.9+) — append-only лог с офсетами: несколько независимых читателей, перечитывание с любой позиции, ретеншен по времени и размеру. По сути так RabbitMQ отвечает Kafka, когда нужно не «раздать задачи», а «многие читают одну ленту».

Kafka vs RabbitMQ: где проходит граница

Хуже всего ответить на собесе «Kafka быстрее». Они различаются не скоростью, а моделью хранения, и из неё выводится всё остальное. Kafka хранит упорядоченный лог, из которого чтение ничего не удаляет: позиция читателя остаётся его собственным числом. RabbitMQ держит набор очередей, из которых сообщение физически исчезает после ack, и мощный маршрутизатор перед ними.

Один и тот же вопрос «добавим ещё консьюмеров» — два разных ответа Kafka: параллелизм = число партиций, порядок по ключу сохраняется p0 p1 p2 consumer 1 consumer 2 consumer 3 consumer 4 Четвёртый не получит ничего: партиций три. Чтобы вырасти — увеличить партиции, а это меняет hash(key) % N и ломает порядок по ключу на границе изменения. Уменьшить нельзя вообще. Зато порядок внутри ключа гарантирован всегда. RabbitMQ: параллелизм не ограничен, порядок теряется сразу q.tasks одна очередь worker 1 worker 2 worker 3 worker 40 Сорок воркеров — просто сорок подов, ничего не надо перенастраивать: брокер раздаёт по кругу с оглядкой на prefetch. Но сообщения обрабатываются параллельно, поэтому сквозного порядка нет уже при двух консьюмерах.
Kafka покупает порядок ценой потолка параллелизма; RabbitMQ покупает неограниченный параллелизм ценой порядка. На этой развилке и делают выбор, остальное из неё следует.
СвойствоKafkaRabbitMQ
МодельРаспределённый лог с офсетами, читатель хранит позициюОчереди + маршрутизатор, сообщение исчезает после ack
Перечитать историюДа, seek на любой офсет, пока жив retentionНет (кроме streams). Прочитано — значит удалено
Много независимых потребителей одного потокаЕстественно: N consumer groups, каждая со своим офсетомN очередей, привязанных к одному exchange; каждая копия хранится отдельно
МаршрутизацияПримитивная: топик и hash(key). Фильтрация — на стороне консьюмера или через отдельный потокБогатейшая: direct / topic-паттерны / fanout / headers, alternate exchange, consistent-hash плагин
Гранулярность подтвержденияОфсет = «прочитано всё до». Одно «плохое» сообщение нельзя пропустить, не сдвинув курсорКаждое сообщение отдельно: ack / nack / reject, любой порядок
ПорядокСтрогий внутри партиции; ключ фиксирует партициюFIFO у очереди, но параллельные консьюмеры и requeue его разрушают
Потолок параллелизмаЧисло партиций в топикеПрактически нет; упирается в производительность самой очереди
Отложенная доставка, приоритетыНет из коробки (делают retry-топиками и внешним планировщиком)TTL+DLX, плагин delayed exchange, x-max-priority
Пропускная способностьСотни тысяч и миллионы сообщений в секунду, батчи, zero-copyДесятки тысяч в секунду на очередь; узкое место — Erlang-процесс очереди
Задержка одного сообщенияЕдиницы–десятки мс (батчинг, linger.ms)Субмиллисекундная при малой нагрузке
Размер сообщенияmax.message.bytes ~1 МБ по умолчанию; для большего — claim checkПо умолчанию до 16 МБ (max_message_size, можно поднять до 512 МБ), практически — держат мелкими: очередь живёт в памяти
ЭксплуатацияТяжелее: KRaft (ZooKeeper — до 4.0), ISR, ребалансировки, партиционированиеЛегче стартовать, сложнее масштабировать одну очередь

Как формулировать выбор

  • Kafka, если поток служит лентой фактов, которую читают несколько независимых систем и хотят перечитать: событийная интеграция между доменами, аналитика и стриминг, CDC из БД, аудит, материализация проекций, метрики. Признак: «нужен replay» или «потребителей будет больше, чем сейчас».
  • RabbitMQ, если по потоку идут задачи: «сконвертируй видео», «отправь письмо», «сходи в 1С». Признаки: нужна сложная маршрутизация по атрибутам, per-message ретраи с задержкой, приоритеты, RPC-стиль request/reply, много мелких разнородных очередей, команда не хочет содержать Kafka.
  • Оба сразу тоже нормальная архитектура, никакого хаоса в этом нет: события домена в Kafka, фоновые задачи в RabbitMQ. Ответ «взяли бы Kafka, потому что она модная» идёт в минус, ответ «зависит от того, лента это или задачи» в плюс.

NATS и JetStream

Core NATS — предельно простой брокер на Go: один статический бинарник, десятки мегабайт памяти, задержки в десятки микросекунд. Модель называется subject-based messaging: публикуют в тему вида orders.eu.created, подписываются по шаблону с * (ровно одно слово) и > (хвост целиком). И в отличие от остальных, Core NATS не хранит ничего. Если в момент публикации подписчика нет, сообщение просто исчезает: at-most-once, никаких ack. Зато из коробки есть то, чего нет у соседей: queue groups (подписчики с одинаковым именем группы делят поток, получаются competing consumers без объявления очередей) и встроенный request/reply с автоматическим inbox-subject для ответа.

JetStream добавляет поверх слой персистентности и включается флагом. С ним появляются:

  • Stream — лог, захватывающий один или несколько subjects, с политикой хранения: limits (как Kafka: держим, пока влезает по max_msgs, max_bytes, max_age), interest (удаляем, когда все подписчики подтвердили) и workqueue (как очередь: сообщение удаляется, как только его подтвердил один консьюмер). Хранить можно в file или memory, реплицирует Raft (replicas: 3).
  • Consumer — именованная позиция в стриме, durable или ephemeral, push или pull, с AckPolicy (none/all/explicit), AckWait, MaxDeliver, BackOff (готовая лестница задержек ретрая, которую в RabbitMQ приходится собирать из очередей) и FilterSubject. MaxAckPending работает как prefetch.
  • Дедупликация из коробки: заголовок Nats-Msg-Id и duplicate_window на стриме. Брокер сам отбросит повторную публикацию с тем же идентификатором внутри окна. В Kafka ради этого включают идемпотентного продюсера, в RabbitMQ пишут руками.
  • KV и Object store надстроены над стримом: key-value с watch на изменения и хранилище блобов. Часто они закрывают потребность, ради которой иначе тянули бы ещё и Redis.
Одной фразой на собесе

«NATS — это шина с микросекундными задержками и почти нулевой эксплуатацией; JetStream добавляет к ней persistence и умеет вести себя и как лог (limits), и как очередь задач (workqueue). Берут, когда нужен лёгкий нервный узел для микросервисов, request/reply и очень много subject-ов; не берут, когда нужна экосистема Kafka: Connect, Streams, устоявшиеся коннекторы и опыт команды».

Вопросы

5
Суть: продюсер публикует не в очередь, а в exchange с routing key. Exchange ничего не хранит — это чистая функция «ключ + заголовки → список очередей», а правило подписки называется binding. Отсюда и гибкость RabbitMQ, и его главная ловушка: не совпало ни с одной привязкой — сообщение исчезло молча.

Пять сущностей

  • Connection — TCP+TLS-соединение до брокера, дорогое, открывается один раз на процесс.
  • Channel — логический подканал внутри соединения. Все команды идут через него. Не потокобезопасен: в Go заводят по каналу на горутину. Ошибка протокола закрывает канал, и приложение должно уметь его пересоздать.
  • Exchange — маршрутизатор без состояния. Четыре типа плюс безымянный default exchange (""), к которому каждая очередь автоматически привязана по своему имени. Поэтому «hello world» с publish("", "my-queue", body) и работает без объявления обменника.
  • Queue — единственное место, где сообщение реально лежит. FIFO, у classic-типа живёт на одном узле кластера.
  • Binding связывает exchange → queue через binding key (или набор заголовков). Между одной парой может быть несколько привязок с разными ключами.

Типы exchange

ТипПравилоКогда берут
directbinding key совпадает с routing key посимвольноАдресная доставка: очередь на задачу, приоритеты через разные ключи
topicКлюч режется точками на слова; * — ровно одно слово, # — ноль или большеСобытийная маршрутизация: order.*.eu, audit.#. Дефолтный выбор в 90 % случаев
fanoutКлюч игнорируется, копия во все привязанные очередиBroadcast: инвалидация кэшей, конфиг-апдейты
headersМатчинг по заголовкам, x-match: any|allРедко: когда критериев несколько и они не укладываются в одну строку. Медленнее topic

Проговори отдельно: fanout означает не «одно сообщение всем», а «копию в каждую очередь». Если три сервиса подписаны на fanout через одну общую очередь, выйдет не рассылка, а конкуренция: сообщение достанется ровно одному. Веерная рассылка на N подписчиков = N очередей.

Топология в проде

Объявления идемпотентны, но с проверкой: повторный QueueDeclare с другими аргументами не «обновит» очередь, а закроет канал с PRECONDITION_FAILED. Поэтому то, что меняется (DLX, TTL, лимиты длины), в проде задают через policy: политики применяются к очередям по regexp и меняются без передекларации:

rabbitmqctl set_policy dlx-all "^q\\.orders\\." \
  '{"dead-letter-exchange":"dlx.orders","message-ttl":600000,"max-length":100000}' \
  --apply-to queues
Куда девается несмаршрутизированное сообщение

В никуда. Не в DLX (DLX работает на выходе из очереди, а в очередь сообщение не попало) и не в «unroutable»: оно просто отбрасывается, а Publish вернёт nil. Ловят двумя способами: mandatory=true плюс ch.NotifyReturn(...), чтобы брокер вернул сообщение продюсеру, или alternate-exchange на обменнике, через который «непонятная почта» стекается в отдельную очередь. Метрика для алерта: rabbitmq_channel_messages_unroutable_dropped_total.

Суть: четыре механизма, которые вместе образуют жизненный цикл сообщения у консьюмера. Ack удаляет, nack решает судьбу, prefetch ограничивает, сколько можно держать в руках, DLX ловит всё, что выбыло не по-хорошему. Главные ошибки: autoAck=true, requeue=true на постоянной ошибке и отсутствие ограничения на число попыток.

Ack, nack, reject

  • basic.ack — «обработано, удаляй». Только после этого сообщение исчезает из очереди. До ack оно в состоянии unacked: невидимо для других консьюмеров, но не удалено.
  • basic.nack(multiple, requeue) — «не смог». С requeue=true вернётся в очередь, с requeue=false уйдёт в DLX (или пропадёт, если DLX не настроен). multiple=true относится ко всем тегам до указанного включительно.
  • basic.reject появился раньше nack и умеет то же самое, но только для одного сообщения.
  • Разрыв канала автоматически возвращает все unacked сообщения этого канала в очередь. На этом и держится базовая гарантия at-least-once: упавший консьюмер ничего не теряет, но всё повторит.
autoAck=true теряет сообщения тихо

ch.Consume(q, "", true, ...) означает «считай доставленным в момент отправки». Брокер удаляет сообщение, ещё не зная, дошло ли оно по сети, и уж точно не зная, обработалось ли. Вдобавок autoAck отключает prefetch, и брокер зальёт в консьюмера всю очередь до OOM. Законный сценарий один: поток, который не жалко (метрики, телеметрия, логи).

Prefetch

basic.qos(prefetchCount=N) работает как кредитное окно: «не отправляй больше N неподтверждённых этому консьюмеру». Без него брокер раздаёт по кругу, но мгновенно, и вся очередь оседает в памяти первого подключившегося консьюмера, а второй под простаивает. Ориентиры: долгая обработка — 1, быстрые мелкие сообщения — 100–300, пул из N воркеров — N или 2N, «не знаю» — 10–50. У большого prefetch есть обратная сторона — head-of-line blocking: набранные в буфер сообщения уже не перераспределятся на новый под, а метрика длины очереди при этом покажет ноль.

Requeue

nack(requeue=true) кладёт сообщение обратно примерно на исходную позицию, то есть почти в голову. При постоянной ошибке и одном консьюмере выходит горячий цикл на 100 % CPU, а не ретрай, и очередь стоит. В AMQP нет счётчика попыток: поле redelivered булево. Считать попытки можно через x-death[0].count (заполняет сам брокер при dead-lettering), своим заголовком с перепубликацией или внешним счётчиком в Redis по message_id.

DLX

Обычный exchange, указанный очереди через x-dead-letter-exchange. Сообщение попадает туда по четырём причинам: rejected (nack/reject с requeue=false), expired (истёк TTL), maxlen (переполнение x-max-length), delivery_limit (quorum-очередь исчерпала x-delivery-limit). Через связку «wait-очередь с TTL и без консьюмеров + DLX обратно» строят отложенные ретраи; ступени делают отдельными очередями, потому что TTL в classic-очереди проверяется только у головы.

Эталонная схема обработки ошибок

Ошибки делят на транзиентные, где ретрай поможет (сеть, таймаут, дедлок БД), и постоянные, где он не поможет никогда (невалидный JSON, нет обязательного поля, бизнес-правило нарушено). Постоянные отправляют сразу в DLQ, транзиентные в retry-лестницу 5s → 30s → 5m, а после исчерпания попыток тоже в DLQ. И обязательно ставят алерт на непустую DLQ: очередь, в которую никто не смотрит, равна /dev/null с задержкой.

Суть: это разные флаги на разных уровнях, и нужны все четыре: durable exchange, durable queue, delivery_mode=2 у сообщения и publisher confirms у продюсера. Забыть одно — потерять всё; классический вопрос-ловушка «сделали очередь durable, почему сообщения пропали после рестарта».

Кто за что отвечает

  • durable exchange сохраняет определение обменника после рестарта. Без него публикация после перезапуска закроет канал с NOT_FOUND.
  • durable queue переносит через рестарт определение очереди, а не сообщения: очередь после рестарта на месте, но окажется пустой, если сообщения были transient.
  • persistent message (DeliveryMode = 2) — брокер обязан записать тело на диск. Transient-сообщение живёт в памяти. При нехватке RAM его могут выгрузить на диск, но это экономия памяти, а не долговечность, и рестарт его всё равно убьёт.
  • С publisher confirms продюсер узнаёт, что брокер принял сообщение. Без них Publish просто пишет в TCP-сокет: ошибки не будет, гарантии тоже.

Что значит confirm

Брокер шлёт basic.ack продюсеру, когда сообщение принято всеми очередями, куда оно смаршрутизировалось, и (для persistent) сброшено на диск; для quorum-очереди ack приходит после репликации на большинство узлов Raft-группы. Между вызовом Publish и fsync есть окно в сотни миллисекунд, и падение брокера в нём съедает сообщение. Синхронно ждать каждый confirm нельзя (round-trip на сообщение), поэтому держат окно неподтверждённых: map deliveryTag → message, вычёркиваем по Ack (не забыв про флаг multiple: он подтверждает всё до тега включительно), переотправляем по Nack и по таймауту.

Цена и что выбирают сегодня

Persistent + confirms снижают пропускную способность в разы: каждый publish упирается в диск. Поэтому долговечность включают выборочно: платежи и заказы — persistent и confirms, телеметрия и «пользователь открыл экран» — transient и fire-and-forget. И помни, что classic mirrored queues (ha-mode) устарели и удалены в RabbitMQ 4.0. Реплицируют сегодня через quorum queues на Raft, где ack означает запись на большинство узлов. Они всегда durable, умеют x-delivery-limit против poison message, но не поддерживают transient-режим, exclusive-очереди и global prefetch, а Raft-лог занимает диск до snapshot.

Что ответить на «а это гарантия от потери?»

Нет, гарантия тут одна: «брокер подтвердил, значит записал». Потерять сообщение по-прежнему можно на трёх других участках: продюсер упал до publish (лечится transactional outbox), консьюмер сделал ack до фактической обработки (лечится ack после коммита работы), весь узел с очередью потерян безвозвратно (лечится quorum-очередями и репликацией). Хороший ответ проходит по всей цепочке и не останавливается на флажке durable.

Суть: не «что быстрее», а разные модели хранения. Kafka — упорядоченный лог, из которого ничего не удаляется по факту чтения; читатель хранит свою позицию. RabbitMQ — очереди, из которых сообщение исчезает после ack, плюс мощный маршрутизатор перед ними. Одной фразой: лента фактов против раздачи задач.

Что следует из модели

  • Replay. В Kafka историю перечитывают штатным seek на офсет (пересобрать проекцию, залить новый сервис историей). В RabbitMQ прочитанного больше нет: чтобы «перечитать», надо было заранее завести ещё одну очередь на тот же exchange.
  • Много потребителей. В Kafka N consumer groups читают один и тот же лог, данные хранятся один раз. В RabbitMQ у каждого подписчика своя очередь, то есть N копий сообщения в памяти и на диске.
  • Гранулярность. Kafka коммитит офсет, то есть «прочитано всё до». Пропустить одно «плохое» сообщение, не сдвинув курсор, нельзя, отсюда retry-топики. RabbitMQ подтверждает каждое сообщение отдельно, в любом порядке — отсюда per-message ретраи, DLX, приоритеты.
  • Параллелизм и порядок. В Kafka потолок консьюмеров в группе = число партиций, зато порядок внутри ключа гарантирован. В RabbitMQ консьюмеров может быть сколько угодно, но сквозного порядка нет уже при двух.
  • Маршрутизация. У Kafka её практически нет: топик и hash(key), фильтрует консьюмер. У RabbitMQ есть topic-паттерны, headers, alternate exchange, consistent-hash, и логика «кому это надо» живёт в брокере.

Развилка выбора

Требование в задачеОтвет
Нужен replay, аналитика, CDC, несколько независимых потребителей потокаKafka
Сотни тысяч сообщений в секунду, поток однородныйKafka
Строгий порядок по бизнес-ключу (по клиенту, по заказу)Kafka, ключ = этот идентификатор
Фоновые задачи, воркеры, «сделай и забудь»RabbitMQ
Ретраи с задержкой, приоритеты, TTL, сложная маршрутизация по атрибутамRabbitMQ
Request/reply поверх брокераRabbitMQ (или NATS)
Неравномерная длительность задач, надо докинуть 50 воркеров под пикRabbitMQ: в Kafka упрёшься в число партиций
Маленькая команда, нет SRE, «лишь бы работало»RabbitMQ или managed-Kafka
Хороший финал ответа

«Держать оба нормально: события домена в Kafka, потому что их читают несколько команд и иногда переигрывают; фоновые задачи в RabbitMQ, потому что там нужны приоритеты и отложенные ретраи. И граница размывается: у RabbitMQ появились streams с офсетами, у Kafka retry-топики закрывают часть сценариев очередей. Но модель хранения у них по-прежнему разная, и выбирать надо по ней».

Суть: Core NATS — сверхлёгкая шина без хранения (at-most-once, нет подписчика — сообщение исчезло), с subject-маршрутизацией, queue groups и встроенным request/reply. JetStream — слой персистентности поверх, который умеет вести себя и как лог (limits), и как очередь задач (workqueue).

Core NATS

Один статический бинарник на Go, десятки мегабайт памяти, задержки в десятки микросекунд, кластеризация через gossip и супер-кластеры между регионами. Маршрутизация строится на subjects вида orders.eu.created с шаблонами * (ровно одно слово) и > (весь хвост). Субъекты не объявляют заранее, их можно заводить по одному на сущность. Такого не выдержит ни Kafka (топик = каталоги и файлы), ни RabbitMQ (очередь = процесс Erlang).

У соседей из коробки нет двух вещей. Это queue groups, где подписчики с одинаковым именем группы делят поток (competing consumers без объявления очередей), и request/reply с автоматическим временным inbox-subject для ответа. Платить за это приходится отсутствием хранения: at-most-once, никаких ack, а без подписчика сообщение просто отбрасывается.

JetStream

  • Stream захватывает набор subject-ов и хранит их в file или memory storage с репликацией на Raft. Политика limits — как Kafka (держим, пока влезает по max_msgs/max_bytes/max_age), interest удаляет, когда все подписчики подтвердили, workqueue работает как очередь: сообщение удаляется после ack одного консьюмера.
  • Consumer — именованная позиция: durable или ephemeral, push или pull, AckPolicy (none/all/explicit), AckWait, MaxDeliver, готовая лестница BackOff для ретраев, FilterSubject. MaxAckPending здесь вместо prefetch.
  • Дедупликация публикации: заголовок Nats-Msg-Id плюс duplicate_window на стриме, и брокер сам отбросит повтор внутри окна. В Kafka ради этого включают идемпотентного продюсера, в RabbitMQ пишут руками.
  • KV и Object store поверх стрима: key-value с Watch на изменения и хранилище блобов. Нередко закрывают потребность, ради которой иначе завели бы Redis.

Как позиционировать

JetStream садится между Kafka и RabbitMQ: он даёт персистентность и ретраи, как RabbitMQ, и офсеты с перечитыванием, как Kafka, а эксплуатировать его заметно дешевле обоих. Берут, когда нужен лёгкий «нервный узел» для микросервисов, request/reply и много subject-ов, когда важен edge/IoT или мультирегион. Не берут, когда нужна экосистема Kafka (Connect, Streams, ksqlDB, готовые коннекторы) или когда команда уже умеет эксплуатировать Kafka и не хочет второй брокер.

Если опыта с NATS нет

Так и скажи, но всё равно назови разделение «core = at-most-once без хранения, JetStream = persistence с политиками retention» и одну отличительную черту NATS: дедупликацию по Nats-Msg-Id из коробки или queue groups без объявления очередей. Этого хватит, чтобы закрыть вопрос: проверяют не опыт, а понимание, зачем существует третий вариант.

8.4Паттерны асинхронной интеграции

Паттерны этой главы лечат одну и ту же беду: у нас две системы, а транзакция только в одной. База данных умеет атомарность, брокер умеет доставку, общего коммита между ними нет. Outbox переносит публикацию внутрь транзакции БД, inbox переносит дедупликацию внутрь транзакции консьюмера. Остальное касается того, что класть в сообщение и что делать, когда сообщений становится больше, чем сил их переварить.

Двойная запись — корень всех бед

Задача выглядит безобидно: сохранить заказ в PostgreSQL и опубликовать событие OrderCreated, чтобы отреагировали склад, биллинг и аналитика. Два действия, две разные системы, и ни одного способа сделать их атомарно. Какой порядок ни выбери, между ними остаётся окно, в котором процесс может умереть — под убьёт OOM-killer, нода уедет на обслуживание, сеть моргнёт.

Двойная запись: оба порядка сломаны, ломаются по-разному Красная молния — процесс умирает между двумя операциями. Это не редкость: это нормальный режим работы кластера. Порядок A: сначала COMMIT в БД, потом publish COMMIT заказ 42 сохранён crash publish не выполнился Заказ есть в базе, события нет никогда. Склад не резервирует товар, биллинг не выставляет счёт, аналитика не видит продажу. Тихая рассинхронизация. Порядок B: сначала publish, потом COMMIT publish событие ушло crash COMMIT откат по таймауту Событие есть, заказа нет. Консьюмеры пошли за деталями в API и получили 404 — или, хуже, зарезервировали товар под заказ-призрак. Отменить опубликованное нельзя. Почему «просто поретраить» не спасает Ретрай живёт в памяти того же процесса, который умер. Сдвинуть окно можно, убрать нельзя: какую бы пару операций вы ни взяли, между ними всегда есть момент «первая прошла, вторая нет». 2PC/XA формально решает, практически — нет: Kafka не поддерживает XA вовсе, координатор блокирующий, а его падение оставляет in-doubt транзакции, держащие блокировки в БД.
От порядка зависит лишь то, как испортятся данные, а не то, испортятся ли они. Начнёшь с БД, потеряешь события; начнёшь с брокера, получишь фантомные. Поэтому нужен паттерн, а не аккуратность.
Чем не решается
  • «Публиковать внутри транзакции». Брокер в транзакции БД не участвует: defer tx.Rollback() откатит строки, но уже отправленное сообщение не отзовёт.
  • «Ретраить публикацию с бэкоффом» уменьшает вероятность, но не убирает её. Ретрай живёт в памяти процесса и умирает вместе с ним.
  • «Проверять периодически и досылать» по сути тот же outbox, только без таблицы: придётся сканировать сам заказ и где-то хранить флаг «опубликован». Этот флаг и есть outbox, просто размазанный по бизнес-таблице.
  • 2PC / XA в теории подходит, а на практике Kafka не умеет XA, координатор становится единой точкой отказа, зависшие prepared-транзакции держат блокировки и не дают VACUUM убирать старые версии строк. В микросервисах его не применяют.

Transactional outbox

Outbox буквально значит «ящик исходящих», как в почтовом клиенте: письмо ты уже написал и нажал «отправить», но оно ещё не ушло. Зато оно точно существует и точно уйдёт — просто позже. Relay («релей», передатчик) — фоновый процесс, который разгребает этот ящик: берёт неотправленное и толкает в брокер. Вся конструкция называется transactional outbox — «транзакционный ящик исходящих»: запись в него идёт внутри той же транзакции БД, что и само бизнес-изменение.

Выход в том, чтобы сделать публикацию частью той же транзакции БД. Раз атомарности между базой и брокером не бывает, событие записывают в ту же базу, в таблицу outbox, тем же COMMIT. Дальше релей вычитывает эту таблицу и публикует в брокер. Атомарность обеспечивает СУБД, доставку берёт на себя релей, и окна между ними нет: закоммитилось либо всё, либо ничего.

CREATE TABLE outbox (
    id             bigserial PRIMARY KEY,
    aggregate_type text        NOT NULL,          -- 'order'  — отсюда имя топика
    aggregate_id   text        NOT NULL,          -- '42'     — станет ключом партиционирования
    event_type     text        NOT NULL,          -- 'OrderCreated'
    payload        jsonb       NOT NULL,          -- тело события
    headers        jsonb       NOT NULL DEFAULT '{}'::jsonb,  -- traceparent, schema_version
    created_at     timestamptz NOT NULL DEFAULT now(),
    published_at   timestamptz                    -- NULL = ещё не отправлено
);

-- частичный индекс маленький: в нём только неотправленные строки
CREATE INDEX outbox_unpublished_idx ON outbox (id) WHERE published_at IS NULL;
func (s *Service) CreateOrder(ctx context.Context, o Order) error {
    tx, err := s.db.BeginTx(ctx, nil)
    if err != nil { return err }
    defer tx.Rollback() //nolint:errcheck

    if _, err = tx.ExecContext(ctx,
        `INSERT INTO orders (id, user_id, total, status) VALUES ($1,$2,$3,'new')`,
        o.ID, o.UserID, o.Total); err != nil {
        return err
    }

    payload, _ := json.Marshal(OrderCreated{OrderID: o.ID, UserID: o.UserID, Total: o.Total})
    if _, err = tx.ExecContext(ctx,
        `INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload, headers)
         VALUES ('order', $1, 'OrderCreated', $2, $3)`,
        o.ID, payload, traceHeaders(ctx)); err != nil {
        return err
    }

    return tx.Commit() // ← всё становится реальным только здесь
}
Главное свойство

После COMMIT событие уже существует и будет опубликовано, даже если процесс умрёт на следующей инструкции: релей работает отдельно, он найдёт строку и отправит. До COMMIT события нет вообще, так что фантомам взяться неоткуда. Двойная запись превратилась в одинарную.

Relay через polling

Самый простой вариант: фоновая горутина или отдельный воркер раз в N миллисекунд забирает пачку неотправленных строк, публикует их и помечает. Нетривиальная деталь тут одна — FOR UPDATE SKIP LOCKED: с ней можно держать несколько экземпляров релея, и они не блокируют друг друга и не берут одни и те же строки.

const batch = `
  SELECT id, aggregate_type, aggregate_id, event_type, payload, headers
    FROM outbox
   WHERE published_at IS NULL
   ORDER BY id
   LIMIT 100
     FOR UPDATE SKIP LOCKED`      // ← конкурентные релеи разбирают разные строки

func (r *Relay) tick(ctx context.Context) error {
    tx, err := r.db.BeginTx(ctx, nil)
    if err != nil { return err }
    defer tx.Rollback() //nolint:errcheck

    rows, err := tx.QueryContext(ctx, batch)
    if err != nil { return err }
    msgs, ids := scanBatch(rows)

    for i, m := range msgs {
        // ключ = aggregate_id: события одного заказа лягут в одну партицию;
        // порядок между ними держится, пока релей один
        if err := r.producer.Publish(ctx, topicFor(m), m.AggregateID, m.Payload, m.Headers); err != nil {
            return err            // строки останутся неотмеченными, отправим в следующий тик
        }
        _ = i
    }

    _, err = tx.ExecContext(ctx,
        `UPDATE outbox SET published_at = now() WHERE id = ANY($1)`, pq.Array(ids))
    if err != nil { return err }
    return tx.Commit()
}
Почему нельзя вести курсор по «последнему обработанному id»

Соблазнительно вместо флага published_at хранить одно число last_id и делать WHERE id > last_id. Подход дырявый, и дырка коварная: bigserial выдаёт номера вне транзакции. Транзакция T1 взяла id = 5 и ещё думает, T2 взяла id = 6 и закоммитилась. Релей в этот момент видит 6, публикует и ставит last_id = 6. Через миллисекунду коммитится T1, но строку с id = 5 релей уже никогда не выберет. Событие пропало молча — ровно та беда, от которой мы уходили. Поэтому отметку ставят на строке (published_at или DELETE), а не в курсоре. По той же причине нельзя фильтровать по created_at > last_time.

Что ещё надо предусмотреть в polling-релее:

  • Уборка. Таблица растёт со скоростью потока событий. Чистят либо пачками через DELETE ... WHERE published_at < now() - interval '3 days' (не одним запросом на миллион строк — иначе выйдет долгая транзакция и раздутый WAL), либо DELETE сразу после публикации, либо секционированием по дню: старую секцию отцепляют (DETACH PARTITION) и удаляют через DROP TABLE. Без уборки деградируют индекс и autovacuum, а за ними и весь сервис, а не один релей.
  • Интервал опроса — компромисс между задержкой и нагрузкой. При 100 мс это 10 лишних запросов в секунду на пустой таблице (дёшево благодаря частичному индексу), при 5 с событие ждёт 5 с. Разумно взять короткий интервал плюс сигнал NOTIFY outbox после коммита, чтобы релей просыпался сразу.
  • Порядок. Один релей с ORDER BY id публикует по порядку, но у нескольких экземпляров со SKIP LOCKED порядка нет даже внутри агрегата: два релея могут взять события 4 и 5 одного заказа в разные пачки, и пятое уйдёт раньше. Ключ aggregate_id положит их в одну партицию, но в том порядке, в каком они пришли. Если порядок внутри агрегата важен, берут схему «один релей, лидер через advisory lock» или CDC.
  • Идемпотентность продюсера. С enable.idempotence уйдут дубликаты от ретраев внутри одной сессии продюсера, но не те, что появятся при перезапуске релея.

Relay через CDC (Debezium)

При CDC (Change Data Capture, «захват изменений данных») изменения в базе вычитывают не запросами к таблицам, а из журнала транзакций самой СУБД (WAL в PostgreSQL, binlog в MySQL). База и так пишет в этот журнал каждое изменение, чтобы восстановиться после сбоя. CDC-инструмент подключается к нему как ещё одна реплика и превращает каждую запись журнала в сообщение, не делая ни одного запроса к рабочей базе.

Второй вариант релея не спрашивает базу, а читает её журнал. Debezium подключается к PostgreSQL через логический слот репликации (к MySQL как реплика по binlog) и превращает каждое изменение строки в сообщение Kafka. Если направить его на таблицу outbox, получится релей без единого SELECT в рабочую базу.

PollingCDC / Debezium
Нагрузка на БДПостоянные запросы + UPDATE каждой строки (второй проход по heap, больше WAL)Чтение WAL, обычной нагрузки почти нет
ЗадержкаПоловина интервала опроса в среднемМиллисекунды
ПорядокВнутри одного релея; между несколькими — не гарантированСтрого порядок коммитов в WAL
ИнфраструктураНичего лишнего: горутина в сервисеKafka Connect кластер, конфиги коннекторов, мониторинг слотов
СхемаПолный контроль над телом сообщенияНужен Outbox Event Router SMT, чтобы вытащить payload и поставить key/topic
РискЗабыть уборку → распухшая таблицаСлот не читается → WAL не удаляется → диск БД кончается
Трюк «INSERT + DELETE в одной транзакции»

С CDC строку в outbox можно не хранить вовсе: вставить и тут же удалить в той же транзакции. В таблице всегда пусто (нечего чистить, индекс не растёт), а в WAL остаются обе записи, INSERT и DELETE — Debezium опубликует событие по INSERT, а DELETE отфильтрует. REPLICA IDENTITY FULL для этого не нужен: INSERT в WAL и так несёт строку целиком, а Outbox Event Router удаления пропускает. Красиво, но намертво привязывает тебя к CDC: включить polling «на всякий случай» уже не выйдет.

В эксплуатации CDC опаснее всего replication slot. Пока слот существует и его не читают (Connect упал, топик недоступен, коннектор на паузе), PostgreSQL обязан хранить весь WAL с последней подтверждённой позиции. Диск заполняется за часы, и падает уже не пайплайн событий, а сама база. Алерты тут обязательны: на pg_replication_slots.active и на лаг слота в байтах.

Transactional outbox: где событие рождается, кто его выносит и что бывает при падении сервис заказов BEGIN INSERT orders INSERT outbox COMMIT PostgreSQL — одна транзакция orders 42 · new · 1200 ₽ outbox id 7 · OrderCreated · 42 published_at = NULL атомарно с orders relay: polling SELECT ... published_at IS NULL FOR UPDATE SKIP LOCKED relay: CDC Debezium читает WAL через replication slot Kafka key = aggregate_id порядок внутри заказа Точки отказа — и что происходит в каждой ① Падение до COMMIT Откатывается всё: ни заказа, ни события. Клиент получил ошибку и повторит — это корректное состояние, а не порча данных. ② Падение после COMMIT, до publish Строка лежит с published_at = NULL. Релей (любой из экземпляров) подхватит её в следующем тике. Задержка, но не потеря. ③ Publish прошёл, UPDATE не успел Строка снова видна как неотправленная — и будет опубликована второй раз. Отсюда берутся дубликаты. Всегда. Итоговая гарантия outbox — at-least-once, и никакая аккуратность её не улучшит: между «брокер записал» и «мы это отметили» окно есть всегда, потому что это снова две системы. Дубликаты добавляют и ретраи продюсера, и перезапуск релея с середины пачки, и повторный запрос клиента. Поэтому outbox почти никогда не живёт один — на приёме его дополняет inbox или идемпотентный обработчик.
Outbox убирает потерю событий, но не дубликаты. Вместо «двух систем без общего коммита» получается «одна транзакция плюс надёжный вынос наружу» — а вынос наружу по-прежнему может повториться.

Inbox и дедупликация на приёме

Inbox («ящик входящих») — зеркальная пара к outbox. В outbox лежит то, что мы обязаны отправить, а в inbox — то, что мы уже приняли и обработали. Магии в нём нет: обычная таблица с ключами сообщений, по которой консьюмер проверяет, не видел ли он это сообщение раньше.

У консьюмера та же история, только зеркальная: надо сделать работу и запомнить, что сообщение обработано, — снова две операции. Если запоминать отдельным запросом после коммита работы, между ними опять появляется окно: упали, потеряли отметку, обработали второй раз. Решение симметрично outbox: отметку кладут в ту же транзакцию, что и бизнес-работу.

CREATE TABLE processed_messages (
    message_id  text        NOT NULL,
    consumer    text        NOT NULL,      -- имя consumer group: каждый обрабатывает своё
    processed_at timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer, message_id)     -- уникальность и есть весь механизм
);
-- Секционировать по processed_at нельзя: ключ обязан включать processed_at,
-- и дедупликация сломается. Старое удаляют пачками по индексу.
CREATE INDEX ON processed_messages (processed_at);
func (c *Consumer) Handle(ctx context.Context, msg Message) error {
    tx, err := c.db.BeginTx(ctx, nil)
    if err != nil { return err }
    defer tx.Rollback() //nolint:errcheck

    // 1) сначала отметка: INSERT блокирует ключ и сериализует
    //    конкурентные доставки одного сообщения
    res, err := tx.ExecContext(ctx,
        `INSERT INTO processed_messages (consumer, message_id) VALUES ($1, $2)
         ON CONFLICT DO NOTHING`, c.name, msg.ID)
    if err != nil { return err }
    if n, _ := res.RowsAffected(); n == 0 {
        return tx.Commit() // дубликат: уже обрабатывали, просто подтверждаем брокеру
    }

    // 2) бизнес-работа в той же транзакции
    if err := c.apply(ctx, tx, msg); err != nil {
        return err // rollback снимет и отметку: сообщение вернётся, обработаем заново
    }

    return tx.Commit() // 3) отметка и результат становятся видимыми вместе
}
Inbox: отметка и работа — одна транзакция, поэтому промежуточного состояния не бывает брокер msg id=A msg id=A повтор консьюмер: BEGIN … COMMIT INSERT processed_messages ON CONFLICT DO NOTHING UPDATE balance бизнес-работа rows = 1 → первое появление: работаем и коммитим rows = 0 → дубликат: работу пропускаем, брокеру шлём ack ошибка в работе → ROLLBACK снимает и отметку: сообщение вернётся COMMIT — отметка и результат становятся видимыми одновременно Так делать нельзя: две транзакции tx1: UPDATE balance … COMMIT ✗ crash здесь tx2: INSERT processed_messages … COMMIT Отметки нет, деньги списаны. Повтор спишет ещё раз. Так правильно: одна транзакция BEGIN INSERT processed_messages → UPDATE balance COMMIT (ack брокеру — только после) Промежуточного состояния не существует в принципе.
Inbox устроен как outbox наоборот. Там мы вносили в транзакцию публикацию, здесь вносим отметку об обработке. Оба паттерна держатся на одном: атомарность есть только у базы, значит всё критичное должно оказаться в ней.

Детали, на которых валятся

  • Ключ дедупликации. Лучше всего подходит message_id, который продюсер сгенерировал один раз и переиспользует при всех ретраях (в outbox это естественно: id строки). Годится и бизнес-ключ вида order_id + event_type + version. Хеш тела не годится: ретрай с обновлённым timestamp даст другой хеш и проедет как новое сообщение, а два легитимных одинаковых события (два одинаковых списания по 100 ₽), наоборот, схлопнутся в одно.
  • Область уникальности задаёт пара «консьюмер + message_id», а не один message_id. Иначе первая же группа «съест» сообщение у всех остальных.
  • Порядок внутри транзакции: сначала INSERT отметки, потом работа. Так конкурентные доставки одного сообщения (две реплики консьюмера после ребалансировки) упрутся в блокировку по первичному ключу и сериализуются, а не выполнят работу дважды.
  • Ack брокеру уходит строго после COMMIT. Ack до коммита возвращает нас к at-most-once со всеми его потерями.
  • Ретеншен. Окно дедупликации должно быть шире самого долгого возможного повтора: retention топика, длительности retry-лестницы, планового replay. Чистят пачками DELETE по индексу processed_at, а не одним запросом на миллион строк. Если чистить слишком агрессивно, повторное проигрывание старого топика создаст полный комплект дубликатов.
  • Побочные эффекты вне БД. Отправленное письмо транзакция не откатит, вызов внешнего API не отменит. Помогает идемпотентный внешний вызов с Idempotency-Key или запись «надо отправить письмо» в свой outbox — тогда эффект уйдёт наружу из релея, а не из середины транзакции.
Когда таблица не нужна

Если операция естественно идемпотентна, дедупликация не нужна вовсе: INSERT ... ON CONFLICT DO NOTHING, UPDATE ... SET status='paid', UPSERT проекции по ключу. Опасны только накопительные операции (balance = balance + 100, counter++, отправка письма) — именно ради них и заводят processed_messages. Есть и лёгкий вариант — SETNX в Redis с TTL: быстро и без нагрузки на БД, но не атомарно с бизнес-работой (упали между SETNX и работой — сообщение считается обработанным, хотя его не обработали). Это компромисс для потоков, где дубликат неприятен, а потеря терпима.

Competing consumers

Несколько одинаковых консьюмеров разбирают одну очередь. Так сервис масштабируют добавлением подов: брокер сам следит, чтобы каждое сообщение досталось ровно одному обработчику, и приложению не нужны ни шардирование, ни координация.

СистемаКак выглядитПотолок
RabbitMQНесколько basic.consume на одну очередь, раздача по кругу с учётом prefetchПрактически нет; упирается в производительность самой очереди (один Erlang-процесс)
KafkaОдин group.id, партиции распределяются между участникамиЧисло партиций: лишние консьюмеры простаивают
NATSQueue group: подписчики с одинаковым именем группыНет
SQSНесколько получателей, visibility timeout прячет сообщение на время обработкиНет

Что ломается при переходе от одного консьюмера к нескольким:

  • Порядок. Это главное. Два консьюмера работают как два независимых потока: сообщение №2 может завершиться раньше №1, и «заказ оплачен» обработается раньше «заказ создан». Если порядок важен, его сохраняют шардированием по ключу: в Kafka партиционируют по order_id, в RabbitMQ берут плагин consistent-hash exchange или явный набор очередей на шард. Порядок тогда гарантирован внутри ключа, а бизнесу большего и не надо: события разных заказов независимы.
  • Гонки на общих данных. Параллельные обработчики одного агрегата дерутся за строку в БД. Спасает шардирование (см. выше), оптимистическая блокировка по версии или SELECT ... FOR UPDATE, но последнее берут осознанно: оно сериализует обработку и ставит потолок пропускной способности.
  • Poison message. Одно вечно падающее сообщение при requeue=true съедает всех консьюмеров по очереди. Ограничение попыток и DLQ обязательны.
  • Долгие задачи. Пока обработчик занят, брокер ждёт признаков жизни: в Kafka это max.poll.interval.ms (превысишь — консьюмера исключат из группы и запустят ребалансировку), в JetStream AckWait, в SQS продление visibility timeout. Долгие задачи лучше дробить или уносить в отдельный поток с большими таймаутами, чтобы не тащить эти настройки на весь сервис.

Fanout / multicast и job distribution

Эти два режима постоянно путают, хотя разница ровно в одном: сколько раз сообщение будет обработано. При fanout каждый подписчик получает свою копию (широковещание). При job distribution оно достаётся ровно одному из группы (конкуренция). Оба режима даёт одна и та же инфраструктура, и ошибка в конфигурации молча превращает один в другой.

СистемаFanout (каждому копия)Job distribution (одному из группы)
RabbitMQfanout/topic exchange + своя очередь на каждого подписчикаОдна очередь + несколько консьюмеров
KafkaРазные group.id у каждого сервисаОдин group.id на все инстансы сервиса
NATSОбычные подписчики на subjectQueue group с общим именем
JetStreamНесколько durable-консьюмеров стримаОдин durable-консьюмер, много клиентов
Две классические ошибки конфигурации
  • «Fanout», у которого одна очередь на всех. Три сервиса подписались на fanout exchange через общую очередь — и каждое событие достаётся одному случайному из них. Снаружи это выглядит как «иногда письмо не приходит» и отлаживается неделями. Правило: сколько подписчиков, столько и очередей.
  • Общий group.id у разных сервисов в Kafka. Ошибка зеркальная: два разных сервиса с одинаковым group.id делят партиции, и каждый видит половину потока. Правило: group.id = имя сервиса, никогда не «app» и не имя топика.

Отдельный вопрос — кто плодит копии. В Kafka копия одна: N групп читают один лог, данные хранятся однократно. RabbitMQ при fanout физически размножает сообщение по очередям: у 10 подписчиков будет 10 копий в памяти и на диске. С крупными сообщениями это быстро становится заметно, и лечится claim check: в сообщение кладут ссылку на объект в S3, а не сам объект.

Что класть в событие: notification, state transfer, event sourcing

Наполнить сообщение можно тремя способами, и каждый по-своему распределяет связанность между сервисами. Этот выбор важнее выбора брокера: от него зависит, сможет ли потребитель работать, когда источник лежит.

Один и тот же факт, две формы сообщения — и разная судьба потребителя Event notification: «случилось, детали спроси сам» orders источник истины { "type":"OrderCreated", "order_id":"42", "version":7 } shipping получил только id обратный синхронный вызов GET /orders/42 — связанность вернулась источник лежит — потребитель встал. И возможна гонка Event-carried state transfer: «случилось, вот всё состояние» orders источник истины { "type":"OrderCreated", "order_id":"42", "version":7, "items":[…], "total":1200, "address":{…}, "user":{…} } сообщение выросло в 20 раз shipping работает автономно локальная реплика заказа источник может лежать сутки Цена notification: N+1 синхронных вызовов, доступность перемножается, и гонка «событие обогнало реплику БД». Цена state transfer: трафик, копии данных у всех, жёсткая связанность по схеме и разъезжающиеся по системе персональные данные.
Тонкое событие экономит трафик, но возвращает синхронную зависимость; толстое покупает автономию потребителя ценой дублирования данных и связанности по схеме. На практике чаще выбирают середину: id, версия и три-четыре поля, которых хватает большинству.
Event notificationEvent-carried state transferEvent sourcing
Что в сообщенииИдентификатор и тип фактаПолное состояние сущности на момент событияВсе изменения; события — сами по себе источник истины
Состояние живётТолько в сервисе-источникеКопия у каждого потребителяВосстанавливается свёрткой лога, плюс снапшоты
РазмерСотни байтКилобайты и десятки килобайтМелкие события, но их много и хранятся долго
СвязанностьОбратная рантайм-зависимость: нужен живой источникЗависимость от схемы: поменял поле — сломал всехСхема событий вечна: старые записи придётся читать всегда
Источник недоступенОбработка встаётПотребитель продолжает работатьНе применимо: лог и есть источник
СогласованностьЧитаем свежее, но можно поймать гонку «событие быстрее реплики»Данные слегка устаревшие, зато согласованные с событиемПолная история, любые temporal-запросы
ПлюсТонкий контракт, нет дублированияАвтономия, устойчивость, легко строить проекцииАудит, replay, пересборка проекций, отладка «как дошли до такого»
МинусN+1 вызовов, доступность перемножаетсяТрафик, копии данных, PII расползаетсяВысокая сложность, версионирование навсегда, тяжело удалять данные
Когда берутСобытия редкие, данные крупные, потребителей малоДефолт для интеграции доменов; кэши и локальные репликиДомены, где ценна история: деньги, склад, юридически значимые действия
Гонка, которую стоит назвать вслух

У event notification есть неочевидный отказ: сервис orders закоммитил заказ, опубликовал событие — и shipping тут же пошёл за деталями через GET /orders/42, а балансировщик отправил запрос на реплику для чтения, куда изменения ещё не доехали. В ответ приходит 404 на объект, о создании которого мы только что получили событие. Лечат это чтением с мастера для этого сценария, ретраем с бэкоффом, номером версии в событии с опросом «пока version < 7 — жди» или переходом на state transfer, где данные уже лежат в сообщении.

Стоит предложить компромисс: тонкое событие плюс самые нужные поля и номер версии. Большинству потребителей хватает того, что пришло, а тем, кому нужны детали, остаётся один вызов вместо всех. Для больших полезных нагрузок есть claim check: в событии ссылка на S3, а тело лежит там.

Event sourcing не сводится к «мы используем Kafka»

Публикация событий в брокер сама по себе не делает систему event-sourced. При event sourcing источником истины служит лог событий, а текущее состояние остаётся производной величиной, которую в любой момент можно выбросить и пересобрать. За это платят: схемы событий живут вечно (десятилетний лог придётся читать сегодняшним кодом, отсюда upcasting), нужны снапшоты, «просто поправить строку в базе» нельзя, а персональные данные по GDPR удаляют через crypto shredding: данные шифруют ключом на субъекта, а по запросу на удаление уничтожают ключ. Здоровая практика: применять подход к одному-двум агрегатам, где история действительно ценна, а не ко всей системе.

Схемы сообщений и их эволюция

Сообщение — такой же публичный API, только без компилятора и без единого места, где его можно посмотреть. HTTP-контракт хотя бы ломается сразу и в тестах, а событие ломается через полгода у команды, о существовании которой ты не знал, — и не падением, а тихим неправильным поведением. Отсюда две практики: описывать схему явно и проверять совместимость машиной.

Protobuf

В protobuf контракт держится на номерах полей, а не на именах: имена в wire-формате не передаются вовсе. Из этого и следуют все правила.

message OrderCreated {
  string order_id = 1;
  string user_id  = 2;
  int64  total    = 3;         // в копейках: менять тип потом будет нельзя

  reserved 4, 7 to 9;          // номера удалённых полей больше не занимать
  reserved "discount_code";    // и имя тоже, чтобы не вернулось с другим смыслом

  repeated Item items    = 5;
  optional string coupon = 6;  // explicit presence: отличаем "нет значения" от ""
}
ИзменениеБезопасно?Что произойдёт
Добавить поле с новым номеромДаСтарый консьюмер положит его в unknown fields и проигнорирует
Переименовать поле (номер тот же)Да для бинарного форматаНо ломает JSON-маппинг и весь код, который читает поле по имени
Удалить полеДа, с reservedСтарый консьюмер увидит дефолт вместо значения — семантику проверить обязательно
Изменить тип поляНетint32 → string меняет wire type: старый читатель отложит поле в unknown fields и увидит ноль, без ошибки
Изменить номер поляНетДля старого читателя поле исчезло, а новое пришло под чужим номером
Переиспользовать номер удалённого поляКатегорически нетСамая опасная поломка: старые данные декодируются успешно, но означают другое
Добавить значение в enumУсловноФормат выдержит, а switch без default у консьюмера — нет
Сделать поле обязательным по смыслуНетСтарые продюсеры его не заполняют; валидация уронит консьюмера
Зачем нужен reserved

Поле 4 было string discount_code, его удалили. Через год другой разработчик добавляет int32 quantity = 4: номер же свободен. Теперь старое сообщение с discount_code = "SALE" прилетает новому консьюмеру, тот ждёт varint, получает length-delimited и молча откладывает поле в unknown fields: quantity = 0, ни одной ошибки. А при совпавшем wire type он прочитает чужое число и пойдёт с ним дальше. С reserved это становится ошибкой компиляции .proto и всплывает в единственный момент, когда её ещё дёшево исправить. reserved "имя" делает то же для JSON-представления и истории.

Avro и Schema Registry

Avro устроен иначе: схема не зашита в код, а сопровождает данные. Читатель делает schema resolution: сопоставляет writer schema и reader schema по именам полей и подставляет default для тех, которых нет. Отсюда правило: у любого добавляемого поля должен быть default, иначе старые данные перестанут читаться.

Schema Registry («реестр схем») — отдельный маленький сервис рядом с брокером: он хранит все версии схем сообщений и раздаёт им номера. Нужен он потому, что сама Kafka о содержимом сообщений ничего не знает и проверить его не может: для неё это просто байты. В реестре договорённость о формате становится машинно проверяемой.

Confluent Schema Registry раскладывает версии по subject (обычно <topic>-value) и выдаёт им числовые id. В сообщение пишется префикс: magic byte 0 + 4 байта schema id, дальше тело. Консьюмер вытаскивает id, тянет схему из реестра (и кэширует) и декодирует. Приятный побочный эффект, который стоит упомянуть: схема не передаётся с каждым сообщением — только 5 байт.

Но главное в реестре не хранение, а проверка совместимости при регистрации: несовместимую схему он просто не примет, и продюсер упадёт на старте, а не на проде через месяц.

РежимЧто гарантируетЧто разрешено менятьПорядок деплоя
BACKWARD (дефолт)Новый консьюмер прочитает данные, записанные старой схемойУдалять поля; добавлять optional с defaultСначала консьюмеры, потом продюсеры
FORWARDСтарый консьюмер прочитает данные новой схемыДобавлять поля; удалять optional с defaultСначала продюсеры, потом консьюмеры
FULLОба направления сразуТолько добавление/удаление optional с defaultЛюбой
*_TRANSITIVEТо же, но против всех прошлых версий, а не только последнейСтроже: цепочка мелких «совместимых» шагов не сможет тихо увести схему далеко от исходнойЛюбой
NONEНичегоВсёМолитва

В Kafka по умолчанию стоит BACKWARD, потому что в топике лежат старые записи (retention неделя, а то и вечность при compaction), и новый консьюмер обязан их прочитать. FORWARD берут, когда консьюмеров обновить нельзя: внешние клиенты, мобильные приложения, партнёрские интеграции. Продюсер уходит вперёд, а старый код должен продолжать работать. FULL нужен, когда порядок выкатки не контролируется вовсе, а в микросервисах это скорее норма, чем исключение.

«Схема совместима» ≠ «семантика совместима»

Реестр проверяет форму, а не смысл. Добавили в enum status значение PARTIALLY_REFUNDED — схема совместима, а консьюмер со switch без default либо паникует, либо тихо считает заказ неоплаченным. Поменяли total с рублей на копейки, не тронув тип, и схема по-прежнему совместима, а счета выросли в сто раз. Закрывают это такие правила: tolerant reader (игнорировать неизвестное, всегда иметь ветку по умолчанию); смысл существующего поля не менять никогда, только добавлять новое; класть schema_version в заголовки сообщения; держать контракты в общем репозитории с проверкой совместимости в CI, а не в реестре постфактум.

Backpressure: консьюмер не успевает

Напомним определение из главы 8.1: при backpressure (обратном давлении) медленный получатель заставляет быстрого отправителя сбавить темп. Здесь разберём, что делать, если этого механизма не хватило и невыполненная работа начала копиться.

Сначала диагноз, потом лечение. Растущий лаг сам по себе ничего не значит: при пике он растёт и возвращается, при проблеме растёт монотонно. Различают их по одной величине: скорость разгребания = consume_rate − produce_rate. Если она устойчиво отрицательная, никакой буфер не поможет — вопрос лишь в том, когда он кончится. На собесе полезно произнести такую оценку: время до катастрофы = свободный объём буфера / (produce − consume). Она сразу переводит разговор из «надо что-то делать» в «у нас четыре часа».

Продюсер быстрее консьюмера — вопрос только в том, что переполнится первым продюсер 30 000 msg/s брокер буфер на диске 4 000 консьюмер 4 000 msg/s Долг растёт на 26 000 сообщений в секунду. Это ≈ 1,6 млн в минуту. Любой буфер — это просто отсрочка, измеряемая в часах. Что переполнится первым — зависит от того, где стоит самое слабое звено ① Продюсер, буфер отправки Kafka: buffer.memory (32 МБ) полон → send() блокируется на max.block.ms, потом TimeoutException. Это хороший исход: давление дошло до кода. ② Брокер, диск и лимиты Kafka: retention.ms/bytes молча удаляет хвост — консьюмер получит OffsetOutOfRange и пропустит данные. Потеря без единой ошибки у продюсера. ③ Консьюмер, память Push без prefetch: брокер заливает всё, что есть → RSS растёт → OOM-killer. Под перезапускается, unacked возвращаются, цикл повторяется — «краш-луп под нагрузкой». ④ RabbitMQ: memory alarm (0.6 от RAM) или disk_free_limit Брокер перестаёт читать TCP от всех продюсеров узла — publish зависает даже у тех, кто не виноват. ⑤ Как надо: давление доходит до источника Quota на брокере, rate limiter в сервисе, 429 на входе HTTP. Отказ на границе честнее, чем OOM в середине пайплайна. Правило: буфер покупает время, а не пропускную способность. Если consume < produce устойчиво — надо менять одно из двух чисел, а буфер использовать только чтобы дожить до момента, когда вы это сделаете.
«Увеличим буфер» ничего не решает, только выбирает место, где рванёт. Полезно заранее знать, что переполнится первым: от этого зависит, увидишь ты аккуратный TimeoutException у продюсера или молчаливую потерю хвоста топика.

Инструменты по порядку применения

  • Горизонтальное масштабирование консьюмеров. Его пробуют первым и в него же первым упираются: в Kafka потолок задаёт число партиций, лишние поды в группе просто простаивают. В RabbitMQ потолка почти нет, но растёт нагрузка на сам процесс очереди. Прежде чем масштабировать, стоит проверить, что консьюмеры заняты работой, а не ждут БД: если ждут, новые поды только добьют базу.
  • Ускорить обработку. Обычно тут и лежит основной резерв: батчить записи в БД вместо построчных, убрать N+1, вынести побочные вызовы, распараллелить внутри одного консьюмера воркер-пулом по hash(key) % N, коммитя только сплошной префикс обработанных офсетов. Так параллелизм растёт без потери порядка по ключу.
  • Увеличить число партиций, понимая цену. (а) Меняется hash(key) % N, ключ уезжает в другую партицию, и порядок по ключу ломается на границе изменения: старые события заказа лежат в старой партиции, новые в новой, и два консьюмера обрабатывают их одновременно. (б) Уменьшить число партиций нельзя вообще. (в) Растут файловые дескрипторы, память на брокерах, время ребалансировки и выборов лидеров. (г) Уже лежащие сообщения не перераспределяются, так что сразу легче не станет.
  • Throttling продюсера, то есть ограничение на источнике: quotas в Kafka (producer_byte_rate на client.id/user: брокер задерживает ответы, и продюсер тормозит сам), rate limiter в сервисе, 429 на HTTP-границе. Обратное давление должно доходить до того, кто порождает трафик, иначе оно копится где-то посередине и рвётся там.
  • Отбрасывание и приоритеты (load shedding). Заранее решить, что можно потерять. Аналитику, телеметрию, «пользователь открыл экран» можно, платежи нельзя. Технически это разные топики или очереди под разные классы с отдельными пулами консьюмеров (bulkhead); в RabbitMQ x-max-length с overflow=drop-head (вытесняем самое старое) или reject-publish (отказываем на входе, что честнее); TTL для данных, которые через минуту бессмысленны, — котировки, статусы онлайна. Осознанное отбрасывание всегда лучше неосознанного: второе всё равно случится, только жертву выберет само.
  • Буферизация и её предел. Брокер и есть буфер, и копить правильно именно в нём: диск дешевле памяти консьюмера. Но он конечен: его ограничивают retention.ms/retention.bytes, x-max-length и размер диска. Буфер даёт время дожить до пятого пункта, но не заменяет его.
Как отвечать на этот вопрос

По ступеням, от дешёвого к дорогому, и обязательно с ценой каждого шага. «Сначала смотрю, пик это или тренд, и считаю скорость разгребания. Потом проверяю, упёрлись ли консьюмеры в CPU или ждут БД: если ждут, новые поды сделают хуже. Дальше по порядку: масштабирование до числа партиций, оптимизация и батчинг, увеличение партиций (с пониманием, что это ломает порядок по ключу и необратимо), троттлинг продюсера и в крайнем случае осознанное отбрасывание того, что не жалко. Буфер брокера при этом не решение, а время на то, чтобы решение применить». Такой ответ показывает, что кандидат думает о системе, а не об одной ручке.

Вопросы

6
Суть: решает двойную запись — «сохранить в БД и опубликовать событие» атомарно нельзя, потому что это две системы без общего коммита. Outbox переносит публикацию внутрь транзакции БД: событие пишется в таблицу тем же COMMIT, а отдельный релей выносит его в брокер. Гарантия — at-least-once, дубликаты остаются.

Почему без него ломается

Порядков два, и обе поломки надо назвать. COMMIT, потом publish: если упасть между ними, заказ в базе есть, а события нет и не будет — склад и биллинг о нём не узнают. Publish, потом COMMIT: событие ушло, а транзакция после падения откатилась, консьюмеры реагируют на заказ-призрак и получают 404 из API. Ретраи в памяти не спасают, они умирают вместе с процессом. 2PC/XA формально решает, но Kafka его не поддерживает, координатор блокирующий, а его падение оставляет in-doubt транзакции, держащие блокировки в базе.

Реализация

Таблица outbox(id, aggregate_type, aggregate_id, event_type, payload, headers, created_at, published_at) в той же базе, что и бизнес-данные, плюс частичный индекс WHERE published_at IS NULL — он маленький, потому что покрывает только хвост. В коде один tx.Commit() на оба запроса, INSERT orders и INSERT outbox. Всё, атомарность обеспечила СУБД.

Relay: два варианта

  • Polling. SELECT ... WHERE published_at IS NULL ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED → publish → UPDATE published_at. Благодаря SKIP LOCKED несколько релеев работают, не блокируя друг друга, но порядок событий одного агрегата тогда не гарантирован. Плюсы: ничего лишнего в инфраструктуре. Минусы: задержка ≈ полинтервала, постоянные запросы и лишний WAL от UPDATE, таблицу надо убирать.
  • CDC (Debezium). Читает WAL через логический слот и публикует изменения outbox-таблицы; Outbox Event Router SMT достаёт payload и ставит топик и ключ по aggregate_id. Плюсы: задержка в миллисекундах, нулевая нагрузка на базу от опроса, порядок = порядок коммитов. Минусы: Kafka Connect как отдельная система и главный риск — нечитаемый replication slot не даёт PostgreSQL удалять WAL, и диск базы заканчивается.
Вопрос на добивание: почему нельзя вести курсор по id

Потому что bigserial выдаёт номера вне транзакции. T1 взяла id=5 и ещё не закоммитилась, T2 взяла id=6 и закоммитилась; релей видит 6, сдвигает курсор — и строку 5, закоммиченную мгновением позже, не выберет никогда. Отметка должна жить на строке (published_at или DELETE), а не в курсоре. Тот же аргумент хоронит фильтр по created_at.

Гарантии и цена

Outbox даёт at-least-once, а не exactly-once. Дубликаты неизбежны: если релей опубликовал и упал до UPDATE published_at, он опубликует снова; плюс ретраи продюсера, плюс перезапуск с середины пачки. Поэтому outbox почти всегда идёт в паре с inbox или идемпотентным обработчиком на приёме. Порядок гарантирован внутри агрегата, если ключом партиционирования взять aggregate_id. Ещё платишь ростом таблицы (нужна чистка партициями или пачками), лишней записью в каждой бизнес-транзакции и ещё одним процессом, за которым надо следить. Главный алерт всей конструкции — метрика «возраст самой старой неопубликованной строки».

Суть: зеркальное отражение outbox. На приёме тоже две операции — «сделать работу» и «запомнить, что сделали», — и их тоже надо объединить в одну транзакцию. Таблица processed_messages с уникальным ключом, INSERT первым, бизнес-работа следом, ack брокеру строго после COMMIT.

Механика

Нужна таблица processed_messages(consumer, message_id, processed_at) с первичным ключом по паре «консьюмер + идентификатор сообщения». Обработчик открывает транзакцию и делает INSERT ... ON CONFLICT DO NOTHING. Ноль затронутых строк означает дубликат: работу пропускаем и просто подтверждаем брокеру. Если строка вставилась, выполняем бизнес-работу в той же транзакции и коммитим. Ошибка в работе откатывает и отметку, поэтому сообщение честно вернётся и обработается снова.

Порядок внутри транзакции важен: сначала отметка. INSERT берёт блокировку по первичному ключу, и две конкурентные доставки одного сообщения (две реплики консьюмера после ребалансировки) сериализуются, а не делают работу дважды.

Выбор ключа

  • Хорошо: message_id, сгенерированный продюсером один раз и неизменный при всех ретраях (в outbox это id строки), либо бизнес-ключ order_id + event_type + version.
  • Плохо: хеш тела. Ретрай с обновлённым timestamp даст другой хеш и проедет как новое, а два легитимных одинаковых события (два списания по 100 ₽), наоборот, схлопнутся в одно. Вместо дедупликации получится потеря данных.
  • Плохо: глобальная уникальность без имени консьюмера — первая группа «съест» сообщение у всех остальных.

Ретеншен

Окно дедупликации должно быть шире максимально возможного повтора: retention топика, длительность retry-лестницы, планируемые replay. Чистят пачками DELETE по индексу processed_at, а не одним запросом на миллион строк. Перестараешься с чисткой, и повторное проигрывание старого топика даст полный набор дубликатов.

Два уточнения, которые отличают сильный ответ
  • Не всегда нужна таблица. Если операция естественно идемпотентна (UPSERT по ключу, SET status='paid'), дедупликация не нужна вовсе. Проблемы бывают только у накопительных операций (balance = balance + 100) и внешних эффектов.
  • Побочные эффекты вне БД транзакция не откатит. Письмо отправлено, платёж в шлюзе проведён. Помогает идемпотентный внешний вызов с Idempotency-Key или запись «надо отправить» в свой outbox, чтобы эффект уходил наружу уже из релея. Redis SETNX с TTL работает быстро, но не атомарно с бизнес-работой: упали между SETNX и работой — сообщение считается обработанным, хотя ничего не произошло.
Суть: два режима доставки, отличающиеся одним — сколько раз сообщение будет обработано. Job distribution (competing consumers): ровно один из группы, масштабирование подами. Fanout/multicast: копия каждому подписчику. Ошибка конфигурации молча превращает один режим в другой.

Competing consumers

Несколько одинаковых обработчиков читают одну очередь, а брокер следит, чтобы сообщение досталось одному. В RabbitMQ это несколько basic.consume на очередь с раздачей по кругу и оглядкой на prefetch, потолка практически нет. В Kafka нужен общий group.id, и потолок жёсткий: число партиций, лишние консьюмеры простаивают. В NATS это queue group, в SQS visibility timeout.

Что ломается при переходе от одного консьюмера к N: порядок (в двух потоках «оплачен» может обработаться раньше «создан»; лечится шардированием по ключу: партиции Kafka или consistent-hash exchange в RabbitMQ), гонки на общих данных (оптимистическая блокировка по версии или тот же шардинг), poison message (нужны лимит попыток и DLQ, иначе одно сообщение обойдёт всех консьюмеров по очереди), долгие задачи (max.poll.interval.ms в Kafka, AckWait в JetStream, продление visibility в SQS).

Fanout vs job distribution

СистемаFanout: копия каждомуJob distribution: одному
RabbitMQfanout/topic exchange + своя очередь на подписчикаОдна очередь, несколько консьюмеров
KafkaРазные group.idОдин group.id
NATSОбычные подписчики на subjectQueue group с общим именем
Две ошибки, которые ищут этим вопросом
  • Три сервиса подписаны на fanout через одну общую очередь. Каждое событие достаётся одному случайному из них. Выглядит это как «иногда письмо не приходит» и отлаживается неделями. Очередей нужно столько же, сколько подписчиков.
  • Одинаковый group.id у разных сервисов в Kafka. Зеркальная беда: сервисы делят партиции и видят по половине потока. В group.id пишут имя сервиса, а не имя топика и не «app».

И про стоимость. В Kafka веерная рассылка бесплатна по хранению: N групп читают один лог, копия одна. RabbitMQ при fanout физически размножает сообщение: десять подписчиков дают десять копий в памяти и на диске. Если сообщения крупные, помогает claim check: в событие кладут ссылку на объект в S3, а не сам объект.

Суть: выбор между «тонким событием + обратный вызов за деталями» и «толстым событием, несущим состояние». Первое экономит трафик, но возвращает синхронную зависимость от источника; второе даёт потребителю автономию ценой дублирования данных и жёсткой связанности по схеме. Третья степень — event sourcing, где события сами и есть источник истины.

Event notification

В сообщении только идентификатор и тип факта: {OrderCreated, order_id: 42}. Плюсы: контракт тонкий, данные не дублируются, потребитель всегда читает свежее состояние. Минусы серьёзнее, чем кажется. Возвращается рантайм-зависимость: источник лёг, и обработка встала, доступность снова перемножается. На поток событий появляется N+1 синхронных вызовов. И есть неприятная гонка — событие обгоняет репликацию, консьюмер идёт в read-реплику за заказом и получает 404 на объект, о создании которого только что узнал.

Event-carried state transfer

Сообщение несёт полное состояние сущности на момент события. Потребитель становится автономным: держит локальную реплику, переживает недоступность источника, не делает обратных вызовов. Минусы: килобайты вместо сотен байт, копия данных у каждого потребителя, связанность по схеме (изменил поле — сломал всех подписчиков) и персональные данные, расползающиеся по системе: из-за них удаление по GDPR превращается в проект (в Kafka через log compaction и tombstone-записи).

Event sourcing

Уже не про форму сообщения, а про хранение: истину хранит лог событий, а текущее состояние — лишь свёртка, которую можно выбросить и пересобрать. Отсюда полный аудит, temporal-запросы, пересборка проекций и честная отладка «как мы дошли до такого». Цена: схемы событий живут вечно (десятилетний лог придётся читать сегодняшним кодом, отсюда upcasting), нужны снапшоты, нельзя «просто поправить строку», а персональные данные удаляют через crypto shredding: шифруют ключом на субъекта и уничтожают ключ. Здоровая практика — применять его к одному-двум агрегатам, где история ценна, а не ко всей системе. Публикация событий в Kafka сама по себе ещё не event sourcing.

Как выбирают на практике

Для интеграции между доменами по умолчанию берут state transfer: автономия потребителя дороже трафика. Notification выбирают, когда данные крупные или чувствительные (не хочется рассылать PII всем подряд) и потребителей мало. Чаще всего живут посередине: id, номер версии и три-четыре самых нужных поля. Большинству этого хватает, остальные делают один вызов вместо всех. Для тяжёлых полезных нагрузок есть claim check: в событии ссылка, тело в объектном хранилище. А version в событии дёшево страхует от гонок и от применения событий не по порядку.

Суть: сообщение — публичный API без компилятора. Protobuf держит контракт на номерах полей (отсюда reserved), Avro — на именах с обязательными default. Schema Registry не столько хранит схемы, сколько не даёт зарегистрировать несовместимую. BACKWARD — обновляем консьюмеров первыми, FORWARD — продюсеров, FULL — порядок не важен.

Protobuf

Имена полей в wire-формате не передаются, передаются номера. Отсюда и все правила. Поле с новым номером добавлять безопасно (старый консьюмер положит его в unknown fields). Переименование безопасно для бинарного формата, но ломает JSON-маппинг. Удалять поле можно только вместе с reserved. Смена типа или номера ломает контракт, как и поле, ставшее обязательным по смыслу: старые продюсеры его не заполняют.

Зачем reserved — главный вопрос на добивание

Поле 4 было string discount_code, и его удалили. Через год кто-то добавляет int32 quantity = 4: номер же свободен. Старое сообщение прилетает новому консьюмеру, тот ждёт varint, получает length-delimited и молча откладывает поле в unknown fields: quantity = 0, ошибки нет. А при совпавшем wire type он прочитает чужое число и пойдёт с ним дальше. Это худший класс поломок: данные декодируются успешно, но означают другое. С reserved 4; и reserved "discount_code"; это становится ошибкой компиляции .proto и всплывает в единственный момент, когда чинить ещё дёшево.

Avro и Schema Registry

В Avro схема сопровождает данные, и читатель делает schema resolution: сопоставляет writer schema с reader schema по именам полей и подставляет default для отсутствующих, поэтому у любого добавляемого поля default обязателен. В Confluent-формате идут magic byte 0 + 4 байта schema id, дальше тело. Консьюмер тянет схему по id из реестра и кэширует, так что по проводу схема не ходит.

BACKWARD vs FORWARD

РежимГарантияПорядок деплояКогда нужен
BACKWARD (дефолт)Новый консьюмер читает старые данныеСначала консьюмерыKafka: в топике лежит история, её надо уметь прочитать
FORWARDСтарый консьюмер читает новые данныеСначала продюсерыКонсьюмеров обновить нельзя: внешние клиенты, мобильные приложения
FULLОба направленияЛюбойПорядок выкатки не контролируется — обычная ситуация в микросервисах
*_TRANSITIVEТо же против всех версий, не только последнейЛюбойЧтобы цепочка мелких совместимых шагов не увела схему далеко от исходной
«Совместимо по схеме» не значит «совместимо по смыслу»

Реестр проверяет только форму. Добавили в enum status значение PARTIALLY_REFUNDED — схема совместима, а консьюмер со switch без default либо падает, либо считает заказ неоплаченным. Перевели total из рублей в копейки, не тронув тип, и схема всё ещё совместима, а счета выросли в сто раз. Спасают такие практики: tolerant reader (игнорировать неизвестное, всегда иметь ветку по умолчанию), смысл существующего поля не менять, а добавлять новое, schema_version в заголовках сообщения, контракты в общем репозитории и проверка совместимости в CI, а не в реестре постфактум.

Суть: сначала диагноз — пик это или тренд, и упирается ли консьюмер в CPU или ждёт БД. Потом ступени от дешёвого к дорогому: масштабирование до числа партиций, оптимизация обработки, увеличение партиций (необратимо и ломает порядок по ключу), троттлинг продюсера, осознанное отбрасывание. Буфер брокера покупает время, а не пропускную способность.

Диагностика

Растущий лаг сам по себе не диагноз. Считаем скорость разгребания = consume_rate − produce_rate: при пике она положительна и лаг вернётся, при проблеме устойчиво отрицательна. Оценка, которую полезно произнести вслух: время до переполнения = свободный объём / (produce − consume) — она сразу переводит панику в «у нас четыре часа». Второй вопрос: консьюмеры действительно заняты работой или ждут БД и внешние вызовы? Если ждут, новые поды только добьют базу.

Ступени

  • Горизонтальное масштабирование в Kafka работает до числа партиций (дальше лишние поды просто простаивают), в RabbitMQ потолка почти нет.
  • Основной резерв обычно в том, чтобы ускорить обработку: батчить записи в БД, убрать N+1, распараллелить внутри консьюмера воркер-пулом по hash(key) % N, коммитя только сплошной префикс офсетов, чтобы не потерять порядок по ключу.
  • Увеличить партиции и назвать цену. Меняется hash(key) % N, ключ уезжает в другую партицию, и порядок по ключу ломается на границе изменения; уменьшить обратно нельзя вообще; растут дескрипторы, память брокеров, время ребалансировок; уже лежащие сообщения не перераспределяются, поэтому мгновенного облегчения не будет.
  • Троттлинг продюсера делают через quotas в Kafka (producer_byte_rate: брокер задерживает ответы, продюсер тормозит сам), rate limiter в сервисе или 429 на HTTP-границе. Давление должно дойти до источника трафика, иначе оно копится посередине и рвётся там.
  • Отбрасывание и приоритеты. Заранее решают, что можно потерять (телеметрия, аналитика), а что нельзя (платежи), и разводят по разным топикам с отдельными пулами консьюмеров. В RabbitMQ есть x-max-length с overflow=drop-head или, честнее, reject-publish, а для данных, бессмысленных через минуту, TTL. Осознанное отбрасывание лучше неосознанного: второе всё равно случится, просто жертву выберет само.
  • Буферизация. Брокер и есть буфер, и копить правильно именно в нём: диск дешевле памяти консьюмера. Но он конечен (retention.ms/ retention.bytes, x-max-length, размер диска) и даёт только время, чтобы успеть применить предыдущие пункты.
Что переполнится первым и почему это надо знать заранее
  • Буфер продюсера. Kafka: buffer.memory полон → send() блокируется на max.block.ms → ошибка. Это хороший исход: давление дошло до кода и его видно.
  • Диск брокера. Kafka молча удаляет хвост по retention, консьюмер получает OffsetOutOfRange и пропускает данные, а продюсер не видит ни одной ошибки.
  • Память консьюмера. Push-модель без prefetch: брокер заливает всё, что есть, RSS растёт, OOM-killer убивает под, unacked возвращаются, и цикл повторяется. Классический краш-луп под нагрузкой.
  • Узел RabbitMQ. Memory alarm (0,6 от RAM, в 3.x было 0,4) или disk_free_limit заставляют брокер перестать читать TCP от всех продюсеров узла: publish зависает даже у тех, кто не создавал нагрузки.