Exactly-once и транзакции в Kafka

Как Kafka даёт exactly-once: идемпотентный producer, транзакции с transactional.id, атомарный consume-process-produce через sendOffsetsToTransaction и consumer read_committed — и где граница EOS (внешние сайд-эффекты), которую закрывает outbox

«Exactly-once» в распределённой системе звучит как нарушение законов физики, и во многом так и есть — но Kafka даёт его в чётко очерченных границах: атомарная запись в топики Kafka плюс сдвиг offset в одной транзакции. Это покрывает классический паттерн «прочитал → обработал → записал» без дублей и потерь. А вот всё, что выходит за пределы Kafka — запись во внешнюю БД, вызов API, — под EOS уже не попадает, и здесь начинается outbox.

Это пятая статья серии «Kafka: глубокое погружение». Опирается на идемпотентный producer, разобранный в статье про репликацию и надёжность, и на модель коммита offset из статьи про consumer groups; сводит механику EOS, разобранную с прикладной стороны в статье про транзакции в брокерах. Все числа ниже — из живого трёхброкерного KRaft-кластера, включая настоящий SIGKILL процесса-клиента посреди транзакции, не симуляцию. Код этого стенда открыт в репозитории digital-cookbook — kafka/go/eos (franz-go) и kafka/java/eos (kafka-clients).

Exactly-once и транзакции: абортнутые записи физически есть, но невидимы для read_committed

В статье

Идемпотентный producer: фундамент, а не exactly-once сам по себе

Идемпотентный producer (enable.idempotence=true) закрывает один конкретный дефект: дубль от retry поверх ненадёжной сети — producer отправил запись, ответ потерялся, клиент решил, что попытка провалилась, и повторил её, хотя запись уже физически на диске. Брокер выдаёт каждому producer’у PID (producer ID) и ведёт sequence-номер на партицию; повтор с уже виденным номером отбрасывается как no-op. Механика подробно разобрана в статье про репликацию и надёжность — здесь важно только одно ограничение: идемпотентность действует строго в пределах одной сессии producer’а на одну партицию. Она не решает дедупликацию между несколькими producer-инстансами, не даёт атомарности между несколькими партициями и не связывает запись с offset’ом, который её породил. Это фундамент, на котором строятся транзакции, а не сама exactly-once semantics.

Транзакции: transactional.id и жизненный цикл

Транзакция в Kafka — это способ атомарно закоммитить (или откатить) произвольный набор записей сразу в несколько партиций и топиков, включая специальный случай — offset потребления как ещё одну «запись» в служебном топике __consumer_offsets. За это отвечает transactional.id: стабильный идентификатор producer’а, который переживает перезапуски процесса и позволяет брокеру узнать «это тот же логический производитель» даже после краха.

Именно здесь между клиентскими библиотеками проявляется реальная асимметрия — не в гарантиях, а в форме API. В kafka-clients (Java) producer.initTransactions() — обязательный явный вызов, который делается один раз при старте: он инициализирует (или восстанавливает) producer ID и, что критично, фенсит любую зависшую транзакцию предыдущего процесса с тем же transactional.id — процедура, описанная в KIP-98. Без этого вызова beginTransaction() просто упадёт с ошибкой.

props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, txnId);
// С transactional.id заданным kafka-clients берёт idempotence=true по
// дефолту; явная попытка enable.idempotence=false + transactional.id
// падает с ConfigException("Cannot set a transactional.id without also
// enabling idempotence.").
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// initTransactions() ОБЯЗАТЕЛЕН — фенсит зависшую транзакцию предыдущего
// процесса с тем же transactional.id (KIP-98).
producer.initTransactions();

