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

Откуда берутся дубли в логе, если продюсер отправил сообщение один раз, и как их убирает enable.idempotence.

Дубли в Kafka бывают двух видов, и их часто путают. Первый — дубли обработки: консьюмер упал до коммита, новый владелец перечитал записи заново. Второй — дубли в самом логе, и их создаёт продюсер.

Механика простая. Продюсер отправил запись, брокер её принял, но подтверждение не вернулось вовремя — потерялось в сети или не уложилось в request.timeout.ms. Продюсер считает отправку неудачной и повторяет её. В логе оказывается две копии, консьюмер честно обработает обе. Это не ошибка приложения: так ведёт себя любой ретрай поверх ненадёжной сети.

enable.idempotence=true закрывает эту дыру. Продюсер получает идентификатор и нумерует записи по каждой партиции, брокер помнит последний принятый номер и молча отбрасывает повтор, отвечая продюсеру успехом. С Kafka 3.0 это включено по умолчанию (KIP-679) вместе с acks=all.

Что идемпотентность не даёт: это не exactly-once. Гарантия действует на отрезке «продюсер → брокер» и в пределах сессии продюсера. Путь от брокера до вашей базы данных она не покрывает — там по-прежнему нужен идемпотентный обработчик или транзакции.

Если в вашем проекте стоит старый клиент или кто-то выключил идемпотентность ради «производительности», вы платите дублями зря: накладные расходы у неё минимальные.

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

Частые вопросы

Идемпотентность замедляет продюсера?

Практически нет. Ограничение на max.in.flight.requests.per.connection ≤ 5 сохраняет пакетирование и порядок.

Работает ли идемпотентность между перезапусками?

Нет, она действует в пределах сессии продюсера. Для сквозной гарантии нужны транзакции с постоянным transactional.id.

Нужен ли всё равно идемпотентный обработчик?

Да. Дубли на стороне консьюмера возникают при ребалансировке и падениях независимо от настроек продюсера.

Дальше по теме

Что такое партиция в Kafka Consumer lag: что это и как с ним работать acks в Kafka: 0, 1 и all ISR и min.insync.replicas Ребалансировка consumer group