Outbox cleanup — pick your poison

Three cleanup strategies trade WAL noise, race risk, and operational complexity; INSERT-then-DELETE in the same transaction collapses cleanup into the write path.

Previously

The outbox closes the dual-write gap — but every business write now also writes a row whose entire purpose is to be deleted, so the table will grow forever unless we pick a cleanup strategy. Three strategies are on the table; each loses something different.

Scene 09

Outbox cleanup — pick your poison

  1. Watch
  2. Try it
  3. Predict
  4. Capture
Outbox cleanup · pick your poison · scrub t = 0dt=01h1d1w30dlane 1Hard-delete after emitoutbox rows · 1,500WAL noise · 18000/hr op=dtopic size · 1h / 24h / 30d1h24h30dspurious op=d events; race window ifconsumer lagslane 2Tombstone + log compactionoutbox rows · 600WAL noise · 0/hr op=dtopic size · 1h / 24h / 30d1h24h30dlane 3Partition by date, drop month…outbox rows · 0WAL noise · 0/hr op=dtopic size · 1h / 24h / 30d1h24h30doperationally simple, but needs adate-partitioned table from day one
What to watch for

Three lanes, one workload. Watch each lane age across 30 days — the outbox row count, the op=d noise in the WAL, and the Kafka topic size diverge. Lane 1 stays small but pumps op=d events; lane 2 stays small via tombstones + compaction; lane 3 climbs daily and drops at month-end.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Connector.hardDeleteAfterEmit
naive lane 1: connector emits, then deletes the row
def hardDeleteAfterEmit(row):
kafka.produce(
topic = route(row),
key = row.aggregate_id,
value = row.payload,
)
# race: if the consumer lags, op=d may arrive
# before the consumer has read op=c
db.execute(
'DELETE FROM outbox WHERE id = %s',
row.id,
)
# WAL now records op=d; Outbox SMT must filter
Connector.tombstoneAndCompact
lane 2: emit a null-valued tombstone; Kafka compacts per key
# topic configured cleanup.policy=compact
def tombstoneAndCompact(key):
kafka.produce(
topic = outbox_topic,
key = key,
value = None, # tombstone
)
# log compaction will eventually drop prior
# values for this key — per-key 'keep latest'
# race: lagging consumer may miss the value
# before compaction reclaims it
Service.insertDeleteSameTx
lane 1 collapsed: cleanup folded into the write path
def insertDeleteSameTx(event):
db.execute('BEGIN')
db.execute(
'INSERT INTO outbox(...) VALUES (...)',
event,
)
db.execute(
'DELETE FROM outbox WHERE id = %s',
event.id,
)
db.execute('COMMIT')
# WAL has both records; Debezium emits op=c
# then op=d atomically — sink ignores op=d

Where this sits in Build a CDC pipeline (Debezium + outbox)

Scene 09 of 12. Hard-delete after emit, tombstone + log compaction, partition-by-date drop — three trade-offs across WAL noise, race risk, and operational complexity. INSERT+DELETE same-tx collapses cleanup into the write.

Up next. Cleanup is sorted; the outbox emits one event per business action — but those events are about to be sharded across Kafka partitions, and the question of which events stay in order matters as soon as there is more than one.

All 12 scenes in Build a CDC pipeline (Debezium + outbox) · Every curriculum

Built with Arqly
Every scene in Build a CDC pipeline (Debezium + outbox) builds on the one before it.All 12 Build a CDC pipeline (Debezium + outbox) scenes