Доставка exactly-once — миф; однократные эффекты реальны¶
Kafka не может гарантировать ровно одно обновление вашей БД. Без общей транзакции этого не может гарантировать ни один брокер: база и брокер — две системы с двумя фиксациями, и обещания брокера заканчиваются на его границе. Но можно получить точнее названный и столь же полезный результат: каждый эффект возникает один раз, сколько бы раз ни доставлялось вызвавшее его сообщение. В статье мы разберём все окна расхождения между фиксацией и публикацией, измерим их на настоящих PostgreSQL и Kafka и закроем по очереди: outbox у отправителя, inbox у получателя и ключ идемпотентности на границе HTTP.
Результаты получены экспериментальным скриптом статьи: он запускает PostgreSQL 17 и Kafka в контейнерах и считает строки и сообщения после каждого сценария. Версии: omni-box 0.2.0, aiokafka 0.14.0, SQLAlchemy 2.0.52, Python 3.13. Отправляющей стороне outbox посвящена отдельная статья; здесь рассматривается весь путь.
Сообщения и эффекты¶
«At-least-once», «at-most-once» и «exactly-once» описывают сообщения: сколько раз потребитель увидит запись. Это язык брокера, и гарантии действуют внутри него. Собственная exactly-once-семантика Kafka — транзакции между топиками с идемпотентным producer — реальна и ограничена Kafka. Цикл consume-transform-produce можно сделать однократным, если и вход, и выход — топики. Как только преобразование записывает строку PostgreSQL, гарантия заканчивается на границе JDBC или, в нашем случае, asyncpg.
Бизнес интересуют не сообщения, а однократное списание с клиента, создание счёта, отправка письма. Это эффекты. Поэтому вопрос к цепочке обработки: не «сколько раз доставлено сообщение?», а «сколько раз возник каждый эффект и возник ли вообще?». Это разные вопросы. На второй можно хорошо ответить даже без желаемого ответа на первый.
Измеряем каждое окно¶
Начнём с простейшего действия сервиса: записать заказ и сообщить о нём остальным. Две записи, две системы и процесс, способный завершиться между ними.
commit, then publish, crash between: orders=1 messages=0
publish, then commit, commit fails: orders=0 messages=1
Первая строка — заказ есть, события нет. Биллинг о нём не узнает; клиент получил заказ, по которому никто не выставит счёт. Во второй порядок обратный: событие есть, строки нет, счёт выставляется за несуществующий заказ. Никакой порядок двух записей не устраняет оба риска: вторая всегда может упасть после успешной первой. Это проблема двойной записи — свойство двух независимых фиксаций, а не ошибка конкретного кода.
Outbox закрывает сторону отправителя¶
Запишите событие в базу в одной транзакции со строкой заказа. Теперь фиксация одна: существуют либо и заказ, и событие, либо ничего. Отдельный relay читает ожидающие события, публикует их в Kafka и помечает завершёнными. Именно к нему переместились вторая фиксация и окно сбоя:
relay cycle 1, crash after the send: outbox rows [('pending', 0)] messages=1
relay cycle 2, normal: outbox rows [('completed', 0)] messages=2
В первом цикле relay отправил событие в Kafka и завершился до фиксации отметки о завершении. Строка ещё ожидает обработки, сообщение уже в топике. Во втором цикле relay находит строку и снова выполняет работу: в топике две копии. Ничего не потерялось, появился дубликат. В этом весь обмен: outbox превращает «может потеряться или продублироваться» в «не потеряется, но может продублироваться». Доставка по устройству at-least-once: отправка предшествует фиксирующей её транзакции, а сбой между ними вызывает повторную публикацию.
Здесь многие проекты останавливаются, оставляя примечание «потребители должны быть идемпотентными», которое каждая команда реализует по-своему.
Inbox закрывает сторону получателя¶
Половина паттерна на стороне потребителя — inbox: таблица с уникальным индексом по (message_id, consumer_group). Для каждого обрабатываемого сообщения потребитель вставляет строку в одной транзакции с его эффектом. Повторная доставка сталкивается с уже существующей строкой; конфликт служит сигналом пропустить эффект:
delivery 1: processed=True duplicate=False committed=True invoices=1
delivery 2: processed=False duplicate=True committed=True invoices=1
Пришли обе копии, счёт один. Смещение второй доставки зафиксировано у брокера, поэтому её больше не доставляют, а обработчик для неё не запускался. Окно дедупликации равно сроку жизни строки inbox, который должен превышать срок хранения сообщений брокером. Ключ включает группу потребителей, поэтому каждая группа получает собственную однократную обработку.
Всё держится на слове транзакция. Строка inbox и счёт должны фиксироваться вместе. Если обработчик записывает счёт через собственную session, а inbox фиксируется отдельно, сбой между ними даст либо счёт без записи об обработке — и повтор создаст второй, либо запись без счёта — и повтор будет пропущен навсегда. Обработчик обязан писать в транзакцию, которую открыл runner, а библиотека должна передать её ему:
async def create_invoice(event: InboxEvent, repo: InboxEventRepository) -> None:
await repo.session.execute(invoices.insert().values(order_id=event.payload["order_id"]))
repo.session — именно транзакция строки inbox. Runner вставляет строку, запускает обработчик, фиксирует оба изменения, затем смещение брокера. Если обработчик выбрасывает исключение, строка откатывается вместе с его изменениями, и брокер доставляет сообщение снова. Повтор начинает с чистого состояния, как и должен.
Отметка inbox и бизнес-запись фиксируются одной транзакцией БД. При повторной доставке бизнес-запись пропускается; подтверждение Kafka идёт после коммита.
Граница HTTP¶
Есть ещё одно окно до всей этой цепочки: клиент, отправивший заказ, получил таймаут и повторил запрос. Два заказа, два события outbox, два счёта — всё обработано правильно, но всё продублировано. Outbox и inbox этого не распознают: для них это разные заказы. Дедупликация нужна на входе по клиентскому ключу идемпотентности; этому посвящена статья об идемпотентности. Три точки дедупликации — по одной на каждую границу запроса.
Весь путь¶
HTTP-запрос ── ключ идемпотентности: тот же запрос обрабатывается один раз
│
PostgreSQL ── одна транзакция: бизнес-запись и строка outbox
│
ретранслятор ── минимум один раз: публикация, затем отметка completed
│
Kafka ── доставляет каждое сообщение минимум один раз
│
inbox ── (message_id, consumer_group): то же сообщение обрабатывается один раз
│
консьюмер ── одна транзакция: строка inbox и результат обработки
Каждая стрелка на схеме означает at-least-once. Каждый блок — место, где однократность обеспечивается уникальным ключом внутри одной транзакции. Гарантия всей цепочки: каждый эффект возникает ровно один раз. Она построена из доставки at-least-once и трёх ключей, поэтому выдерживает сбои на каждой стрелке, а отдельным компонентам не приходится обещать больше, чем они могут.
Почему библиотека не должна владеть транзакцией¶
На каждом шаге выше сказано «в одной транзакции», и в каждой есть ваши данные: заказ или счёт. Если библиотека сама открывает и фиксирует отдельную транзакцию, строка outbox окажется в одной, а заказ — в другой: та же двойная запись, только сложнее. Поэтому репозиторий outbox принимает вашу session и пишет в неё; чтение, публикация и отметка relay выполняются в транзакции, которую открываете и фиксируете вы; runner inbox организует транзакцию и передаёт её обработчику. Библиотека не может отделить свои записи от вашей транзакционной границы: гарантия существует именно в ней.
Как выглядит код¶
Отправитель: строка и событие вместе:
async with session_factory() as session, session.begin():
await session.execute(orders.insert().values(id=order_id))
await PostgresOutboxRepository(session, model_class=OutboxEventDB).create(
domain.create_outbox_event(
aggregate_type="order", aggregate_id=order_id, event_type="order.created",
topic="orders.events", partition_key=str(order_id), payload={"order_id": str(order_id)},
idempotency_key=f"order.created:{order_id}",
)
)
Один цикл relay в транзакции, которой управляете вы:
async with session_factory() as session, session.begin():
repo = PostgresOutboxRepository(session, model_class=OutboxEventDB)
await OutboxPublisher(repo, KafkaEventPublisher(producer, EnvelopeEventConverter())).publish_batch(worker_id="relay-1", batch_size=100)
И потребитель, чей обработчик пишет через переданную транзакцию:
runner = InboxConsumerRunner(
consumer=KafkaEventConsumer(kafka_consumer),
transaction_provider=InboxTxProvider(session_factory), # opens the transaction, yields the repository
handler=create_invoice, # writes through repo.session
worker_id="billing-1",
consumer_group="billing",
)
Эти три фрагмента соответствуют трём измерениям выше по порядку. Это omni-box: примитивы outbox и inbox поверх внешней транзакции. На входной границе используется idempotency-kit. Ни одна библиотека не обеспечивает однократную доставку. Вместе с вашими транзакциями они обеспечивают однократный эффект — то, ради чего exactly-once и требовалась.
Суть — во второй таблице: два сообщения, один счёт.