В franz-go (Go) отдельного вызова, аналогичного initTransactions(), нет вовсе — это подтверждено чтением исходников pkg/kgo/txn.go, а не по памяти. BeginTransaction() сама вызывает maybeRecoverProducerID() внутри себя при первом обращении, и та же фенсинг-процедура (InitProducerIdRequest под капотом) отрабатывает синхронно, до того как BeginTransaction() вернёт управление. Разница чисто в форме API — framework-level протокол один и тот же, franz-go просто не выделяет инициализацию в отдельный публичный метод.

cl, err := kgo.NewClient(
    kgo.SeedBrokers(seeds...),
    kgo.TransactionalID(txnID),
    kgo.TransactionTimeout(30*time.Second),
    kgo.RequiredAcks(kgo.AllISRAcks()),
    // Идемпотентность НЕЛЬЗЯ выключить при заданном TransactionalID —
    // kgo.DisableIdempotentWrite() вместе с TransactionalID(...) вернёт
    // ошибку конфигурации ("cannot both disable idempotent writes and
    // use transactional IDs"). Отдельного initTransactions() нет —
    // BeginTransaction() восстанавливает producer id сама.
)

Обе библиотеки при этом симметрично запрещают отключить идемпотентность при заданном transactional.id — это не асимметрия, а общий инвариант протокола: транзакционный producer — надмножество идемпотентного, а не независимая настройка, которую можно случайно забыть включить. Реальная асимметрия между клиентами лежит не здесь, а в автокоммите потребителя (раздел про consume-process-produce ниже).

Жизненный цикл одной транзакции — beginTransaction() → одна или несколько записей → commitTransaction() либо abortTransaction() при сбое. На живом стенде eos/txn.go и eos/Txn.java гоняют три батча по 5 записей подряд одним и тем же транзакционным producer’ом: батч A коммитится штатно, батч B на первой попытке падает посередине (симулированный сбой обработки на третьей записи) и абортится, затем batch B пересылается заново и коммитится успешно.

Живой прогон: 13 физически, 10 логически

Абортнутая транзакция — не «как будто её не было». Записи batch B, отправленные до сбоя, физически уходят на брокер и занимают место в логе партиции — транзакция помечает их control-record’ом ABORT, но не стирает с диска. Именно эта разница между «физически на диске» и «логически видно» — прямое доказательство того, что изоляция транзакций Kafka работает через фильтрацию на чтении, а не через откат записи назад.

[txn] батч batchA: commit — 5 сообщений подтверждено
[txn] батч batchB: сбой в середине -> abort — 3 записи физически ушли, логически ничего
[txn] батч batchB: commit — 5 сообщений подтверждено
[txn] физически отправлено: 13, логически подтверждено: 10
[txn-verify] read_committed=10 read_uncommitted=13

Арифметика: batch A — 5 записей, коммит. Batch B, попытка 1 — обрыв на третьей записи (индексы 0, 1, 2 уже отправлены), abort; это 3 записи физически на диске и 0 логически. Batch B, попытка 2 — штатные 5, коммит. Физически: 5 + 3 + 5 = 13. Логически: 5 + 5 = 10. Consumer с isolation.level=read_committed насчитывает ровно 10, а с read_uncommitted — ровно 13, побайтово идентично на Go (franz-go) и Java (kafka-clients).

При первом живом прогоне на Java нашлась тонкость, ради которой стоит держать этот сценарий на синхронном producer’е: KafkaProducer.abortTransaction() откатывает клиентски ещё не отправленные на брокер (буферизованные) записи, а не только помечает их control-record’ом ABORT постфактум — так и написано в javadoc: «Any unflushed produce messages will be aborted». Первая версия стенда слала записи асинхронно, без .get(), и abortTransaction() успевал отменить их до отправки — итоговый read_uncommitted получался 10 вместо ожидаемых 13. Исправление — синхронный producer.send(record).get() (в Go — эквивалентный ProduceSync), который гарантирует, что запись реально дошла до брокера физически, прежде чем сценарий перейдёт к следующей. После правки 13 совпало на обоих клиентах, воспроизведено дважды.

