Брокеры сообщений и асинхронность
Секция, где проверяют не знание кнопок в 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-й ошибке теряется целиком — если только вызывающий сам не построил себе очередь ретраев, то есть тот же брокер, только плохой.
Ответ «асинхронность нужна для развязки и надёжности» звучит как заученный. В сильном ответе всегда есть и цена: клиенту теперь нельзя сказать «готово» (нужен статус/поллинг/вебхук), появились дубликаты (нужна идемпотентность), порядок больше не глобальный, отладка стала многошаговой (нужны correlation id и трассировка через заголовки сообщений), и сам брокер стал единой точкой отказа, которую надо кластеризовать и мониторить.
Push vs pull
Вопрос в том, кто начинает передачу: брокер или потребитель. В модели push («толкать») брокер сам шлёт сообщение, как только оно появилось. В pull («тянуть») консьюмер сам приходит и спрашивает: «есть что-нибудь?». Это не деталь реализации: от неё зависит, где именно копится работа, когда потребитель не успевает.
Через backpressure (обратное давление) медленный получатель заставляет быстрого отправителя сбавить темп. Бытовой пример: раздача на кухне ресторана. Если повар ставит тарелки на стойку быстрее, чем официанты их разносят, стойка переполнится, и дальше возможны только два исхода: либо повар видит, что места нет, и притормаживает (обратное давление сработало), либо тарелки летят на пол (потери). В push тормозить отправителя приходится специально, в 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» |
Почему exactly-once невозможен: проблема двух генералов
Классическая формулировка: две армии на холмах по обе стороны от долины, в долине противник. Атака удастся только при одновременном ударе. Гонцы ходят через долину, и их могут перехватить. Доказано, что никакого протокола с конечным числом сообщений, гарантирующего согласованное решение, не существует. Доказательство от противного: возьмём кратчайший корректный протокол; его последнее сообщение может потеряться, но раз протокол корректен и без него, это сообщение лишнее, и мы получили протокол короче. Повторяем, пока сообщений не останется вовсе, и приходим к противоречию.
Прикладной перевод: отправитель никогда не может достоверно узнать, был ли получен и обработан его последний запрос. Не получив ответа, он вынужден выбирать между «повторить» (риск дубликата) и «не повторять» (риск потери). Брокер этот выбор не отменяет, а лишь берёт на себя. Сюда же относится FLP-теорема (Fischer, Lynch, Paterson, 1985): в асинхронной системе, где может отказать даже один узел, детерминированный консенсус за конечное время невозможен.
«Exactly-once доставка невозможна. Достижима exactly-once обработка, и только как сочетание at-least-once доставки с идемпотентным потребителем или транзакцией, которая атомарно фиксирует и результат работы, и факт потребления. То есть проблема решается не на транспорте, а end-to-end, на прикладном уровне.» Именно это и хотят услышать.
Две отдельные вещи, которые маркетинг слепил в одну. Во-первых, идемпотентный продюсер
(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 означает «все реплики из текущего ISR подтвердили».
Если ISR схлопнулся до одной реплики (остальные отстали или упали), acks=all
превращается в acks=1, и падение этого брокера означает потерю. Ровно для этого есть
min.insync.replicas=2: при меньшем ISR продюсер получит
NotEnoughReplicasException и запись не пройдёт. Правильная пара:
RF=3 и min.insync.replicas=2. С ней кластер переживает падение одного брокера
без потерь и без остановки записи.
Очередь vs топик
Базовых моделей обмена две, и различаются они тем, сколько потребителей увидит одно сообщение.
От выбора зависит, кто знает о потребителях. Через очередь идёт команда («сделай задачу»), и отправитель обычно знает, что её выполнят. В топик уходит факт («заказ создан»), и отправителю всё равно, кто на него отреагирует. На собесе по архитектуре из этого вырастает вопрос про commands vs events: команды маршрутизируют в конкретную очередь и валидируют, события публикуют в топик и не валидируют, потому что событие уже случилось и отменять нечего.
Вопросы
8У синхронной цепочки HTTP-вызовов два свойства, которые под нагрузкой становятся смертельными. Доступность перемножается: четыре сервиса по 99,9 % дают 99,6 %, то есть примерно три часа простоя в месяц вместо сорока минут. А задержка суммируется, и клиент ждёт худшую ветку, даже если ему интересна только первая.
Три эффекта, которые даёт брокер
- Развязка (decoupling). Продюсер публикует факт в топик и не знает ни числа подписчиков, ни их состояния. Появился пятый потребитель события «заказ создан» — команда заказов не выкатывает ни строчки кода. Это разница между «изменение в N местах» и «изменение в одном».
- Сглаживание пиков (load levelling). Брокер держит буфер на диске. В чёрную пятницу продюсер выдаёт 30 000 сообщений/с, а консьюмеры продолжают жевать по 4 000/с; растёт лаг, но ничего не падает. Без буфера пришлось бы держать железо под пик, который бывает два раза в год.
- Надёжность. Сообщение лежит в персистентном логе. Потребитель может лежать двадцать минут, подняться и доесть накопившееся. Синхронный запрос при недоступности адресата просто теряется, если только вызывающий сам не построил себе очередь ретраев, то есть плохой брокер.
Чем добить ответ
Обязательно назвать цену, иначе ответ звучит как реклама. За асинхронность платят
конечной согласованностью вместо мгновенной (клиент нажал кнопку, а эффект наступит через
секунду или через минуту), дубликатами, которые надо обрабатывать, потерей сквозного порядка,
отсутствием простого стектрейса (нужен трейсинг с пробросом traceparent в заголовках
сообщения) и ещё одной инфраструктурной системой, которую надо мониторить и обновлять.
Хорошая формулировка: «асинхронность — это не оптимизация, а смена контракта: вместо
ответа мы отдаём обещание».
Когда вызывающему нужен результат прямо сейчас, чтобы показать его пользователю, это работа для синхронного вызова, а «запрос-ответ через две очереди» превращается в RPC с ручной корреляцией, таймаутами и мусорными очередями. Асинхронность оправдана там, где ответ не нужен немедленно: побочные эффекты (письма, индексация, аналитика), долгие задачи, интеграция с внешними системами и веерная рассылка событий.
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). Правильный вывод для собеса: модель доставки определяет,
где копится очередь при перегрузке — в памяти потребителя или на диске брокера.
Второе почти всегда лучше.
Три уровня
- При 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 на стороне записи.
Путь первый: операция идемпотентна сама по себе
Он лучше, потому что не требует хранить состояние. 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. Без confirmsbasic.publishработает как fire-and-forget: TCP-запись прошла, а брокер мог не принять, и ты об этом не узнаешь. Quorum queues добавляют репликацию через Raft и пишут на диск всегда.
«Очередь durable — значит, сообщения не потеряются?» Нет. durable у классической
очереди означает, что переживёт перезапуск сама очередь как объект, а не её содержимое.
Нужны обе галочки: durable queue плюс persistent message. И даже этого мало без publisher
confirms — потерять можно на участке «продюсер отправил, брокер ещё не записал».
Что ломает порядок
- Несколько партиций. Между ними нет общих часов и нет координации. Лечится
ключом партиционирования:
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 контроллера стал на порядок быстрее, а кластер держит миллионы партиций вместо пары сотен тысяч.
__consumer_offsets,
а лаг считается как разница между концом лога и этим указателем.Почему Kafka быстрая
Вопрос любят задавать, потому что ответ проверяет понимание работы с диском и сетью, а не знание Kafka. Четыре причины, по убыванию значимости.
- Последовательная запись append-only. Нет обновлений на месте, нет B-дерева, нет случайных seek. Последовательная запись на обычный HDD даёт сотни МБ/с, что сопоставимо со случайным доступом к памяти. Kafka нарочно спроектирована так, чтобы обращаться к диску только последовательно.
- Page cache вместо своего кэша. Kafka не держит данные в JVM heap, а пишет и читает через страничный кэш ядра. GC от этого не страдает, кэш переживает перезапуск брокера, а «горячие» читатели (те, кто у хвоста лога) вообще не касаются диска. Отсюда практическое правило: брокеру дают маленький heap (6–8 ГБ), а всю остальную память отдают ОС.
- Zero-copy. Консьюмеру данные отдаются через
sendfile(2): page cache → сокет, без копирования в user space и обратно. Это экономит два копирования и два переключения контекста. Оговорка: zero-copy отваливается, если включён TLS (данные надо шифровать в user space) или если брокеру приходится конвертировать формат записи под старого клиента. Отсюда и ответ на вопрос «почему включение TLS уронило пропускную способность вдвое». - Батчинг и единый формат. Продюсер собирает 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) | Дефолт для новых сервисов: нет глобальной паузы |
Обработка одного сообщения замедлилась (тормозит внешний 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 запоминается уже при выдаче, так что коммит может обогнать
обработку и без всяких горутин.
// 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 блокирует и ретраит до успеха. Надёжно, но каждый батч получает лишний
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 дней».
- Восстановление состояния с нуля. Новый сервис подписался на compacted-топик
user-profiles, прочитал его целиком и получил актуальные профили всех пользователей, не трогая чужую БД. На этом построены CQRS-проекции и KTable в Kafka Streams. - Служебные топики самой Kafka.
__consumer_offsetsи__transaction_statecompacted, иначе они росли бы вечно. - 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 storm | max.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
}
}
}
- Ошибки бывают двух видов, и путать их нельзя. Транзиентные (таймаут БД, 503 от соседа) ретраят. Постоянные (невалидный JSON, нет обязательного поля, бизнес-запрет) ретраить бессмысленно, их сразу отправляют в DLQ. Без этого разделения retry-цепочка просто откладывает неизбежное и тратит ресурсы.
- Retry-топики ломают порядок. Сообщение, ушедшее на ретрай, вернётся позже своих соседей по ключу. Если порядок критичен, придётся останавливать обработку всего ключа до разбора или осознанно принять расхождение.
- DLQ без алерта и без процесса разбора превращается в мусорку. Нужны метрика количества, алерт на «первое сообщение за N минут», ретеншен побольше основного топика и инструмент повторной подачи (reprocess из DLQ обратно в основной топик).
Транзакции и exactly-once
Транзакции Kafka закрывают ровно один сценарий: read-process-write внутри Kafka. Прочитали из топика A, посчитали, записали в топики B и C, и всё это либо видно потребителям целиком, либо не видно вовсе.
Механика по шагам:
- Продюсер задаёт
transactional.id, стабильный идентификатор, переживающий рестарт (например,billing-worker-p3, привязанный к партиции, а не случайный UUID). initTransactions()находит transaction coordinator и получает PID с новым epoch. Все старые продюсеры с тем жеtransactional.idи меньшим epoch получаютProducerFenced. Так работает zombie fencing: подвисший старый инстанс не сможет дописать в транзакцию после того, как его заменили.- Записи внутри транзакции пишутся в целевые партиции сразу, но помечены как транзакционные.
sendOffsetsToTransaction(offsets, consumer.groupMetadata())кладёт коммит входных offset в ту же транзакцию, и связка «прочитал и записал» становится атомарной.commitTransaction(): координатор пишет решение в__transaction_stateи рассылает во все затронутые партиции control batch (маркер COMMIT или ABORT).
На стороне чтения нужен isolation.level=read_committed. Такой консьюмер видит
записи только до LSO (last stable offset, offset первой ещё не завершённой транзакции)
и отфильтровывает записи прерванных транзакций по маркерам. А по умолчанию стоит
read_uncommitted, отсюда классический подвох «включили транзакции, а дубликаты и
грязные данные остались»: консьюмера просто забыли переключить.
Как только обработчик делает что-то за пределами 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Снизу вверх
- Сегмент складывается из файла
.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{}. Возьмёшь не тот — ключ вроде бы задан,
а сообщения одного заказа расползаются по партициям относительно соседнего сервиса,
и «порядок гарантирован» превращается в тыкву. Проверяется это за минуту: прогнать
сто сообщений с одним ключом и посмотреть распределение по партициям.
Кто и как распределяет
Координатор группы (брокер, выбранный по хешу 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,
сохраняя порядок на ключ. Так пробивают потолок, не меняя топик.
Где хранится
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.
«А если консьюмер работает в несколько горутин и коммитит по завершении каждой?» Тогда можно закоммитить offset 50, пока обработка 47 ещё идёт, и падение потеряет 47–49. При параллельной обработке коммитить надо сам минимальный незавершённый offset: он и есть следующая запись для чтения; держать map незавершённых и двигать «водяной знак» только непрерывно. Такое проговорит только тот, кто реально писал консьюмеры.
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прыгнул в конец, молча потеряв данные. Это инцидент, а не самоизлечение.
Чем мерить
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), явно зафиксировав,
что часть событий не будет обработана. Последнее уже управленческое решение, а не
техническое.
Почему «просто не коммитить» не работает
Партиция встаёт колом на одном ядовитом сообщении — 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 это всегда пишут руками, и такой код стоит вынести в общую библиотеку на всю компанию.
Механика
- Продюсер задаёт стабильный
transactional.id, переживающий рестарт. initTransactions()находит transaction coordinator и получает PID с новым epoch. Старые продюсеры с тем же id и меньшим epoch получаютProducerFenced, и zombie fencing отсекает подвисший предыдущий инстанс.- Записи пишутся в партиции сразу, но помечены транзакционными.
sendOffsetsToTransaction(offsets, consumer.groupMetadata())включает коммит входных offset в ту же транзакцию; на этом и держится атомарность связки «прочитал и записал».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, куда стекается вся «непонятная почта».
Механика routing key на примерах
Чаще всех работает topic: почти в любом проекте хватает одного topic-обменника на домен.
Ключ состоит из слов через точку, это иерархия длиной до 255 байт. Спецсимволов всего два:
* заменяет ровно одно слово, # — ноль или больше слов.
| Binding key | order.created | order.created.eu | order.paid.eu.b2b | user.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-соединения.
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 не ограничен, и хуже дефолта не придумать: первый подключившийся консьюмер выгребает всю очередь себе в память, остальные простаивают.
Как подобрать значение? Каждое подтверждение стоит сетевого round-trip, поэтому
prefetch=1 даёт идеально ровное распределение и низкий throughput: консьюмер
простаивает, пока ack идёт туда, а следующее сообщение обратно. Ориентиры такие:
| Профиль задачи | prefetch | Почему |
|---|---|---|
| Долгая обработка (секунды и минуты): конвертация видео, отчёт, вызов внешнего API | 1 | Round-trip в единицы миллисекунд на фоне секунд не виден, зато задача не «залипнет» за спиной другой |
| Быстрая обработка (единицы миллисекунд), много мелких сообщений | 100–300 | Иначе сеть станет узким местом раньше, чем CPU |
| Пул из N воркеров внутри одного процесса | N или 2N | Держать все воркеры занятыми и один-два в очереди на вход, не больше |
| Неизвестно | 10–50 | Безопасный старт: и память ограничена, и сеть не пилит round-trip на каждое |
Память тут не единственная цена. Сообщения из буфера консьюмера уже не будут перераспределены: их не увидит новый под, который ты только что поднял под пик. При 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 | Что произошло | Типичная причина в коде |
|---|---|---|
rejected | basic.reject или basic.nack с requeue=false | Консьюмер решил, что ретрай не поможет |
expired | Истёк TTL — сообщения (expiration) или очереди (x-message-ttl) | Задержка в retry-очереди или протухшие данные |
maxlen | Очередь упёрлась в x-max-length / x-max-length-bytes | Перегрузка; при overflow=drop-head вытесняется самое старое |
delivery_limit | Quorum-очередь: превышен x-delivery-limit перед доставкой | Poison message, брокер сам отправляет его в DLX |
При каждом прохождении брокер дописывает запись в заголовок x-death. Там лежит
массив структур с полями queue, reason, count,
exchange, routing-keys, time. Другого
штатного счётчика попыток в AMQP нет: x-death[0].count честно говорит,
сколько раз сообщение уже возвращалось через dead-lettering.
- 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-очередь и класть в каждое сообщение свой
expiration». Получится система, где сообщение с задержкой 1 секунда,
попавшее за сообщение с задержкой 10 минут, выйдет через 10 минут. Classic-очередь
проверяет TTL только у головы: она FIFO, и «протухшее в середине»
физически не может выйти раньше. Решений два: лестница из очередей с фиксированным TTL
(по одной на ступень) или плагин
rabbitmq_delayed_message_exchange, который держит отложенные сообщения
в Mnesia и публикует по таймеру.
Durable, persistent, confirms: четыре независимых условия
«Сообщение переживёт рестарт брокера» требует не одной галочки, а конъюнкции четырёх условий, и каждое отвечает за свой кусок пути. Провалить достаточно одно.
| Что настраиваем | Где | Что будет, если забыть |
|---|---|---|
durable=true у exchange | ExchangeDeclare | После рестарта обменника нет; публикация в него закрывает канал с NOT_FOUND |
durable=true у очереди | QueueDeclare | Определение очереди не восстановится; очередь исчезает вместе со всем содержимым |
DeliveryMode = 2 (persistent) | у каждого сообщения | Тело живёт только в памяти. Очередь durable — но пустая после рестарта |
| Publisher confirms | ch.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 {
// брокер не принял: публикуем заново, дубликаты тут норма
}
Для 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 | RabbitMQ |
|---|---|---|
| Модель | Распределённый лог с офсетами, читатель хранит позицию | Очереди + маршрутизатор, сообщение исчезает после 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Пять сущностей
- Connection — TCP+TLS-соединение до брокера, дорогое, открывается один раз на процесс.
- Channel — логический подканал внутри соединения. Все команды идут через него. Не потокобезопасен: в Go заводят по каналу на горутину. Ошибка протокола закрывает канал, и приложение должно уметь его пересоздать.
- Exchange — маршрутизатор без состояния. Четыре типа плюс безымянный
default exchange (
""), к которому каждая очередь автоматически привязана по своему имени. Поэтому «hello world» сpublish("", "my-queue", body)и работает без объявления обменника. - Queue — единственное место, где сообщение реально лежит. FIFO, у classic-типа живёт на одном узле кластера.
- Binding связывает exchange → queue через binding key (или набор заголовков). Между одной парой может быть несколько привязок с разными ключами.
Типы exchange
| Тип | Правило | Когда берут |
|---|---|---|
direct | binding 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.
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: упавший консьюмер ничего не теряет, но всё повторит.
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 с задержкой.
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.
Что следует из модели
- 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-топики закрывают часть сценариев очередей. Но модель хранения у них по-прежнему разная, и выбирать надо по ней».
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 и не хочет второй брокер.
Так и скажи, но всё равно назови разделение «core = at-most-once без хранения,
JetStream = persistence с политиками retention» и одну отличительную черту NATS:
дедупликацию по Nats-Msg-Id из коробки или queue groups
без объявления очередей. Этого хватит, чтобы закрыть вопрос: проверяют
не опыт, а понимание, зачем существует третий вариант.
8.4Паттерны асинхронной интеграции
Паттерны этой главы лечат одну и ту же беду: у нас две системы, а транзакция только в одной. База данных умеет атомарность, брокер умеет доставку, общего коммита между ними нет. Outbox переносит публикацию внутрь транзакции БД, inbox переносит дедупликацию внутрь транзакции консьюмера. Остальное касается того, что класть в сообщение и что делать, когда сообщений становится больше, чем сил их переварить.
Двойная запись — корень всех бед
Задача выглядит безобидно: сохранить заказ в PostgreSQL и опубликовать событие
OrderCreated, чтобы отреагировали склад, биллинг и аналитика. Два действия,
две разные системы, и ни одного способа сделать их атомарно. Какой порядок ни
выбери, между ними остаётся окно, в котором процесс может умереть — под убьёт OOM-killer,
нода уедет на обслуживание, сеть моргнёт.
- «Публиковать внутри транзакции». Брокер в транзакции БД не участвует:
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()
}
Соблазнительно вместо флага 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 в рабочую
базу.
| Polling | CDC / Debezium | |
|---|---|---|
| Нагрузка на БД | Постоянные запросы + UPDATE каждой строки (второй проход по heap, больше WAL) | Чтение WAL, обычной нагрузки почти нет |
| Задержка | Половина интервала опроса в среднем | Миллисекунды |
| Порядок | Внутри одного релея; между несколькими — не гарантирован | Строго порядок коммитов в WAL |
| Инфраструктура | Ничего лишнего: горутина в сервисе | Kafka Connect кластер, конфиги коннекторов, мониторинг слотов |
| Схема | Полный контроль над телом сообщения | Нужен Outbox Event Router SMT, чтобы вытащить payload и поставить key/topic |
| Риск | Забыть уборку → распухшая таблица | Слот не читается → WAL не удаляется → диск БД кончается |
С 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 и на лаг слота в байтах.
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) отметка и результат становятся видимыми вместе
}
Детали, на которых валятся
- Ключ дедупликации. Лучше всего подходит
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, партиции распределяются между участниками | Число партиций: лишние консьюмеры простаивают |
| NATS | Queue 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(превысишь — консьюмера исключат из группы и запустят ребалансировку), в JetStreamAckWait, в SQS продление visibility timeout. Долгие задачи лучше дробить или уносить в отдельный поток с большими таймаутами, чтобы не тащить эти настройки на весь сервис.
Fanout / multicast и job distribution
Эти два режима постоянно путают, хотя разница ровно в одном: сколько раз сообщение будет обработано. При fanout каждый подписчик получает свою копию (широковещание). При job distribution оно достаётся ровно одному из группы (конкуренция). Оба режима даёт одна и та же инфраструктура, и ошибка в конфигурации молча превращает один в другой.
| Система | Fanout (каждому копия) | Job distribution (одному из группы) |
|---|---|---|
| RabbitMQ | fanout/topic exchange + своя очередь на каждого подписчика | Одна очередь + несколько консьюмеров |
| Kafka | Разные group.id у каждого сервиса | Один group.id на все инстансы сервиса |
| NATS | Обычные подписчики на subject | Queue 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 | Event-carried state transfer | Event sourcing | |
|---|---|---|---|
| Что в сообщении | Идентификатор и тип факта | Полное состояние сущности на момент события | Все изменения; события — сами по себе источник истины |
| Состояние живёт | Только в сервисе-источнике | Копия у каждого потребителя | Восстанавливается свёрткой лога, плюс снапшоты |
| Размер | Сотни байт | Килобайты и десятки килобайт | Мелкие события, но их много и хранятся долго |
| Связанность | Обратная рантайм-зависимость: нужен живой источник | Зависимость от схемы: поменял поле — сломал всех | Схема событий вечна: старые записи придётся читать всегда |
| Источник недоступен | Обработка встаёт | Потребитель продолжает работать | Не применимо: лог и есть источник |
| Согласованность | Читаем свежее, но можно поймать гонку «событие быстрее реплики» | Данные слегка устаревшие, зато согласованные с событием | Полная история, любые temporal-запросы |
| Плюс | Тонкий контракт, нет дублирования | Автономия, устойчивость, легко строить проекции | Аудит, replay, пересборка проекций, отладка «как дошли до такого» |
| Минус | N+1 вызовов, доступность перемножается | Трафик, копии данных, PII расползается | Высокая сложность, версионирование навсегда, тяжело удалять данные |
| Когда берут | События редкие, данные крупные, потребителей мало | Дефолт для интеграции доменов; кэши и локальные реплики | Домены, где ценна история: деньги, склад, юридически значимые действия |
У event notification есть неочевидный отказ: сервис orders закоммитил
заказ, опубликовал событие — и shipping тут же пошёл за деталями через
GET /orders/42, а балансировщик отправил запрос на реплику для
чтения, куда изменения ещё не доехали. В ответ приходит 404 на объект, о создании
которого мы только что получили событие. Лечат это чтением с мастера для этого сценария,
ретраем с бэкоффом, номером версии в событии с опросом «пока version < 7 — жди» или
переходом на state transfer, где данные уже лежат в сообщении.
Стоит предложить компромисс: тонкое событие плюс самые нужные поля и номер версии. Большинству потребителей хватает того, что пришло, а тем, кому нужны детали, остаётся один вызов вместо всех. Для больших полезных нагрузок есть claim check: в событии ссылка на S3, а тело лежит там.
Публикация событий в брокер сама по себе не делает систему 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 у консьюмера — нет |
| Сделать поле обязательным по смыслу | Нет | Старые продюсеры его не заполняют; валидация уронит консьюмера |
Поле 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). Она сразу переводит разговор из «надо что-то делать» в «у нас четыре часа».
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 или ждут БД: если ждут, новые поды сделают хуже. Дальше по порядку: масштабирование до числа партиций, оптимизация и батчинг, увеличение партиций (с пониманием, что это ломает порядок по ключу и необратимо), троттлинг продюсера и в крайнем случае осознанное отбрасывание того, что не жалко. Буфер брокера при этом не решение, а время на то, чтобы решение применить». Такой ответ показывает, что кандидат думает о системе, а не об одной ручке.
Вопросы
6COMMIT,
а отдельный релей выносит его в брокер. Гарантия — 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, и диск базы заканчивается.
Потому что 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.
Ещё платишь ростом таблицы (нужна чистка партициями или пачками), лишней записью в каждой
бизнес-транзакции и ещё одним процессом, за которым надо следить.
Главный алерт всей конструкции — метрика «возраст самой старой неопубликованной
строки».
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, чтобы эффект уходил наружу уже из релея. RedisSETNXс TTL работает быстро, но не атомарно с бизнес-работой: упали между SETNX и работой — сообщение считается обработанным, хотя ничего не произошло.
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: одному |
|---|---|---|
| RabbitMQ | fanout/topic exchange + своя очередь на подписчика | Одна очередь, несколько консьюмеров |
| Kafka | Разные group.id | Один group.id |
| NATS | Обычные подписчики на subject | Queue group с общим именем |
- Три сервиса подписаны на fanout через одну общую очередь. Каждое событие достаётся одному случайному из них. Выглядит это как «иногда письмо не приходит» и отлаживается неделями. Очередей нужно столько же, сколько подписчиков.
- Одинаковый
group.idу разных сервисов в Kafka. Зеркальная беда: сервисы делят партиции и видят по половине потока. Вgroup.idпишут имя сервиса, а не имя топика и не «app».
И про стоимость. В Kafka веерная рассылка бесплатна по хранению: N групп читают один лог, копия одна. RabbitMQ при fanout физически размножает сообщение: десять подписчиков дают десять копий в памяти и на диске. Если сообщения крупные, помогает claim check: в событие кладут ссылку на объект в S3, а не сам объект.
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 в событии дёшево страхует от гонок и от применения событий не
по порядку.
reserved), Avro — на именах
с обязательными default. Schema Registry не столько хранит схемы, сколько
не даёт зарегистрировать несовместимую. BACKWARD — обновляем консьюмеров первыми,
FORWARD — продюсеров, FULL — порядок не важен.Protobuf
Имена полей в wire-формате не передаются, передаются номера. Отсюда и все правила.
Поле с новым номером добавлять безопасно (старый консьюмер положит его в unknown
fields). Переименование безопасно для бинарного формата, но ломает JSON-маппинг.
Удалять поле можно только вместе с 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, а не в реестре постфактум.
Диагностика
Растущий лаг сам по себе не диагноз. Считаем скорость разгребания = 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 зависает даже у тех, кто не создавал нагрузки.