Как реализовать exactly-once обработку событий в Python-сервисе с Kafka (consume → обработка → запись в БД/другую тему) при возможных падениях сервиса?
Строим at-least-once + идемпотентность. Варианты:
Transactional outbox: пишем результат и outbox-запись в одной БД-транзакции; отдельный паблишер надёжно читает outbox и шлёт в Kafka, после чего помечает запись.
Или Kafka EOS: confluent-kafka с enable.idempotence=true и transactional.id; оборачиваем read-process-write и commit offset в одну Kafka-транзакцию.
В обоих случаях делаем операции идемпотентными (ключи/версионирование), коммитим оффсеты только после durable-записи результата, храним processed keys (или используем upsert).
Библиотека собеса по Python
Post #1059
900
- 👍 1