Exactly-Once Semantics (EOS) в Apache Kafka

Гарантия Exactly-Once Semantics (EOS) в Apache Kafka — это «святой Грааль» распределённых систем. Она означает, что результат обработки сообщения запишется в систему строго один раз, даже если вокруг падают серверы, рвётся сеть или дублируются запросы.

Академически гарантировать 100% доставку строго одного пакета по сети в любых условиях невозможно. Классическая иллюстрация этой проблемы — задача двух генералов. Поэтому Kafka решает задачу иначе: физические повторные отправки (retries) могут происходить, но их бизнес-эффект применяется ровно один раз.

Чтобы это сработало, Kafka объединяет два независимых механизма:

  1. Идемпотентный продюсер.
  2. Транзакции.

Шаг 1. Идемпотентный продюсер: борьба с дубликатами

Представьте, что приложение отправляет в Kafka сообщение о переводе денег. Брокер сохранил его у себя на диске, но в этот же миг моргнула сеть, и ответное подтверждение (ACK) до продюсера не дошло.

Продюсер думает: «Сбой!» — и отправляет сообщение повторно. Без идемпотентности в топике оказалось бы две записи о списании.

Как Kafka решает это при включённом параметре:

enable.idempotence=true

1. Producer ID (PID)

При старте брокер выдаёт продюсеру уникальный идентификатор — Producer ID, или PID.

2. Sequence Number

Для каждого сообщения продюсер назначает порядковый номер:

0, 1, 2, 3, ...

Этот номер привязан к конкретной партиции.

Когда брокер получает сообщение, он смотрит на пару:

PID=42, Seq=10

Если брокер видит, что запись с PID=42 и Seq=10 уже сохранена, он просто игнорирует дубликат и отправляет продюсеру успешный ответ:

Всё окей, я это уже записал.

Данные не дублируются.


Шаг 2. Транзакции в Kafka: атомарность

Идемпотентность защищает от дублей при записи в рамках продюсера и партиций. Но в реальной жизни часто нужно сделать цепочку действий:

  1. Прочитать сообщение из одного топика, например payment-requests.
  2. Обработать его.
  3. Отправить результат в другой топик, например notifications.
  4. Закоммитить прочитанный offset.

Нужно, чтобы эти действия произошли атомарно:

Либо всё записалось вместе, либо ничего.

Для этого используются транзакции Kafka, появившиеся в рамках KIP-98: Exactly Once Delivery and Transactional Messaging.

Процесс устроен следующим образом.

1. Инициализация транзакции

Продюсер инициализирует транзакцию на сервере.

Ключевой параметр:

transactional.id=my-unique-app-id

transactional.id — это уникальное фиксированное имя приложения, по которому Kafka связывает разные запуски одного и того же transactional producer.

2. Отправка сообщений

Продюсер отправляет сообщения в разные топики.

На брокере эти сообщения помечаются как неподтверждённые, то есть uncommitted. В этот момент обычные потребители с корректной настройкой изоляции их ещё не видят.

3. Добавление offset в транзакцию

Продюсер отправляет в транзакцию ещё и прочитанный offset.

Это возможно, потому что __consumer_offsets — тоже внутренний топик Kafka.

4. Завершение транзакции

Продюсер вызывает:

commitTransaction();

Специальный серверный компонент — Transaction Coordinator — ставит на всех сообщениях и offset маркер Committed.

Если приложение упадёт до завершения транзакции, координатор по таймауту выполнит abort. Отправленные сообщения будут проигнорированы потребителями с read_committed, а offset не будет считаться успешно зафиксированным в рамках этой транзакции.


Что нужно настроить разработчику, чтобы включить Exactly-Once

EOS нельзя получить бесплатно. Его нужно явно активировать в конфигурации компонентов.

1. Producer

enable.idempotence=true
transactional.id=my-unique-app-id
acks=all

Что это означает:

  • enable.idempotence=true включает уникальные номера сообщений и защиту от дублей при retry.
  • transactional.id включает механизм транзакций и позволяет Kafka идентифицировать приложение между перезапусками.
  • acks=all заставляет продюсер ждать подтверждения от всех необходимых реплик для максимальной надёжности записи.

2. Consumer

isolation.level=read_committed

По умолчанию потребители могут читать неподтверждённые транзакционные данные. Настройка read_committed заставляет consumer выдавать бизнес-логике только те записи, транзакции которых официально подтверждены.

3. Kafka Streams

Если используется Kafka Streams, значительная часть работы автоматизирована.

Достаточно настроить:

processing.guarantee=exactly_once_v2

В более старых конфигурациях можно встретить:

processing.guarantee=exactly_once

Kafka Streams сам настроит необходимые транзакции, идемпотентность и работу с offset.


Обратная сторона медали: плата за надёжность

Интервьюеры любят спрашивать:

Если Exactly-Once такой крутой, почему не включать его всегда по умолчанию?

Ответ — latency и накладные расходы.

Основные причины:

  • Transaction Coordinator тратит время на управление транзакциями и отправку маркеров Commit / Abort.
  • Потребители с isolation.level=read_committed вынуждены ждать завершения транзакций, что увеличивает end-to-end latency.
  • EOS работает строго в границах Kafka-транзакций. Классический сценарий — «прочитали из Kafka → записали в Kafka → закоммитили offset».
  • Если consumer внутри обработки отправляет REST-запрос во внешнюю CRM или пишет данные в стороннюю базу данных, например PostgreSQL, Kafka не сможет откатить этот внешний эффект при сбое.

Для внешних систем всё равно нужны дополнительные механизмы:

  • идемпотентные операции;
  • уникальные ключи;
  • outbox pattern;
  • ручная дедупликация;
  • transactional sink connector, если он доступен и подходит под задачу.

Короткая формула для собеседования

Exactly-Once в Kafka — это не магическая доставка одного сетевого пакета ровно один раз.

Это гарантия того, что при корректной настройке идемпотентного producer, транзакционного producer и consumer с read_committed результат обработки внутри Kafka будет виден ровно один раз, даже если физически были retries, сбои и повторные попытки.


Важное продолжение: Zombie Fencing

Следующая связанная тема — Zombie Fencing.

Kafka использует transactional.id, producer epoch и координацию транзакций, чтобы «отстреливать» зависшие старые инстансы приложения. Это защищает систему от ситуации, когда старый producer внезапно оживает и пытается продолжить транзакцию, хотя вместо него уже работает новый экземпляр приложения.

Прокрутить вверх