Kafka и RabbitMQ: consumer groups, ребалансировка, exactly-once
Продолжение статьи «Защита тестового Avito на реальном проекте: Avito.Кухня» — там уже разобраны transactional outbox, DLQ и идемпотентность на уровне одного конкретного проекта. Здесь — то, что осталось за кадром: как Kafka устроена под капотом на уровне consumer group и партиций, и чем в этом смысле отличается RabbitMQ.
Consumer group и ребалансировка партиций
Проще всего представить так: есть один писатель (продюсер, кладёт сообщения в топик) и группа читателей (consumer group), которые вместе разбирают этот топик. Партиции топика распределяются между читателями одной группы так, что каждую партицию в любой момент читает ровно один читатель этой группы — не два, не ноль.
Именно это даёт горизонтальное масштабирование чтения: добавил ещё одного читателя в группу — партиций на каждого стало меньше, разбирают быстрее. Group coordinator (один из брокеров) следит, какие читатели ещё живы, через heartbeat — короткий сигнал "я ещё здесь, партиции у меня забирать не надо".
Ребалансировка — это момент, когда координатор заново раздаёт партиции читателям, потому что состав группы изменился: один читатель упал (перестал слать heartbeat), пришёл новый читатель, или кто-то сам вышел из группы.
На классической (eager) схеме ребалансировки на время этого пересчёта все читатели группы останавливают обработку и сдают свои партиции разом — это буквально stop-the-world пауза для ВСЕЙ группы, а не только для того читателя, из-за которого всё началось.
Более новые протоколы (incremental cooperative rebalancing) уменьшают эту паузу, отбирая только те партиции, которые реально должны переехать к другому читателю, но полностью убрать паузу для самих переезжающих партиций всё равно нельзя — кто-то должен закончить читать партицию, прежде чем её отдать.
Почему это важно для at-least-once. Читатель мог обработать сообщение (например, применить переход статуса заказа в БД), но не успеть закоммитить offset до того, как начался ребаланс и партиция перешла к другому читателю. Новый читатель партиции продолжит с последнего закоммиченного offset — то есть то же самое сообщение придёт на обработку повторно, но уже другому читателю, который ничего не знает про попытку первого.
Это ещё один источник дублей поверх обычного at-least-once, и он не лечится ничем, кроме идемпотентности на стороне обработки (ровно то, что уже описано в статье про outbox — SELECT ... FOR UPDATE + no-op при повторе того же статуса).
Exactly-once, at-least-once, at-most-once
Три режима гарантии доставки, различающиеся тем, что происходит на границе "сбой ровно между чтением и подтверждением":
- At-most-once — offset коммитится до обработки сообщения. Если процесс упадёт между коммитом и обработкой, сообщение потеряно навсегда: следующий читатель продолжит с уже сдвинутого offset и это сообщение больше никогда не увидит.
- At-least-once — offset коммитится после успешной обработки. Если процесс упадёт после обработки, но до коммита (или партиция уйдёт другому читателю при ребалансе, см. выше), сообщение обработается повторно. Потерь нет, зато возможны дубли.
- Exactly-once — сообщение гарантированно обработано ровно один раз, без потерь и без дублей.
Обычный цикл consumer'а Kafka (consume → обработать → commit offset) даёт at-least-once, не exactly-once — потому что "обработать" и "закоммитить" это два раздельных действия, между которыми процесс может упасть, а откатить транзакционно сам факт обработки (если это, например, запись в другую БД, а не в саму Kafka) Kafka не умеет: у неё нет представления о том, что произошло снаружи её собственного протокола.
У Kafka есть встроенные transactions (идемпотентный продюсер + transactional API), которые дают exactly-once, но только для цепочки "consume из одного топика → produce в другой топик Kafka" внутри одного и того же кластера — как только в цепочке появляется внешняя система (Postgres, HTTP-вызов), настоящий exactly-once перестаёт существовать в принципе: нельзя атомарно закоммитить транзакцию в Postgres и транзакцию в Kafka одной операцией, это два разных ресурса.
Как transactional outbox эмулирует "эффективно ровно один раз" без настоящих Kafka-транзакций. Запись в outbox происходит в той же БД-транзакции, что и бизнес-изменение — значит на стороне отправителя either both happen or neither (атомарность там, где она реально нужна: между "заказ создан" и "событие поставлено на отправку").
На стороне получателя Kafka всё равно даёт at-least-once (с учётом дублей из ребалансировки), но consumer идемпотентен — повторная обработка того же сообщения не меняет результат. Сумма "at-least-once с одной стороны + идемпотентность с другой" даёт тот же наблюдаемый эффект, что exactly-once (сообщение гарантированно применится, и применится только один раз по сути), не будучи им технически.
Партиционирование по ключу
Продюсер выбирает партицию для сообщения хэшированием ключа (hash(key) % количество_партиций) — при отсутствии ключа сообщения распределяются round-robin или случайно. Kafka гарантирует порядок доставки строго в пределах одной партиции: если два сообщения с одним и тем же ключом (например order_id) отправлены в этом порядке, consumer гарантированно прочитает их в том же порядке, потому что оба физически лежат в одной партиции и партиция — это упорядоченный лог. Порядок между сообщениями с разными ключами (то есть, как правило, в разных партициях) не гарантирован вообще — параллельные consumer'ы разных партиций читают независимо.
Почему увеличение числа партиций уже созданного топика опасно. Формула hash(key) % N меняет результат при изменении N — после увеличения количества партиций тот же order_id, что раньше всегда попадал в партицию 3, может начать попадать в партицию 7.
Все сообщения, отправленные до смены N, остаются на старых местах, а новые с тем же ключом идут уже в другую партицию — гарантия "все сообщения одного заказа лежат в одной партиции" на этой границе ломается. Практический вывод: число партиций топика планируется заранее с запасом, а не увеличивается по факту роста нагрузки, если порядок по ключу критичен.
RabbitMQ как альтернатива
RabbitMQ устроен принципиально иначе: продюсер публикует сообщение в exchange, а не напрямую в очередь, exchange по правилам маршрутизации решает, в какие очереди сообщение попадёт:
- direct — точное совпадение routing key с binding key очереди (один получатель на конкретный ключ).
- topic — routing key матчится по маске с wildcard'ами (
orders.*,orders.#) — гибкая маршрутизация по паттерну. - fanout — сообщение уходит во все привязанные очереди без учёта routing key — классический broadcast.
Подтверждение обработки — ack/nack на уровне отдельного сообщения (в отличие от Kafka, где подтверждение — это сдвиг offset, всегда "всё до этой точки"): consumer может nack конкретное сообщение и вернуть его в очередь или в dead-letter exchange, не трогая соседние сообщения.
Это и есть источник разницы в порядке: RabbitMQ не хранит сообщения как упорядоченный лог партиций — после nack и повторной постановки в очередь сообщение легко может обработаться позже сообщений, отправленных после него, тогда как Kafka-партиция такого не допускает в принципе (сообщения читаются строго по offset).
Когда уместнее RabbitMQ, а не Kafka:
- невысокий объём сообщений, где не нужна историческая ретенция лога (Kafka хранит сообщения заданное время/объём независимо от того, прочитаны они или нет — это осознанная модель "лог", а не "очередь"; RabbitMQ же по умолчанию удаляет сообщение после успешного ack, ближе к классической модели очереди);
- нужна сложная маршрутизация по паттернам (topic exchange) или широковещательная рассылка (fanout) без ручной реализации поверх партиций;
- per-message контроль подтверждения важнее строгого порядка.
Когда уместнее Kafka:
- высокий устойчивый объём сообщений;
- нужен строгий порядок в пределах ключа;
- нужна возможность консьюмеру перечитать историю с произвольного offset (replay);
- несколько независимых consumer group читают один и тот же поток данных для разных целей.
Где именно коммитить offset
func consumeLoop(c *kafka.Consumer) {
for {
msg := c.Poll()
c.CommitOffset(msg.Offset) // подтвердили ДО обработки
if err := process(msg); err != nil {
// сообщение уже подтверждено — повторить попытку негде,
// при падении процесса ЗДЕСЬ сообщение теряется навсегда
log.Error(err)
}
}
}func consumeLoop(c *kafka.Consumer) {
for {
msg := c.Poll()
if err := process(msg); err != nil {
// offset не сдвинут — при рестарте сообщение придёт снова
log.Error(err)
continue
}
c.CommitOffset(msg.Offset) // подтвердили ПОСЛЕ успешной обработки
}
}Левый вариант рискует потерять сообщение (упал между коммитом и process), правый рискует обработать его дважды (упал между process и коммитом, или партиция ушла на ребалансе) — но не теряет данных. Для бизнес-событий вроде "заказ создан" почти всегда выбирают правый вариант и решают проблему дублей идемпотентностью на стороне обработки, а не пытаются устранить дубли на стороне доставки — устранить их там средствами голой Kafka в общем случае нельзя.