Exactly-once semantics: at-most-once / at-least-once / exactly-once. Three tiers of delivery semantics. Effective EOS = at-least-once delivery + idempotent consumer + atomic commit. Kafka transactional API (transactional.id, sendOffsetsToTransaction, isolation.level=read_committed, transaction coordinator with __transaction_state). Flink TwoPhaseCommitSinkFunction with pre-commit on checkpoint barrier and commit on notifyCheckpointComplete. Idempotent consumer with Redis dedup. Two Generals myth: exactly-once delivery невозможен, но exactly-once effects реален. ADRs: when EOS critical vs at-least-once + idempotency enough; effective EOS = three ingredients (delivery + dedup + atomic commit). Scenarios: at-least-once duplicate (double billing), Kafka EOS happy path, idempotent producer retry, transaction abort, Flink 2PC commit on checkpoint, Flink failure recovery, idempotent consumer dedup, EOS impossible without sink cooperation.
«Exactly-once» — самая хайповая и одновременно самая недопонимаемая тема в distributed systems. Половина инженеров считает это маркетингом («delivery exactly-once невозможен в network — Two Generals Problem»), другая половина думает «в Kafka есть флаг, включил и всё». Истина посередине: exactly-once delivery невозможен, но exactly-once effects (processing) — реальная и достижимая цель.
Цена ошибки огромная. Платежи на at-least-once без идемпотентности — двойные списания. Inventory — overselling и отрицательные остатки. Notifications — 5 SMS вместо одной. На больших масштабах 0.1% дублей = тысячи обращений в support и прямой денежный убыток. С другой стороны, exactly-once не бесплатно: транзакции Kafka добавляют 100-200 ms latency, режут throughput на 10-20%, 2PC sink стоит CPU и удваивает checkpoint cost.
Понимание EOS = понимание где именно ты закрываешь границу ошибок: producer (idempotent send), broker (transactions), consumer (manual offset commit), sink (idempotent write / 2PC). Без этой ментальной модели ты либо overengineer'ишь там, где at-least-once + dedup достаточно, либо теряешь и дублируешь данные там, где это критично.
«Exactly-once delivery в network невозможен (Two Generals); но exactly-once processing effects = at-least-once delivery + idempotent consumer (operation_id-based dedup) + atomic commit (output вместе с offset/state в одной транзакции). Все три ингредиента обязательны.»
Топология показывает полный pipeline EOS и его участников:
Producer (idempotent + txn) с PID/epoch и Flink Source,
читающий из Kafka.Tx Coordinator,
держащий state в __transaction_state.JobManager рассылает checkpoint barriers,
Operator обрабатывает поток, 2PC Sink коммитит атомарно с external
storage.Consumer (read_committed), Redis (dedup set, TTL 24h), External API (Idempotency-Key), PostgreSQL (MERGE / XA).Edges — это физические соединения, ответы идут по тем же edges в обратном
направлении (reverse animation). Например, prod → tc отвечает за
initTransactions и за commitTransaction; abort markers едут от координатора
во все partitions через те же edges.
В скрипте восемь сценариев — от «как сломать» до «как правильно»:
initTransactions (PID, epoch), send в
несколько partitions, sendOffsetsToTransaction (consumer offsets в той же
txn), двухфазный commit (PrepareCommit в txn log → commit markers во все
partitions), read_committed consumer видит данные только после marker, dedup
через Redis, idempotent POST.(PID, seq=2), broker узнаёт дубль и тихо возвращает success без
записи. Это решает только проблему ретрая одного инстанса.transaction.timeout.ms
истекает → координатор пишет Abort marker → consumer с read_committed
фильтрует aborted batch (LSO не двигается) → ноль downstream effects.notifyCheckpointComplete → sink делает commit → данные
visible downstream.notifyCheckpointComplete. JobManager детектит, sink делает abort/rollback
pending транзакций, job рестартует с последнего completed checkpoint,
процессинг повторяется — но без duplicate effects благодаря 2PC.SISMEMBER → skip / SADD ... EX 86400 → process).
Работает с любым sink, цена — extra round-trip и dedup storage.Idempotency-Key. Решение — требовать idempotency на стороне
sink, либо Saga + compensations.ADR-001: Когда EOS критично vs at-least-once + idempotency достаточно.
Транзакционный producer + read_committed добавляют 100-200 ms latency на
commit и режут throughput на 10-20%. Idempotent consumer через Redis — extra
round-trip и storage. Не везде это оправдано. Дерево выбора:
Idempotency-Key на consumer.MERGE/UPSERT в sink. Нет dedup-state, нет
отдельного хранилища.Правило: чем дороже дубль (в деньгах или доверии), тем выше уровень
гарантии. EOS через Kafka transactions включаем только для inter-Kafka
pipelines (Streams, Flink read-process-write); для external sinks
предпочитаем idempotent writes (MERGE, Idempotency-Key, conditional
update) — они проще, дешевле, не зависят от distributed coordinator.
ADR-002: Effective EOS = at-least-once delivery + idempotent consumer + atomic commit. Two Generals Problem доказывает: exactly-once delivery в network невозможен. Producer не может узнать, дошёл ли пакет с ACK обратно. Любая EOS-система на самом деле гарантирует exactly-once effects через дедупликацию + atomic state commit. Маркетинг от Kafka/Flink этот нюанс скрывает, разработчики думают, что «поставил флаг — получил магию».
Mental model команды — EOS состоит из трёх компонентов, все обязательны:
operation_id, либо UPSERT/MERGE по
unique key, либо conditional update) — гарантирует no duplicate effect.sendOffsetsToTransaction) — гарантирует, что после crash replay не
создаст phantom output.Если хоть один компонент отсутствует — это не EOS, не называйте это так в дизайн-доках. В code review требуем явного указания всех трёх для любого flow с monetary side effect.
Idempotency-Key header — golden standard для public APIs. Каждый
POST принимает Idempotency-Key: <uuid>, Stripe сохраняет результат (status +
body) на 24 часа, повторный запрос возвращает закэшированный ответ.processing.guarantee=exactly_once_v2
по умолчанию для внутренних pipeline'ов. v2 объединяет транзакции per app
instance вместо per task — меньше overhead.event_id). Сознательный выбор: EOS не везде
оправдан.event_id в Druid/ClickHouse.MERGE INTO — sink использует MERGE с unique key,
получая idempotent insert. Популярный pattern для ETL без транзакций.enable.idempotence=true — теперь у меня EOS». Нет. Idempotent
producer защищает только от повторов одного producer'а на retry в окне
последних 5 сообщений per partition. Если producer упадёт и перезапустится
без transactional.id — новый PID, дубли возможны.acks=1 + idempotence. Не работает. Idempotence требует acks=all.transaction.timeout.ms (default
60 сек) сработает, координатор сделает abort, вся работа потеряна. Дизайн:
одна txn на маленький batch, commit часто.isolation.level=read_committed на consumer. Producer делает
transactions, consumer читает с дефолтным read_uncommitted → видит aborted
данные → EOS сломан незаметно.transactional.id на много инстансов. Epoch fencing убьёт
работающего producer'а каждый раз, когда другой стартует. Используй уникальный
txid per instance (app-${pod_name}).Idempotency-Key. Если внешний
сервис не поддерживает дедупликацию — EOS-effects невозможен. Не делай
вид, что есть. Требуй idempotency от вендора, либо строй Saga с явными
compensating actions.MERGE / UPSERT /
conditional PUT — это уже idempotent write. Транзакции Kafka сверху не
нужны, проще и дешевле обойтись.Idempotency-Key, dedup tables, conditional updates.