Contents
1. The Problem: The Async Fire-and-Forget 2. What Could Go Wrong 3. The Transactional Outbox Pattern 4. Implementation 5. The event.id Gotcha 6. The Outbox Relay Worker 7. Verifying Resilience 8. The Cost Question 9. Lessons Learned1. The Problem: The Async Fire-and-Forget
The original gateway handler followed a straightforward pattern: accept the webhook, validate it, persist the event to PostgreSQL, then fire off an async task to publish to Kafka.
# Original code (simplified)
db.add(event)
await db.commit()
async def _publish():
producer = get_kafka()
if producer is None:
return
try:
await producer.send_and_wait(...)
# update status to 'queued'
except Exception:
# update status to 'failed'
asyncio.create_task(_publish()) # ๐จ fire-and-forget
At a glance this looks fine. The event is safely in PostgreSQL. The
Kafka publish happens asynchronously. The 202 Accepted
response is not blocked. But there's a critical vulnerability:
๐จ Crash after db.commit() but before Kafka publish โ event lost silently
The event is persisted as pending. The background task
never runs. Kafka never receives the message. The destination never
gets the webhook. And nobody knows โ the event sits in the database
with a status of pending forever.
2. What Could Go Wrong
Four distinct failure modes, ordered by severity:
| Scenario | What Happens | Silent Data Loss? |
|---|---|---|
Gateway crashes right after db.commit() |
Async task never created or killed mid-flight | โ Yes |
| Kafka broker is unreachable | send_and_wait raises โ status = 'failed' set, but one-shot retry with
30s delay only |
No, but needs manual replay |
| Startup race โ Kafka starts after gateway | init_kafka() returns None, publish is silently skipped |
โ
Yes (if producer is None, _publish just returns) |
| Multiple gateway replicas | Each replica creates competing background tasks for the same events | Duplicate deliveries |
The fundamental problem is the dual-write pattern: one write to PostgreSQL, one write to Kafka, coordinated only by application code in a single process. If that process dies between the two writes, one of them is lost.
3. The Transactional Outbox Pattern
The outbox pattern eliminates the dual-write by making the Kafka publish someone else's problem.
Before: Gateway Owns Both Writes
After: Gateway Only Writes to DB
4. Implementation
The OutboxRecord Model
A single new table tracks all events waiting to be published:
class OutboxRecord(Base):
__tablename__ = "outbox_records"
id: Mapped[uuid.UUID] # PK
event_id: Mapped[uuid.UUID] # FK โ events.id
status: Mapped[str] # pending | completed | failed
publish_key: Mapped[str] # endpoint_id (kafka key)
publish_topic: Mapped[str] # kafka topic
attempts: Mapped[int] # retry counter
last_error: Mapped[str | None] # last failure message
created_at: Mapped[datetime]
processed_at: Mapped[datetime | None]
The Gateway Handler (After)
The refactored handler is simpler and safer โ no Kafka imports, no background tasks:
# Refactored handler โ single transaction, zero Kafka dependency
event_id = uuid.uuid4()
event = Event(
id=event_id, # explicit! see ยง5
endpoint_id=endpoint.id,
request_body=payload,
request_headers=headers,
status="pending",
)
db.add(event)
outbox = OutboxRecord(
event_id=event_id,
publish_key=str(endpoint.id),
publish_topic=settings.kafka_topic_raw_events,
)
db.add(outbox)
await db.commit()
TheEventandOutboxRecordare created in the same database transaction. Either both are persisted, or neither is.
5. The event.id Gotcha
This one cost us a NOT NULL constraint violation on the
first run. The foreign key from outbox_records.event_id
to events.id was being set to NULL.
# โ BROKEN โ SQLAlchemy default=uuid.uuid4 fires at flush, not construction
event = Event(
endpoint_id=...,
request_body=...,
)
db.add(event) # event.id is None here
outbox = OutboxRecord(
event_id=event.id, # ๐ฅ FK = NULL โ constraint violation
)
SQLAlchemy's default=uuid.uuid4 on the column definition
is only invoked at flush time, not at object
construction. When OutboxRecord(event_id=event.id) runs,
event.id is None.
# โ
FIXED โ generate the ID explicitly before construction
event_id = uuid.uuid4() # now event_id is available โฆ
event = Event(
id=event_id, # โฆ before the Event object exists
endpoint_id=...,
request_body=...,
)
db.add(event)
outbox = OutboxRecord(
event_id=event_id, # โ
valid FK
)
Rule of thumb: When you need an object's PK before it's flushed, generate the UUID explicitly. Don't rely on SQLAlchemy's column default.
6. The Outbox Relay Worker
The relay is a standalone process that runs in a tight loop:
# Poll loop (simplified)
while not _stop.is_set():
async with async_session_factory() as db:
result = await db.execute(
select(OutboxRecord)
.where(OutboxRecord.status == "pending")
.order_by(OutboxRecord.created_at)
.limit(BATCH_SIZE) # โ 10
.with_for_update(skip_locked=True) # โ key detail
)
records = result.scalars().all()
for record in records:
try:
await _publish(record)
record.status = "completed"
except Exception as exc:
record.attempts += 1
record.last_error = f"{type(exc).__name__}: {exc}"
if record.attempts >= 10:
record.status = "failed" # give up
await db.commit()
Key Design Decisions
| Decision | Choice | Rationale |
|---|---|---|
| Concurrency control | FOR UPDATE SKIP LOCKED |
PostgreSQL-native row-level locking. Two relay workers won't fight over the same row.
SKIP LOCKED means a worker simply skips rows locked by another worker instead of
waiting. |
| Poll interval | 1 second | Balance between latency and DB load. A webhook relay has sub-second latency requirements for most use cases. |
| Batch size | 10 | Small enough that a single transaction completes quickly. Large enough to handle burst traffic without excessive round trips. |
| Max attempts | 10 | Arbitrary but generous. Each failed attempt records the error in last_error for
debugging. After 10, the record is failed permanently and needs manual
intervention. |
| Transaction scope | One transaction per batch, commit() at end |
Lock is held during DB operations only (select, then update). Network calls (Kafka publish) happen outside the lock โ the row is locked, published, then updated and committed. If the worker crashes between publish and commit, the next poll picks up the still-pending record and re-publishes. |
Kafka Unavailability
If Kafka is down at relay startup, init_kafka() returns
None. The relay exits its poll loop early:
producer = await init_kafka()
if producer is None:
logger.warning("Kafka unavailable โ outbox records will queue")
return
Records stay at pending in PostgreSQL. When Kafka comes
back, a container restart picks them up. No data loss.
7. Verifying Resilience
We tested three crash scenarios against the live system:
Scenario A: Normal Flow
curl -X POST /hooks/<endpoint_id> \
-H "X-Hub-Signature-256: sha256=<sig>" \
-d '{"action":"push"}'
# โ {"status": "accepted", "event_id": "0df9e4db-..."}
# Outbox record:
outbox_records: "completed"
delivery_worker: delivered to collector โ 200 (39ms)
Scenario B: Kill Outbox Relay Mid-Stream
docker stop outbox-relay
# Send event while relay is down
curl -X POST /hooks/<endpoint_id> ...
# โ {"status": "accepted", "event_id": "4e28b9ca-..."}
# Outbox record stays pending:
outbox_records: "pending" โ queued in DB
docker start outbox-relay
# After 1-2 seconds:
outbox_records: "completed" โ drained
delivery_worker: delivered โ 200 (11ms)
Scenario C: Kafka Crash at Relay Start
# Kafka goes down โ relay starts โ init_kafka() returns None # โ relay exits โ records stay pending # Kafka comes back โ restart relay # โ picks up pending records โ publishes โ marks completed
8. The Cost Question
"Isn't a DB table costly if the service goes boom?"
The outbox table is the cheapest insurance policy you can buy. Here's the breakdown:
| Resource | Cost Per Event | Annual Cost (10M events) |
|---|---|---|
| Storage | ~200 bytes (one row in outbox_records) |
~2 GB โ negligible |
| CPU | 1 ร SELECT FOR UPDATE + 1 ร UPDATE |
Trivial โ two indexed lookups |
| Poll overhead | 1 query/second when idle | ~86K queries/day, all hitting an empty table โ index-only scan |
| Network | Zero โ relay is on the same Docker network | โ |
| Complexity | One new table + one worker process | ~300 lines of code total |
Compare this to the cost of losing events:
- Stripe webhook missed โ customer not notified of payment failure โ chargeback fees
- GitHub push event lost โ CI pipeline doesn't run โ production bug deployed
- Slack alert dropped โ on-call engineer not paged โ incident response time blows up
The outbox pattern trades ~200 bytes of disk and one extra process for provable at-least-once delivery guarantees. It's not costly โ it's insurance.
9. Lessons Learned
- Fire-and-forget is a fire hazard. Any
asyncio.create_taskthat touches external infrastructure is a gamble. The process can die between the DB commit and the network call. The outbox pattern makes this gap disappear. - Dual-write is the root of most distributed systems bugs. If you need to write to two systems atomically, you need either a distributed transaction (expensive, fragile) or an outbox (simple, proven). Choose the outbox.
FOR UPDATE SKIP LOCKEDis a superpower. PostgreSQL 9.5+ lets you implement a reliable work queue with zero infrastructure. No Redis streams, no RabbitMQ โ just a SQL query that's safe for concurrent consumers.- SQLAlchemy defaults are deceptive.
default=uuid.uuid4on a mapped column does NOT set the attribute at construction time. If you need the PK before the object is flushed, generate it yourself. - Test the crash, not just the happy path. Kill containers mid-operation. Verify events queue in the database. Restart and watch them drain. A system that works perfectly when everything is up is not a production system.
- Decoupling is not free โ but it's cheaper than recovery. The outbox table and relay worker added ~300 lines of code and one container. The alternative โ recovering lost events from database logs or customer complaints โ is infinitely more expensive.
The full source code is available on GitHub.
The outbox relay lives in workers/outbox_relay.py (112 lines).