Классическая проблема распределенных систем: запись в БД и отправка в Kafka или RabbitMQ не атомарны. Если после коммита транзакции падает сеть или сервис, событие теряется навсегда. Паттерн Outbox решает это, сохраняя событие в той же транзакции, что и бизнес-данные.
Реализация с asyncpg и FOR UPDATE SKIP LOCKED
В бизнес-логике событие пишется в outbox-таблицу в рамках той же транзакции, что и основные данные. Несколько воркеров конкурентно читают строки через
FOR UPDATE SKIP LOCKED, что позволяет параллельно обрабатывать разные записи без deadlock'ов.# Вставка события внутри бизнес-транзакции
async with conn.transaction():
await conn.execute("INSERT INTO orders ...")
await conn.execute(
"INSERT INTO outbox (event_type, payload) VALUES ($1, $2)",
"order.created", json.dumps(data)
)
# Конкурентный воркер с batch-обработкой
async def outbox_worker(pool, broker):
while True:
async with pool.acquire() as conn:
rows = await conn.fetch(
"DELETE FROM outbox WHERE id IN ("
"SELECT id FROM outbox ORDER BY id LIMIT 10 FOR UPDATE SKIP LOCKED"
") RETURNING *"
)
for row in rows:
try:
await broker.send(row['event_type'], row['payload'])
except Exception:
await asyncio.sleep(1)
# Возвращаем запись обратно для ретрая
await conn.execute(
"INSERT INTO outbox (event_type, payload) VALUES ($1, $2)",
row['event_type'], row['payload']
)
await asyncio.sleep(0.1)
Ключевые trade-offs и типичные ошибки
* FOR UPDATE SKIP LOCKED обязателен для конкурентных воркеров. Без него получите сериализацию или deadlock.
* Не теряйте события при сбое отправки. Если брокер недоступен, запись должна быть либо возвращена в outbox, либо помечена для retry. Иначе данные потеряны.
* Читайте batch'ами. По одному событию за запрос — путь к перегрузке БД. Делайте
LIMIT 10..100.* Не используйте CDC (Debezium) вместо outbox без необходимости. CDC видит все изменения таблиц, включая временные состояния, и добавляет зависимость от Kafka Connect. Outbox чище для бизнес-событий.
Вывод: Outbox с asyncpg и FOR UPDATE SKIP LOCKED дает гарантию exactly-once доставки для микросервисов, но требует явного учета ретраев и batch-обработки, чтобы не стать узким местом системы.