Атомарный consume-process-produce: SIGKILL посреди транзакции

Это ядро EOS и самый содержательный сценарий стенда: прочитать из input-топика в составе consumer group, обработать, записать результат в output-топик — и сдвинуть offset потребления input — всё в одной транзакции. Если процесс падает между записью output и коммитом offset, восстановление не должно оставить систему в промежуточном состоянии: либо ничего не случилось (с точки зрения читателя output и следующего запуска той же группы), либо случилось целиком.

За коммит offset внутри транзакции в Java отвечает явный вызов producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata()) перед commitTransaction() — нужно вручную собрать Map<TopicPartition, OffsetAndMetadata> из последней прочитанной позиции по каждой партиции. В franz-go этого шага нет как отдельного публичного вызова: GroupTransactSession сама отслеживает, что было вычитано этим клиентом с последнего Begin() (через внутренние CommittedOffsets()/UncommittedOffsets()), и коммитит именно эту дельту внутри одного session.End(ctx, kgo.TryCommit) — franz-go-эквивалент sendOffsetsToTransaction, но неявный. Гарантия та же, форма API — другая.

session, err := kgo.NewGroupTransactSession(
    kgo.SeedBrokers(seeds...),
    kgo.ConsumeTopics(inputTopic),
    kgo.ConsumerGroup(group),
    kgo.TransactionalID(txnID),
    kgo.Balancers(kgo.CooperativeStickyBalancer()),
    kgo.DisableAutoCommit(), // явно; franz-go и сам не автокоммитит при заданном TransactionalID
)
session.Begin()
fetches := session.PollFetches(ctx)               // consume
// ... process ...
session.ProduceSync(ctx, outs...)                 // produce (физически на диске, транзакция ещё открыта)
committed, err := session.End(ctx, kgo.TryCommit)  // output + offset input атомарно
// enable.auto.commit=false — в отличие от franz-go, где autocommitDisable
// форсируется автоматически при заданном TransactionalID, в Java это нужно
// выставить руками: дефолт (true) может продвинуть offset группы фоновым
// автокоммитом мимо транзакции — реальный footgun, которого нет в franz-go.
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
// ...
producer.initTransactions();
var input = consumer.poll(...);                    // consume
producer.beginTransaction();
for (var r : input) producer.send(...);             // produce
producer.flush();
Map<TopicPartition, OffsetAndMetadata> offsets = /* последний offset+1 по каждой партиции */;
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();                       // output + offset input атомарно
sequenceDiagram participant C as Consumer (группа eos-cpp-group) participant P as Transactional Producer participant K as Kafka (output-топик + __transaction_state) Note over P,K: Штатный путь: consume -> process -> produce + sendOffsets в одной транзакции P->>K: initTransactions() / BeginTransaction() C->>K: poll() input-топик K-->>C: n записей P->>K: beginTransaction() loop обработка каждой записи P->>K: send(output, processed) end P->>K: sendOffsetsToTransaction(offsets, groupMetadata) P->>K: commitTransaction() K-->>C: offset input продвинут атомарно вместе с output Note over P,K: Abort-путь: SIGKILL строго между produce и commit P->>K: send(output, processed) — физически на диске, транзакция ОТКРЫТА Note over P: SIGKILL (host убивает процесс) Note over K: транзакция зависла — offset input НЕ продвинут, output НЕВИДИМ read_committed participant P2 as Producer (повтор, тот же transactional.id) P2->>K: initTransactions() -> фенсит зависшую транзакцию (KIP-98), abort P2->>K: consume-process-produce повтор (те же записи, offset не был продвинут) P2->>K: sendOffsetsToTransaction + commitTransaction K-->>C: output ровно один раз, offset продвинут

