MBL
Go / avito / Тестовое задание / Kafka и RabbitMQ: consumer groups, ребалансировка, exactly-once
Go сложный

Kafka и RabbitMQ: consumer groups, ребалансировка, exactly-once

kafkarabbitmqmessaging

Продолжение статьи «Защита тестового 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

Commit до обработки (at-most-once)
func consumeLoop(c *kafka.Consumer) {
	for {
		msg := c.Poll()
		c.CommitOffset(msg.Offset) // подтвердили ДО обработки

		if err := process(msg); err != nil {
			// сообщение уже подтверждено — повторить попытку негде,
			// при падении процесса ЗДЕСЬ сообщение теряется навсегда
			log.Error(err)
		}
	}
}
Commit после обработки (at-least-once)
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 в общем случае нельзя.

Самопроверка 0 / 7
Могу объяснить, почему ребалансировка партиций — это пауза обработки, а не бесшовное переключение
Знаю, что Kafka из коробки даёт at-least-once, а не exactly-once, и почему
Понимаю, как transactional outbox + идемпотентный consumer вместе эмулируют "эффективно ровно один раз"
Помню, что порядок сообщений в Kafka гарантирован только внутри одной партиции (одного ключа), не по топику целиком
Знаю, почему увеличение числа партиций уже созданного топика может сломать гарантию порядка по ключу
Могу навскидку сказать, когда уместнее RabbitMQ, а когда Kafka
Понимаю разницу между коммитом offset до и после обработки сообщения и какой риск несёт каждый вариант
Как усвоено?