Гарантии доставки и идемпотентность: at-least-once и миф exactly-once
Есть два предложения, которые звучат почти одинаково, но описывают разные вселенные:
- «Сообщение будет доставлено ровно один раз».
- «Эффект сообщения будет применён ровно один раз».
Первое — невозможно. Второе — рутинная инженерная задача, которую решают тысячи систем каждый день. Вся эта статья — про то, почему граница проходит именно здесь, и что конкретно нужно построить, чтобы жить на правильной стороне границы.
Ключевая мысль, которую стоит унести, даже если дальше вы ничего не прочитаете: гарантия доставки — это свойство не транспорта, а пары «транспорт + обработчик». RabbitMQ, Kafka, SQS, gRPC-ретраи, retry в HTTP-клиенте — все они дают одно и то же: at-least-once. Разница между надёжной системой и системой, которая дважды списывает деньги, находится не в брокере, а в двадцати строках кода обработчика. Брокер не может знать, что значит «применить платёж дважды» — это знаете только вы.
Мы опираемся на модели отказов (главное оттуда: отказ узла неотличим от медленной сети), на время и часы (тайм-аут — это решение, а не наблюдение) и на репликацию (там гарантии доставки уже мелькали как частный случай применения записи на реплике). Дальше — детально.
Что вообще значит «доставлено»
Слово «доставка» скрывает как минимум шесть разных событий, и путаница между ними — источник половины споров об exactly-once. Отправитель сформировал запрос; брокер принял его в память; запись стала долговечной и реплицированной; обработчик её получил; эффект применён; подтверждение вернулось к отправителю. Между каждой парой соседних событий может произойти отказ.
Главное свойство этой картинки — наблюдательная эквивалентность. Отправитель во всех пяти случаях видит ровно одно: тайм-аут. У него нет способа отличить «не дошло» от «дошло, применилось, ответ потерялся». Это не ограничение реализации, а фундаментальное свойство асинхронного канала: любое подтверждение само нуждается в подтверждении, и рекурсия не заканчивается за конечное число сообщений.
| Точка отказа | Эффект был? | Ретрай помогает? | Что видно в логах |
|---|---|---|---|
| 1. Запрос не дошёл | нет | да, обязателен | context deadline exceeded, connection reset by peer |
| 2. Брокер принял в память и упал | нет | да | у клиента NotEnoughReplicasException или тишина; у брокера — ничего, он мёртв |
| 3. Обработчик упал до эффекта | нет | да | CommitFailedException, ребаланс группы |
| 4. Эффект применён частично | частично | даёт частичный дубль | половина строк изменена, внешний вызов не сделан |
| 5. Эффект применён, ack потерян | да | даёт полный дубль | у отправителя тайм-аут, у получателя 200 OK в access-логе |
Именно строки 4 и 5 делают exactly-once невозможным на уровне канала. Всё остальное лечится ретраем.
Три гарантии — и почему их ровно три
Отправитель, не получивший подтверждения, имеет ровно два варианта поведения. Не повторять — тогда сообщения из точек 1–3 теряются: это at-most-once. Повторять — тогда сообщения из точек 4–5 применяются дважды: это at-least-once. Третьего варианта поведения не существует, потому что для него нужна информация, которой у отправителя физически нет.
это отсутствие идемпотентности на стороне сервиса
Формально это та же структура, что у задачи двух генералов: соглашение по ненадёжному каналу за конечное число сообщений недостижимо. Первая формулировка — Jim Gray, «Notes on Data Base Operating Systems» (1978), строгий анализ через логику знания — Halpern, Moses, «Knowledge and Common Knowledge in a Distributed Environment» (1990). Прикладной разбор, который стоит дать коллеге в споре, — «You Cannot Have Exactly-Once Delivery».
| Гарантия | Механика | Когда допустима | Реальные примеры |
|---|---|---|---|
| at-most-once | отправил и забыл, ретраев нет | метрики, телеметрия, логи, кадры видео | StatsD по UDP, fire-and-forget продюсер с acks=0 |
| at-least-once | ретраить до подтверждения | практически всё остальное | SQS, Kafka по умолчанию, RabbitMQ с ack, gRPC retry |
| exactly-once delivery | не существует | — | — |
| effectively-once processing | at-least-once + идемпотентное применение | всё, где дубль стоит денег | Stripe, Kafka EOS внутри Kafka, Flink с 2PC-стоками |
Термин effectively-once (иногда «exactly-once semantics», EOS) — это честное название того, что реально строят: доставили несколько раз, применили один. Наблюдаемый эффект «как будто ровно один раз», механика — дедупликация или идемпотентность на стороне получателя. Отсюда две ошибки терминологии, которые стоит различать. «Exactly-once delivery» — невозможно, это утверждение про канал. «Exactly-once processing» — возможно, это утверждение про наблюдаемое состояние системы после обработки. Когда вендор пишет «exactly-once», почти всегда имеется в виду второе, ограниченное границами его продукта.
Идемпотентность: определение точнее, чем «повтор безопасен»
Операция f идемпотентна, если f(f(x)) = f(x) для любого состояния x. Термин из алгебры: идемпотентный элемент равен своему квадрату. Не путайте с двумя соседними свойствами, которые часто нужны одновременно, но это разные вещи:
- идемпотентность — повтор не меняет результат:
set(x, 5)дважды даёт то же, что один раз; - коммутативность — порядок не важен:
set(x,5); set(y,7)=set(y,7); set(x,5), а вотset(x,5); set(x,7)— нет; - ассоциативность — группировка не важна; нужна для деревьев слияния и CRDT.
Дедупликация решает проблему повторов. Она не решает проблему переупорядочивания — это отдельный разговор ниже.
Естественно идемпотентные операции:
UPDATE users SET email = 'a@b.c' WHERE id = 42; -- идемпотентно: присваивание
DELETE FROM sessions WHERE id = 'sess_7'; -- идемпотентно: второй раз удалит 0 строк
INSERT ... ON CONFLICT (id) DO NOTHING; -- идемпотентно: upsert
UPDATE orders SET status='paid' WHERE id=42 AND status='pending'; -- идемпотентно: условный переход
Естественно НЕ идемпотентные:
UPDATE accounts SET balance = balance - 100 WHERE id = 42; -- инкремент: каждый повтор списывает снова
INSERT INTO audit_log (...) VALUES (...); -- append: каждый повтор добавляет строку
Рецепт превращения второго в первое всегда один из трёх: (а) переписать инкремент как присваивание с известным целевым значением; (б) добавить условие на текущее состояние — compare-and-set по версии; (в) добавить внешний ключ идемпотентности и проверять его. Вариант (а) чаще всего невозможен в конкурентной среде, (б) требует версии в модели данных, (в) универсален и потому используется чаще всего.
HTTP: идемпотентность по спецификации и идемпотентность на практике
RFC 9110 §9.2.2 объявляет GET, HEAD, PUT, DELETE, OPTIONS, TRACE идемпотентными, а POST и PATCH — нет. Здесь два подвоха, на которых регулярно горят.
Первый: спецификация описывает намерение, а не гарантию. Если ваш PUT /orders/42 внутри делает INSERT INTO audit_log, метод перестал быть идемпотентным — и никакой RFC этого не исправит. HTTP-клиенты, прокси и service mesh (Envoy, nginx, gRPC) при этом опираются на спецификацию и автоматически ретраят идемпотентные методы при 503, разрыве соединения или тайм-ауте. Вы получите повтор, о котором не просили.
Второй: POST не идемпотентен по спецификации, поэтому его не ретраит инфраструктура — но его ретраит ваш собственный клиентский код, мобильное приложение и пользователь, нажимающий кнопку второй раз. Для POST идемпотентность делается явным ключом в заголовке. Канонический пример — Idempotency-Key у Stripe; есть и черновик стандарта IETF Idempotency-Key header field. Подробнее про дизайн API — в стилях API.
Сценарий отказа: сервис отдаёт 502 из-за перегрузки апстрима, но запрос при этом уже обработан. Envoy с retry_on: gateway-error повторяет PUT, обработчик пишет вторую строку в журнал начислений. В логах это выглядит идеально безобидно:
[access] PUT /v1/bonus/42 502 upstream_reset_before_response_started{connection_termination} 8134ms
[access] PUT /v1/bonus/42 200 - 122ms x-envoy-attempt-count: 2
Никакой ошибки. Просто два бонуса вместо одного, и обнаружится это через месяц при сверке.
Ключ идемпотентности
Ключ идемпотентности — это идентификатор логической операции, стабильный между всеми её повторами. Три требования, каждое из которых нарушают регулярно:
- Стабильность. Ключ генерируется один раз в момент формирования намерения и переиспользуется всеми ретраями.
uuid4()внутри функцииretry()— самая частая ошибка в этой области: каждый повтор приносит новый ключ, дедупликация не срабатывает, а код при этом выглядит правильным. - Детерминированная область действия. Ключ уникален в пределах пары (отправитель, тип операции). Один и тот же
event_idможет законно обрабатываться тремя разными обработчиками — значит, первичный ключ дедупликации это(handler, key), а неkey. - Привязка к содержимому. Если по тому же ключу пришёл другой payload — это ошибка вызывающего, а не повтор. Stripe в этом случае отвечает
400сidempotency_error; так же стоит делать и вам:
{"error": {"type": "idempotency_error",
"message": "Keys for idempotent requests can only be used with the same parameters they were first used with."}}
Откуда брать ключ, в порядке предпочтения: естественный бизнес-идентификатор (order_id + тип операции — лучший вариант, он существует независимо от транспорта); идентификатор события от продюсера (event_id в конверте сообщения); координаты в логе (topic-partition-offset — работает, но ломается при пересоздании топика или миграции); хеш канонизированного payload (последнее средство: два законных одинаковых платежа станут одним).
Дедупликация: хранилище, окно и его цена
Псевдокод протокола идемпотентного применения — он же формулировка того, что нужно реализовать:
APPLY(key, handler, payload):
atomically: # одна транзакция, без исключений
row := INSERT (handler, key, state=IN_PROGRESS) IF NOT EXISTS
if row не вставлена:
existing := SELECT по (handler, key)
if existing.state = DONE: return existing.response # дубль → отдаём сохранённый ответ
if existing.state = IN_PROGRESS: return RETRY_LATER (409) # конкурентный ретрай
response := EXECUTE(payload) # бизнес-эффект здесь же
UPDATE row SET state=DONE, response=response
return response
Сложность: одна вставка по первичному ключу — O(log n) обращений к страницам B-дерева, на практике 3–4, то есть O(1) по времени на сообщение. По памяти O(K), где K — число ключей внутри окна хранения. Последнее критично: без политики очистки таблица дедупликации растёт линейно по трафику и однажды перестаёт помещаться в кэш, после чего самой дорогой операцией системы становится проверка на дубликат.
| Хранилище | Плюс | Минус | Где встречается |
|---|---|---|---|
| Таблица в той же БД, что и эффект | атомарность бесплатно | нагрузка на основную БД, нужна чистка | большинство сервисов |
Redis SET key NX EX |
быстро, TTL встроен | отдельная система → эффект и отметка не атомарны | кэш-слой, антифрод |
| RocksDB в состоянии стрим-процессора | локально, переживает рестарт | привязано к партиционированию | Kafka Streams, Flink |
| Bloom-фильтр | O(1) памяти |
false positive = потерянное сообщение | только как negative cache |
Про Bloom стоит отдельно: вероятностный фильтр даёт ложноположительные срабатывания, а ложноположительное «уже видели» означает молча выброшенное сообщение. Использовать его можно исключительно как быстрый отрицательный ответ («точно не видели — обрабатываем без похода в БД»), а положительный обязан подтверждаться точной проверкой.
Окно хранения ключей выбирается как максимум из: retention топика, максимального времени жизни ретрая у клиента, времени разбора инцидента вручную. Практический ориентир — 7 дней, у Stripe ключи живут 24 часа. Если окно короче, чем возможная задержка ретрая, дедупликация работает «почти всегда» — а это худший вид гарантии, потому что нарушается только под нагрузкой и в инцидентах.
CREATE TABLE processed_messages (
handler TEXT NOT NULL, -- один ключ может обрабатываться разными потребителями
idempotency_key TEXT NOT NULL, -- стабильный ключ от продюсера
state TEXT NOT NULL, -- IN_PROGRESS | DONE | FAILED
payload_hash TEXT, -- чтобы поймать «тот же ключ, другой payload»
response JSONB, -- сохранённый ответ для повторов
lease_until TIMESTAMPTZ, -- защита от зависшего IN_PROGRESS
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
PRIMARY KEY (handler, idempotency_key)
);
-- Чистка: обязательна. Без неё таблица растёт линейно по трафику.
CREATE INDEX ON processed_messages (created_at);
DELETE FROM processed_messages WHERE created_at < now() - INTERVAL '7 days';
Атомарность отметки и эффекта — единственное место, где всё решается
Если запись «обработано» и сам эффект коммитятся раздельно, вы не убрали окно дубля, а переместили его. Между двумя коммитами процесс может умереть — и вы получаете либо дубль (упал после эффекта, до отметки), либо потерю (упал после отметки, до эффекта). Единственное лекарство — одна транзакция.
import hashlib, json, psycopg
from psycopg.errors import UniqueViolation
from psycopg.rows import dict_row
class Conflict(Exception): """Ключ обрабатывается прямо сейчас другим воркером — наружу 409."""
class KeyReuse(Exception): """Тот же ключ, другой payload — ошибка вызывающего, наружу 400."""
def apply_once(conn: psycopg.Connection, handler: str, key: str, payload: dict) -> dict:
"""Идемпотентное применение эффекта. Возвращает ответ — новый или сохранённый.
Инвариант: отметка о ключе и бизнес-эффект попадают в ОДНУ транзакцию.
Откатится она — откатится и то, и другое, и повтор пройдёт честно.
"""
h = hashlib.sha256(json.dumps(payload, sort_keys=True, ensure_ascii=False).encode()).hexdigest()
try:
with conn.transaction(): # BEGIN ... COMMIT / ROLLBACK
conn.execute(
"INSERT INTO processed_messages (handler, idempotency_key, state, payload_hash,"
" lease_until) VALUES (%s, %s, 'IN_PROGRESS', %s, now() + interval '60 seconds')",
(handler, key, h),
)
row = conn.execute( # бизнес-эффект — здесь же, в той же транзакции
"UPDATE orders SET status = 'paid', paid_at = now()"
" WHERE id = %s AND status = 'pending' RETURNING id",
(payload["order_id"],),
).fetchone()
result = {"order_id": payload["order_id"], "applied": row is not None}
conn.execute(
"UPDATE processed_messages SET state = 'DONE', response = %s"
" WHERE handler = %s AND idempotency_key = %s",
(json.dumps(result), handler, key),
)
return result
except UniqueViolation:
pass # ключ занят: разбираемся в НОВОЙ транзакции, текущая помечена как сбойная
with conn.cursor(row_factory=dict_row) as cur: # кем именно занят
prev = cur.execute(
"SELECT state, payload_hash, response, lease_until FROM processed_messages"
" WHERE handler = %s AND idempotency_key = %s", (handler, key)).fetchone()
if prev["payload_hash"] != h:
raise KeyReuse(f"key {key} already used with different parameters")
if prev["state"] == "DONE":
return prev["response"] # честный повтор — отдаём сохранённый ответ
raise Conflict(f"key {key} is being processed right now")
Три детали, которые отличают рабочий код от кода из статьи в блоге:
- Ответ сохраняется. Повтор должен вернуть то же самое, что вернул оригинал, иначе клиент увидит
404там, где ждал201, и вы получите баг, воспроизводимый только под ретраями. IN_PROGRESSимеет лизу. Воркер, который взял ключ и умер, оставит запись навсегда. Полеlease_untilпозволяет другому воркеру перехватить ключ по истечении срока — это ровно тот же механизм лиз, что и в координации.- Конкурентный ретрай отдаёт
409, а не ждёт. Ожидание под блокировкой превращает всплеск ретраев в очередь на коннекты к БД и дальше в метастабильный отказ.
лиза 60 с NEW --> DUPLICATE: ключ уже есть DUPLICATE --> RETURN_SAVED: state = DONE
отдаём сохранённый ответ DUPLICATE --> CONFLICT_409: state = IN_PROGRESS
лиза жива DUPLICATE --> IN_PROGRESS: state = IN_PROGRESS
лиза истекла — перехват DUPLICATE --> KEY_REUSE_400: payload_hash не совпал IN_PROGRESS --> DONE: эффект применён,
COMMIT одной транзакцией IN_PROGRESS --> NEW: ROLLBACK — записи нет,
повтор пройдёт честно IN_PROGRESS --> FAILED: бизнес-ошибка
терминальная, ретрай не поможет DONE --> [*] RETURN_SAVED --> [*] FAILED --> [*] note right of DONE Строки живут ровно окно дедупликации. Нет чистки — таблица растёт линейно. end note
Если эффект физически не может быть в той же БД, что и отметка (эффект — вызов внешнего API), атомарность восстанавливается паттерном inbox/outbox: в транзакции пишется только намерение, а отдельный процесс доставляет его наружу с ретраями и своим ключом идемпотентности. Это уже территория распределённых транзакций и saga.
Порядок — вторая половина задачи
Дедупликация убирает повторы, но не переупорядочивание. Классический сценарий: обработчик получил set(x, 1), завис на GC-паузе в 8 секунд, за это время пришло и применилось set(x, 2), потом ожил и применил своё set(x, 1). Каждое сообщение применено ровно один раз — и состояние всё равно неверное.
Лечится это не дедупликацией, а монотонностью: у каждой записи есть версия, и применяется только та, чья версия строго больше текущей.
// Условная запись с фенсинг-токеном: эффект применяется, только если токен не устарел.
// Токен монотонно растёт и выдаётся координатором (etcd, Raft-лидер).
func (s *Store) Apply(ctx context.Context, key string, val []byte, fence uint64) error {
res, err := s.db.ExecContext(ctx, `
UPDATE items SET value = $1, fence = $2 WHERE key = $3 AND fence < $2`, val, fence, key)
if err != nil {
return err
}
if n, _ := res.RowsAffected(); n == 0 {
// Не ошибка сети и не дубль: наш токен устарел, пока мы были в GC-паузе.
return ErrFenced // правильная реакция — сдаться, а не ретраить
}
return nil
}
Это тот же приём, что описан у Мартина Клеппманна в «How to do distributed locking»: владение ресурсом подтверждается не фактом «я держу блокировку», а номером, который проверяет получатель эффекта. Механику выдачи токенов разбираем в координации, а логическую основу монотонных версий — в часах.
Практический вывод: ретрай и переупорядочивание — конфликтующие требования. Если вы хотите строгий порядок в партиции Kafka и одновременно параллельные in-flight запросы, вы обязаны включить идемпотентный продюсер: без него retries > 0 при max.in.flight.requests.per.connection > 1 переставляет сообщения местами при повторе.
Kafka: что именно означает её exactly-once
Kafka даёт две разные вещи, которые в маркетинге сливаются в одно слово.
Идемпотентный продюсер (KIP-98). Каждая партия сообщений несёт (producer_id, epoch, sequence_number). Брокер держит последние 5 последовательностей на партицию и отбрасывает повтор. Это устраняет дубли, вызванные ретраями продюсера — и только их. С Kafka 3.0 включено по умолчанию.
# Продюсер: минимальный набор, при котором ретрай не порождает дубль и не ломает порядок
enable.idempotence: true # включает (pid, epoch, seq); с 3.0 по умолчанию
acks: all # ждать все ISR-реплики, иначе точка отказа 2 из схемы выше
max.in.flight.requests.per.connection: 5 # >5 несовместимо с идемпотентностью
retries: 2147483647 # ретраить до delivery.timeout.ms, а не N раз
delivery.timeout.ms: 120000 # реальная граница попыток
Границы, о которых нужно знать заранее:
- Гарантия действует в пределах сессии продюсера. Перезапуск процесса без
transactional.idвыдаёт новыйproducer_id— брокер видит нового отправителя и повторы больше не распознаёт. - Гарантия действует в пределах партиции. Тот же ключ, отправленный в другую партицию после ребаланса или смены партиционера, не дедуплицируется.
- Состояние
producer_idна брокере живётtransactional.id.expiration.ms(по умолчанию 7 дней) и не переживает удаление сегментов. Продюсер, молчавший дольше, получит:
WARN [Producer clientId=billing-emitter] Got error produce response with correlation id 91243
on topic-partition payments-3, retrying (2147483646 attempts left). Error: UNKNOWN_PRODUCER_ID
Транзакции делают атомарной связку «записать в выходные топики + закоммитить оффсеты входных». Продюсер объявляет transactional.id, координатор транзакций фенсит старую эпоху этого id (защита от зомби-инстанса после ребаланса), консьюмер с isolation.level=read_committed не видит записей незавершённых транзакций — он читает только до Last Stable Offset. Это и есть «exactly-once» из документации: атомарность цикла consume → process → produce внутри Kafka.
Что она НЕ покрывает: любой побочный эффект наружу. HTTP-вызов платёжного шлюза, письмо, запись в чужую БД. Kafka не участвует в вашей транзакции с миром — и это не недоработка, а следствие первой части статьи.
Сценарий отказа: дубль из-за ребаланса
Обработчик читает пачку из 500 сообщений, обрабатывает 40 секунд, max.poll.interval.ms — 30 секунд. Брокер считает консьюмера мёртвым, запускает ребаланс, партиция уходит другому инстансу. Первый инстанс успел применить эффекты, но коммит оффсета не проходит:
[2026-03-14 02:11:47,113] WARN [Consumer clientId=billing-3, groupId=billing]
Synchronous auto-commit of offsets {payments-7=OffsetAndMetadata{offset=88134021}} failed:
Offset commit cannot be completed since the consumer is not part of an active group
for auto partition assignment; it is likely that the consumer was kicked out of the group.
[2026-03-14 02:11:47,240] INFO [Consumer clientId=billing-5, groupId=billing]
Setting offset for partition payments-7 to the committed offset 88133521
Разница в 500 сообщений — это ровно те записи, что будут обработаны вторично. В логах бизнес-обработчика при рабочей дедупликации это выглядит так, и это норма, а не инцидент:
{"lvl":"info","msg":"duplicate suppressed","handler":"payment_captured",
"key":"evt_01HXQ8K2","age_ms":1873,"partition":7,"offset":88133612}
Метрика duplicate_suppressed_total должна существовать и быть ненулевой. Если она строго ноль — почти наверняка дедупликация не работает, а не дубликатов нет.
Порядок коммита оффсета определяет гарантию
for {
msgs := consumer.Poll(200 * time.Millisecond)
for _, m := range msgs {
// at-least-once: сначала эффект, потом коммит.
// Падение между ними → повтор. Это правильный порядок.
if err := applyOnce(ctx, db, "payment_captured", keyOf(m), m.Value); err != nil {
return err // НЕ коммитим — сообщение придёт снова
}
}
consumer.CommitSync() // коммитим только полностью обработанную пачку
}
Поменяйте две строки местами — получите at-most-once и тихую потерю сообщений при каждом падении. Именно поэтому enable.auto.commit=true опасен по умолчанию: он коммитит по таймеру, независимо от того, применён эффект или нет.
Flink, Spanner и прочие носители слова exactly-once
Flink обеспечивает exactly-once для состояния: распределённые снапшоты по алгоритму Chandy — Lamport (оригинальная статья 1985, реализация во Flink) записывают согласованный срез, и после сбоя оператор откатывается к нему. Для внешних приёмников этого мало, поэтому есть TwoPhaseCommitSinkFunction: запись предварительно фиксируется, а коммит происходит вместе с чекпойнтом. Требование к приёмнику — уметь транзакции или идемпотентную запись. Kafka-сток использует транзакции Kafka, файловый — атомарный rename, JDBC-сток без транзакций даёт at-least-once, что бы ни было написано в конфиге.
Spanner гарантирует линеаризуемость коммита, но не однократность его инициации: клиент, не получивший ответ, всё равно не знает, применилась транзакция или нет. Рекомендация Google в этой ситуации ровно та же, что у нас: делать транзакцию идемпотентной или проверять состояние по бизнес-ключу перед повтором.
Общее правило для чтения документации: находите фразу, определяющую границу системы, и приписывайте exactly-once только к тому, что внутри. Если границы не сформулировано — гарантии тоже нет.
Эффекты, которые нельзя откатить
Механика идемпотентности зависит от природы эффекта, а их всего три класса.
в ТОЙ ЖЕ транзакции, что и эффект"] B -->|"Во внешнем API"| D{"API принимает
ключ идемпотентности?"} B -->|"В реальном мире:
письмо, SMS, отгрузка"| E["Дубль неустраним.
Сужаем окно и делаем его безвредным"] D -->|"Да — Stripe, PayPal, Adyen"| F["Стабильный ключ = f(бизнес-событие),
провайдер дедуплицирует сам"] D -->|"Нет"| G{"Есть операция чтения
по бизнес-ключу?"} G -->|"Да"| H["Двухшаговый протокол:
GET статус → если нет, POST"] G -->|"Нет"| I["Локальный журнал попыток
+ ручная сверка + алерт"] C --> J["Effectively-once:
эффект применён ровно один раз"] F --> J H --> K["Окно дубля сузилось до RTT,
но не исчезло — гонка остаётся"] E --> K I --> K J --> L["Дубли логируем как suppressed —
это норма, а не ошибка"] K --> M["Дубль виден пользователю.
Делаем его понятным:
«код уже отправлен» вместо второго SMS"]
Двухшаговый протокол из ветки «нет ключа идемпотентности» честно описать так: он сужает окно, а не закрывает его. Между GET и POST есть промежуток, в который параллельный воркер успеет сделать свой POST. Если цена дубля высока, окно закрывается локальной блокировкой по бизнес-ключу (лиза в БД или etcd) — но тогда вы получаете зависимость от доступности координатора и все её последствия.
Про эффекты в физическом мире стоит помнить простую вещь: выбор не между «дубль» и «без дубля», а между «дубль» и «потеря». Записать «отправлено» до отправки — это выбор в пользу возможной потери; после отправки — в пользу возможного дубля. Выбирать нужно осознанно и по цене ошибки.
Ретраи: как не превратить лечение в отказ
At-least-once означает ретраи, а ретраи — это усилитель нагрузки, который включается ровно в тот момент, когда система и так деградировала.
- Экспоненциальный backoff с джиттером — обязателен. Без джиттера все клиенты повторяют синхронно и создают периодические пики. Каноническая формулировка — AWS: Exponential Backoff and Jitter; практический выбор — full jitter.
- Бюджет ретраев. Ограничивайте не «число попыток на запрос», а долю ретраев от общего трафика (в Envoy —
retry_budget, у Google в SRE-книге — retry budget 10%). Без бюджета трёхуровневая цепочка сервисов с тремя попытками на уровне даёт 27-кратное усиление нагрузки на нижний сервис. - Ретраить только то, что имеет смысл.
400,403,422не пройдут и со второго раза; ретрай на них — чистая трата бюджета. Ретраить стоит408,429,5xx, тайм-ауты и разрывы соединения — и обязательно уважатьRetry-After. - Poison message и DLQ. Сообщение, которое всегда падает, будет ретраиться вечно и заблокирует партицию (в Kafka — буквально: партиция упорядочена). Нужен счётчик попыток и перенос в dead-letter после N неудач, с сохранением исходного ключа идемпотентности — иначе после ручного разбора вы обработаете его дважды.
import random, time
def retry_with_full_jitter(fn, *, attempts: int = 6, base: float = 0.2, cap: float = 20.0):
"""Full jitter: sleep = random(0, min(cap, base * 2**i)).
Ключ идемпотентности создаётся ВЫШЕ по стеку и не меняется между попытками —
иначе дедупликация на приёмнике никогда не сработает.
"""
for i in range(attempts):
try:
return fn()
except RetryableError as exc: # ваш класс: 408, 429, 5xx, тайм-аут, разрыв
if i == attempts - 1:
raise
delay = random.uniform(0, min(cap, base * (2 ** i)))
if exc.retry_after is not None: # сервер знает лучше нас
delay = max(delay, exc.retry_after)
time.sleep(delay)
Сценарий отказа, который стоит увидеть один раз: апстрим замедлился с 50 мс до 900 мс, клиентские тайм-ауты — 1 с, ретраев три. Нагрузка на апстрим утроилась, он замедлился ещё, тайм-ауты начали срабатывать всегда, и система не вернулась в норму даже после устранения исходной причины — очередь ретраев сама себя поддерживает. Это метастабильный отказ; выход — только сброс трафика. Подробнее в паттернах устойчивости.
Как это тестировать и как за этим следить
Идемпотентность — свойство, которое не проверяется обычными тестами: они не повторяют вызовы. Нужны специальные приёмы.
- Property-based тест на дубли. Для случайной последовательности операций проверяем: применение мультимножества
Sи применениеS, в котором каждый элемент случайно продублирован 1–3 раза, дают одинаковое финальное состояние. Это буквальная проверкаf(f(x)) = f(x)на уровне системы. - Инъекция дубликатов в staging. Прокси, который с вероятностью 5% отправляет запрос дважды. Через неделю все неидемпотентные пути будут найдены — до продакшна, а не после.
- Убийство воркера в окне между эффектом и коммитом. Самый ценный тест:
kill -9внутри обработчика междуCOMMITэффекта иCommitSync()оффсета. Ровно точка 5 из схемы. - Jepsen-подобные проверки с сетевыми разделениями и последующей проверкой инвариантов на истории. Подробно — в тестировании распределённых систем.
Что должно быть в метриках: duplicates_suppressed_total по обработчикам (ноль = подозрительно), idempotency_key_conflicts_total (гонки параллельных ретраев), retry_ratio (доля ретраев в трафике; выше 10% — тревога), размер и возраст таблицы дедупликации, лаг DLQ. Ключ идемпотентности стоит класть в trace-контекст: тогда все повторы одной логической операции собираются в одну картину — см. наблюдаемость.
Типичные ошибки
- Генерировать ключ идемпотентности внутри функции ретрая. Каждая попытка приносит новый ключ, дедупликация не работает, код при этом выглядит корректно.
- Коммитить отметку и эффект раздельно. Окно дубля не исчезло, а переехало. Это распределённая транзакция с двумя участниками, со всеми последствиями.
- Считать
PUTидемпотентным потому, что так написано в RFC. Идемпотентен обработчик, а не глагол. - Не сохранять ответ. Повтор возвращает
404или409там, где оригинал вернул201, и клиент делает неверный вывод. - Забыть про чистку таблицы дедупликации. Через год она больше основной таблицы, а проверка дубля — самый дорогой запрос в системе.
- Верить, что Kafka EOS покрывает вызов платёжного шлюза. Не покрывает: граница проходит по краю Kafka.
- Лечить переупорядочивание дедупликацией. Не лечится; нужны монотонные версии и фенсинг-токены.
- Ретраить
400и422. Пустая трата бюджета ретраев в момент, когда он нужнее всего. - Обрабатывать дубликат как ошибку — с алертом и записью в error-лог. Дубликат при at-least-once — это норма; шум от него приучает игнорировать алерты.
- Писать в архитектурном решении «мы обеспечиваем exactly-once». Это непроверяемое утверждение; проверяемое — «доставка at-least-once, обработка effectively-once за счёт ключей идемпотентности, окно дедупликации 7 дней».
Мини-итог
- Гарантия доставки — свойство пары «транспорт + обработчик». Транспорт умеет только at-least-once или at-most-once.
- Отправитель не может отличить «не дошло» от «дошло, но ack потерян». Из этого следует ровно два варианта поведения и ровно две гарантии; exactly-once delivery не существует.
- Работающая замена — effectively-once: at-least-once плюс идемпотентное применение. Наблюдаемо «один раз», механически «несколько доставок, одно применение».
- Идемпотентность делается тремя способами: присваивание вместо инкремента, условие на версию (CAS), внешний ключ идемпотентности. Третий универсален, и отметка о ключе обязана коммититься одной транзакцией с эффектом — всё остальное детали.
- Дедупликация не решает переупорядочивание. Для него нужны монотонные версии и фенсинг-токены.
- Kafka EOS и Flink exactly-once — реальные и полезные гарантии, но строго внутри границы своей системы. Первый внешний вызов гарантию обрывает.
- Ретраи — усилитель нагрузки. Backoff с джиттером, бюджет ретраев и DLQ обязательны, иначе лечение станет причиной метастабильного отказа.
Источники
- Jim Gray, Notes on Data Base Operating Systems (1978) — исходная формулировка задачи двух генералов в инженерном контексте.
- Halpern, Moses, Knowledge and Common Knowledge in a Distributed Environment (JACM 1990) — почему общее знание недостижимо по ненадёжному каналу.
- Chandy, Lamport, Distributed Snapshots: Determining Global States of Distributed Systems (1985) — основа чекпойнтов Flink.
- KIP-98: Exactly Once Delivery and Transactional Messaging и KIP-129 (EOS в Kafka Streams) — первоисточники по механике.
- Confluent, Exactly-Once Semantics Are Possible: Here’s How Kafka Does It — и критический ответ «You Cannot Have Exactly-Once Delivery». Читать оба: спор в них терминологический, и он полезен.
- Carbone et al., Lightweight Asynchronous Snapshots for Distributed Dataflows (2015) — как Flink делает согласованные снапшоты без остановки потока.
- RFC 9110, HTTP Semantics §9.2.2 и Stripe: Idempotent requests — стандарт и образцовая реализация.
- Martin Kleppmann, How to do distributed locking — фенсинг-токены; «Designing Data-Intensive Applications», глава 11 — потоки и гарантии обработки.
- AWS Builders’ Library, Timeouts, retries and backoff with jitter — почему без джиттера становится хуже.
Что дальше
Мы несколько раз упирались в одну и ту же потребность: кому-то нужно выдать монотонный токен, кто-то должен решить, что лиза протухла и ключ можно перехватить, а зомби-инстанс после ребаланса нужно отсечь до того, как он применит устаревший эффект. Всё это — задачи координации: Координация: распределённые блокировки, лизы, etcd и ZooKeeper.