sequenceDiagram
    participant C as Consumer (группа eos-cpp-group)
    participant P as Transactional Producer
    participant K as Kafka (output-топик + __transaction_state)

    Note over P,K: Штатный путь: consume -> process -> produce + sendOffsets в одной транзакции
    P->>K: initTransactions() / BeginTransaction()
    C->>K: poll() input-топик
    K-->>C: n записей
    P->>K: beginTransaction()
    loop обработка каждой записи
        P->>K: send(output, processed)
    end
    P->>K: sendOffsetsToTransaction(offsets, groupMetadata)
    P->>K: commitTransaction()
    K-->>C: offset input продвинут атомарно вместе с output

    Note over P,K: Abort-путь: SIGKILL строго между produce и commit
    P->>K: send(output, processed) — физически на диске, транзакция ОТКРЫТА
    Note over P: SIGKILL (host убивает процесс)
    Note over K: транзакция зависла — offset input НЕ продвинут, output НЕВИДИМ read_committed
    participant P2 as Producer (повтор, тот же transactional.id)
    P2->>K: initTransactions() -> фенсит зависшую транзакцию (KIP-98), abort
    P2->>K: consume-process-produce повтор (те же записи, offset не был продвинут)
    P2->>K: sendOffsetsToTransaction + commitTransaction
    K-->>C: output ровно один раз, offset продвинут
Consume-process-produce в одной транзакции: штатный путь и abort-путь при SIGKILL до commit

Живой сценарий (ops/eos-kill.sh, топики demo-eos-cpp-input/demo-eos-cpp-output, n=10) устроен так: попытка A запускается фоновым контейнером, читает 10 записей, обрабатывает, синхронно пишет результат в output — записи физически на диске, транзакция ещё открыта — печатает маркер READY-TO-COMMIT и засыпает на 60 секунд. Host-скрипт ловит маркер в логах и убивает контейнер docker kill (SIGKILL) строго в этом окне — между produce и commit, не раньше и не позже.

[cpp-attempt txn=cookbook-eos-cpp-producer] вычитано из demo-eos-cpp-input: 10 записей
[cpp-attempt txn=...] обработано и записано в demo-eos-cpp-output: 10 записей
  (ФИЗИЧЕСКИ на диске, транзакция ещё ОТКРЫТА — офсет input ещё НЕ закоммичен)
READY-TO-COMMIT n=10
--- маркер READY-TO-COMMIT пойман — docker kill попытки A (SIGKILL, ДО commit) ---
--- exit code попытки A: 137 (SIGKILL, реальный крах) ---

verify ДО повтора:
  output read_committed=0, read_uncommitted=10(физически), committed-offset группы=0
  [assert] OK: частичного состояния нет

--- попытка B: ТОТ ЖЕ transactional.id (фенсит зависшую A), ТА ЖЕ группа (перечитывает те же 10) ---

verify ПОСЛЕ повтора:
  output read_committed=10 (РОВНО ОДИН РАЗ, без дублей),
  read_uncommitted=20(физически, 10 от A + 10 от B), committed-offset группы=10
  [assert] OK: частичного состояния нет

Exit code попытки A — 137, то есть реальный SIGKILL, не мягкое завершение. Verify до повтора показывает именно то, что должна показывать честная атомарность: read_committed=0 (output от убитой попытки не виден закоммиченному consumer’у), а offset группы всё ещё 0 — попытка A не продвинула чтение input ни на одну запись, хотя физически уже записала весь output. Попытка B стартует с тем же transactional.id — первый вызов initTransactions() (Java) или BeginTransaction() (Go) на стороне брокера фенсит и абортит зависшую транзакцию попытки A — и с той же consumer group, которая из-за незакоммиченного offset отдаёт ей те же самые 10 записей заново. После штатного коммита B: read_committed=10 ровно, каждое сообщение единожды.

