Распределённые системы Гарантии доставки и идемпотентность: at-least-once и миф exactly-once
0%

Гарантии доставки и идемпотентность: at-least-once и миф exactly-once

Гарантии доставки и идемпотентность: at-least-once и миф exactly-once

Есть два предложения, которые звучат почти одинаково, но описывают разные вселенные:

  1. «Сообщение будет доставлено ровно один раз».
  2. «Эффект сообщения будет применён ровно один раз».

Первое — невозможно. Второе — рутинная инженерная задача, которую решают тысячи систем каждый день. Вся эта статья — про то, почему граница проходит именно здесь, и что конкретно нужно построить, чтобы жить на правильной стороне границы.

Ключевая мысль, которую стоит унести, даже если дальше вы ничего не прочитаете: гарантия доставки — это свойство не транспорта, а пары «транспорт + обработчик». 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

Никакой ошибки. Просто два бонуса вместо одного, и обнаружится это через месяц при сверке.

Ключ идемпотентности

Ключ идемпотентности — это идентификатор логической операции, стабильный между всеми её повторами. Три требования, каждое из которых нарушают регулярно:

  1. Стабильность. Ключ генерируется один раз в момент формирования намерения и переиспользуется всеми ретраями. uuid4() внутри функции retry() — самая частая ошибка в этой области: каждый повтор приносит новый ключ, дедупликация не срабатывает, а код при этом выглядит правильным.
  2. Детерминированная область действия. Ключ уникален в пределах пары (отправитель, тип операции). Один и тот же event_id может законно обрабатываться тремя разными обработчиками — значит, первичный ключ дедупликации это (handler, key), а не key.
  3. Привязка к содержимому. Если по тому же ключу пришёл другой 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, а не ждёт. Ожидание под блокировкой превращает всплеск ретраев в очередь на коннекты к БД и дальше в метастабильный отказ.

Если эффект физически не может быть в той же БД, что и отметка (эффект — вызов внешнего 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

Границы гарантии exactly-once в Kafka

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 обеспечивает exactly-once для состояния: распределённые снапшоты по алгоритму Chandy — Lamport (оригинальная статья 1985, реализация во Flink) записывают согласованный срез, и после сбоя оператор откатывается к нему. Для внешних приёмников этого мало, поэтому есть TwoPhaseCommitSinkFunction: запись предварительно фиксируется, а коммит происходит вместе с чекпойнтом. Требование к приёмнику — уметь транзакции или идемпотентную запись. Kafka-сток использует транзакции Kafka, файловый — атомарный rename, JDBC-сток без транзакций даёт at-least-once, что бы ни было написано в конфиге.

Spanner гарантирует линеаризуемость коммита, но не однократность его инициации: клиент, не получивший ответ, всё равно не знает, применилась транзакция или нет. Рекомендация Google в этой ситуации ровно та же, что у нас: делать транзакцию идемпотентной или проверять состояние по бизнес-ключу перед повтором.

Общее правило для чтения документации: находите фразу, определяющую границу системы, и приписывайте exactly-once только к тому, что внутри. Если границы не сформулировано — гарантии тоже нет.

Эффекты, которые нельзя откатить

Механика идемпотентности зависит от природы эффекта, а их всего три класса.

Двухшаговый протокол из ветки «нет ключа идемпотентности» честно описать так: он сужает окно, а не закрывает его. Между 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-контекст: тогда все повторы одной логической операции собираются в одну картину — см. наблюдаемость.

Типичные ошибки

  1. Генерировать ключ идемпотентности внутри функции ретрая. Каждая попытка приносит новый ключ, дедупликация не работает, код при этом выглядит корректно.
  2. Коммитить отметку и эффект раздельно. Окно дубля не исчезло, а переехало. Это распределённая транзакция с двумя участниками, со всеми последствиями.
  3. Считать PUT идемпотентным потому, что так написано в RFC. Идемпотентен обработчик, а не глагол.
  4. Не сохранять ответ. Повтор возвращает 404 или 409 там, где оригинал вернул 201, и клиент делает неверный вывод.
  5. Забыть про чистку таблицы дедупликации. Через год она больше основной таблицы, а проверка дубля — самый дорогой запрос в системе.
  6. Верить, что Kafka EOS покрывает вызов платёжного шлюза. Не покрывает: граница проходит по краю Kafka.
  7. Лечить переупорядочивание дедупликацией. Не лечится; нужны монотонные версии и фенсинг-токены.
  8. Ретраить 400 и 422. Пустая трата бюджета ретраев в момент, когда он нужнее всего.
  9. Обрабатывать дубликат как ошибку — с алертом и записью в error-лог. Дубликат при at-least-once — это норма; шум от него приучает игнорировать алерты.
  10. Писать в архитектурном решении «мы обеспечиваем 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 обязательны, иначе лечение станет причиной метастабильного отказа.

Источники

Что дальше

Мы несколько раз упирались в одну и ту же потребность: кому-то нужно выдать монотонный токен, кто-то должен решить, что лиза протухла и ключ можно перехватить, а зомби-инстанс после ребаланса нужно отсечь до того, как он применит устаревший эффект. Всё это — задачи координации: Координация: распределённые блокировки, лизы, etcd и ZooKeeper.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов