Exactly-once в Kafka: idempotent producer и транзакции
Exactly-once в Kafka: idempotent producer и транзакции
Гарантии доставки - это договор между producer-ом, брокером и consumer-ом. Kafka поддерживает все три уровня, но для каждого нужна конкретная конфигурация.
Теорию мы уже разобрали в event-driven - гарантии доставки. Здесь - как это выглядит в коде.
At-most-once: отправить и забыть
client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.RequiredAcks(kgo.NoAck()), // не ждать подтверждения
)
client.Produce(ctx, &kgo.Record{
Topic: "events",
Value: payload,
}, nil) // callback nil - не обрабатываем ошибки
Producer не ждёт подтверждения. Максимальная скорость, нулевая надёжность. Для метрик и аналитики, где потеря события допустима.
At-least-once: стандартный режим
client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.RequiredAcks(kgo.AllISRAcks()), // ждать ack от всех in-sync реплик
kgo.RecordRetries(10), // retry при временных ошибках
kgo.RecordDeliveryTimeout(30 * time.Second),
)
При таймауте или сбое сети producer повторяет отправку. Возможны дубликаты, если первая отправка дошла до брокера, но ack потерялся.
Consumer должен уметь обрабатывать дубли - idempotent consumer.
Idempotent Producer
client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.RequiredAcks(kgo.AllISRAcks()),
// franz-go включает idempotent producer автоматически
// при AllISRAcks - не нужна явная настройка
)
Idempotent producer присваивает каждой записи уникальный sequence number. Если брокер получает повтор (тот же producer + та же партиция + тот же seq), он отбрасывает дубль. Это защищает от дублей при сетевых retry.
Ограничение: идемпотентность гарантируется только в рамках одной сессии producer-а. После перезапуска producer получает новый producer ID и дубли снова возможны.
Транзакции Kafka: exactly-once
Транзакции нужны для сценария read-process-write: прочитал из топика A, обработал, записал в топик B + commit offset. Всё это - атомарно.
// Producer с транзакционным ID
producer, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.TransactionalID("analytics-processor-1"), // уникальный per-instance
kgo.RequiredAcks(kgo.AllISRAcks()),
)
defer producer.Close()
// Consumer читает uncommitted записи по умолчанию
// Для exactly-once нужен ReadCommitted:
consumer, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.ConsumerGroup("analytics"),
kgo.ConsumeTopics("lesson-events"),
kgo.FetchIsolationLevel(kgo.ReadCommitted()),
kgo.DisableAutoCommit(),
)
defer consumer.Close()
for {
fetches := consumer.PollFetches(ctx)
if ctx.Err() != nil {
return
}
// Начинаем транзакцию
if err := producer.BeginTransaction(); err != nil {
log.Printf("begin transaction: %v", err)
continue
}
var produceErr error
fetches.EachRecord(func(r *kgo.Record) {
result := processRecord(r) // бизнес-логика
out := &kgo.Record{
Topic: "analytics-output",
Value: result,
}
if err := producer.ProduceSync(ctx, out).FirstErr(); err != nil {
produceErr = err
}
})
if produceErr != nil {
// Откатываем транзакцию
producer.EndTransaction(ctx, kgo.TryAbort)
continue
}
// Commit offset consumer-а в той же транзакции
offsets := consumer.MarkedOffsets()
if err := producer.AddOffsetsToTxn(ctx, offsets, "analytics"); err != nil {
producer.EndTransaction(ctx, kgo.TryAbort)
continue
}
// Коммит транзакции
if err := producer.EndTransaction(ctx, kgo.TryCommit); err != nil {
log.Printf("commit transaction: %v", err)
}
}
AddOffsetsToTxn связывает offset consumer-а с транзакцией. Если транзакция откатится - offset тоже не применится.
Acks: параметр надёжности
// Ack уровни
kgo.NoAck() // 0 - не ждать ответа брокера
kgo.LeaderAck() // 1 - ждать подтверждения от лидера партиции
kgo.AllISRAcks() // -1 - ждать от всех in-sync реплик
ISR (In-Sync Replicas) - реплики, которые синхронизированы с лидером. При AllISRAcks запись считается сохранённой только когда все ISR-реплики её подтвердили. Это защита от потери при падении лидера сразу после записи.
В dev-окружении с одним брокером AllISRAcks = LeaderAck (одна реплика в ISR). В продакшне с replication factor 3 - более строгая гарантия.
Когда что выбирать
| Гарантия | Применение | Конфигурация |
|---|---|---|
| At-most-once | Метрики, некритичные логи | NoAck, callback nil |
| At-least-once | События, нотификации | AllISRAcks + idempotent consumer |
| Exactly-once | Финансовые операции, Kafka Streams | Транзакции + FetchIsolationLevel(ReadCommitted) |
На практике at-least-once + idempotent consumer покрывает 95% случаев. Транзакции нужны только для read-process-write потоков.
Типичная ошибка
TransactionalID один на несколько инстансов - два процесса с одним TransactionalID фенсируют друг друга: при старте второго первый получает ошибку PRODUCER_FENCED. TransactionalID должен быть уникальным per-pod, обычно включает hostname или pod-индекс.
FetchIsolationLevel не установлен - consumer читает uncommitted записи из транзакций, которые позже могут быть отменены. Включай ReadCommitted на consumer-стороне.
Транзакции там, где не нужны - транзакции снижают throughput. Для простых событий достаточно at-least-once + idempotency key в payload.
Смотри также
- Event-Driven - Гарантии доставки - теория at-least-once и exactly-once
- Event-Driven - Idempotency - idempotent consumer на уровне бизнес-логики
- Message Brokers - Consumer - commit стратегии на стороне consumer-а
Мини-задание
- Напиши producer с
AllISRAcksи отправь 10 событий, убедись что все дошли - Намеренно создай дубль: отправь одну запись дважды с одинаковым ключом - проверь что idempotent producer не создаёт дублей в брокере
- Реализуй транзакционный processor: читай из
input-topic, обрабатывай, пиши вoutput-topicс commit offset атомарно - Попробуй запустить два процесса с одним
TransactionalID- что происходит? - Для своего проекта: определи нужный уровень гарантии для каждого типа событий