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.

Смотри также

Мини-задание

  • Напиши producer с AllISRAcks и отправь 10 событий, убедись что все дошли
  • Намеренно создай дубль: отправь одну запись дважды с одинаковым ключом - проверь что idempotent producer не создаёт дублей в брокере
  • Реализуй транзакционный processor: читай из input-topic, обрабатывай, пиши в output-topic с commit offset атомарно
  • Попробуй запустить два процесса с одним TransactionalID - что происходит?
  • Для своего проекта: определи нужный уровень гарантии для каждого типа событий

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.