Физически на диске к этому моменту 20 записей output — 10 от абортнутой A и 10 от закоммиченной B — и это видно не только в счётчике: под одним и тем же ключом (cpp-key-4) реально лежат две разные записи, offset 0 от попытки A и offset 3 от попытки B, в одной партиции. read_committed-consumer видит только вторую. Это наглядное доказательство, что «атомарность» здесь — не просто сошедшийся счётчик, а реальное физическое дублирование данных на диске, полностью скрытое изоляцией транзакций от читателя.

Идентично на Go и Java: ни разу за оба прогона не наблюдалось промежуточное состояние — offset входного топика и видимость output продвигаются строго вместе, одним прыжком, а не по частям.

read_committed: что видит consumer

isolation.level=read_committed — настройка consumer’а, не producer’а: она определяет, что попадает в poll(), а что фильтруется. Записи внутри незакоммиченной или абортнутой транзакции физически лежат в логе партиции наравне с обычными, но read_committed-consumer их не отдаёт — он ждёт control-record с исходом транзакции (COMMIT или ABORT) и либо показывает записи, либо пропускает их целиком. read_uncommitted (дефолт и в franz-go, и в Java) отдаёт всё как есть, включая содержимое ещё не завершённых или уже отменённых транзакций — тот самый режим, которым выше в статье считали физические 13 и 20.

cl, err := kgo.NewClient(
    kgo.SeedBrokers(seeds...),
    kgo.ConsumeTopics(topic),
    kgo.ConsumeResetOffset(kgo.NewOffset().AtStart()),
    kgo.FetchIsolationLevel(level), // kgo.ReadCommitted() или kgo.ReadUncommitted()
)
props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, isolationLevel); // "read_committed" | "read_uncommitted"
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");

Для сервисов, читающих выход транзакционного producer’а (например, следующий шаг конвейера, связывающего output одного этапа с input другого), read_committed — практически всегда правильный выбор: без него consumer видит промежуточные и заведомо отменённые данные, а разбираться с этим в бизнес-логике вручную — то же самое, что не иметь транзакций вовсе.

Граница EOS: где заканчивается гарантия Kafka

Всё, что показано выше, — гарантия строго внутри Kafka: топик-к-топику, топик-к-offset’у. Она не распространяется дальше брокера. Запись во внешнюю систему — реляционную БД, HTTP-вызов, кеш — той же транзакцией не покрывается: у Kafka нет способа откатить чужую транзакцию в PostgreSQL или отменить уже отправленный HTTP-вызов, если commit в Kafka не удался, и наоборот. Consume-process-produce атомарен ровно в тех границах, где «produce» — это запись в топик Kafka, а не куда угодно.

Практический мост через эту границу — паттерн outbox: пишем в БД и в outbox-таблицу в одной локальной транзакции СУБД, а отдельный relay уже асинхронно и at-least-once публикует из outbox в Kafka, полагаясь на идемпотентность или дедупликацию на стороне consumer’а. Разбор паттерна целиком, включая inbox-сторону и дедупликацию по idempotency key, — в статье «Transactional Outbox/Inbox»; прикладная сторона EOS в Kafka, включая то же consume-process-produce, но с прицелом на бизнес-сценарии, а не на протокол, — в статье «Транзакции в брокерах: RabbitMQ и Kafka».


Собранная воедино картина: идемпотентный producer убирает дубли от retry внутри одной партиции, транзакции расширяют это до атомарности сразу по нескольким партициям и топикам (включая offset потребления как частный случай записи), а read_committed — это сторона consumer’а, без которой транзакции producer’а бессмысленны: писать атомарно некому, если читатель всё равно видит незакоммиченный мусор. Три механизма работают только вместе.

Дальше в серии: как эксплуатировать транзакционный кластер — метрики, reassignment, история KRaft — в статье про эксплуатацию и тюнинг; как устроена экосистема вокруг Kafka — Connect и Schema Registry — в статье про Kafka Connect и Schema Registry. Для навигации по всей теме messaging — карта messaging-landscape-map.

Обсуждение в Telegram

Присоединиться →

Комментарии