Soumabrata

Engineering 10 min read June 2026

The Transactional Outbox Pattern
Why Fire-and-Forget Is a Fire Hazard

How replacing an async Kafka publish with a database-level outbox table turned a fragile dual-write into a crash-resilient, at-least-once delivery pipeline โ€” and the two subtle bugs we found along the way.

1. 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

Gateway Database Kafka Gateway Database Kafka 1 Add event 2 Commit transaction 3 Create task 4 Publish event alt crash occurs here 5.1 Event lost
Before: Gateway handles DB writes and async publishing. A crash between step 2 and 4 leaves the DB updated but Kafka missing the message.

After: Gateway Only Writes to DB

Gateway DB Outbox Kafka Gateway DB Outbox Kafka 1 Store event 2 Add to outbox 3 Commit transaction Loop continuous polling 4.1 Poll for pending messages opt messages found 4.2.1 Publish event 4.2.2 Mark as completed
After: Gateway writes to database transactionally (event + outbox record). A background worker polls for pending outbox records, publishes to Kafka, and marks as completed in the 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()
The Event and OutboxRecord are 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:

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

  1. Fire-and-forget is a fire hazard. Any asyncio.create_task that 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.
  2. 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.
  3. FOR UPDATE SKIP LOCKED is 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.
  4. SQLAlchemy defaults are deceptive. default=uuid.uuid4 on a mapped column does NOT set the attribute at construction time. If you need the PK before the object is flushed, generate it yourself.
  5. 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.
  6. 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).