Why We Picked a Polling Outbox Over Debezium for 4,000 Writes per Second
How we weighed a polling transactional outbox against Debezium CDC for order events, with the tradeoffs and production numbers.
Last spring, the order service started dropping events during a load test. The database commit succeeded, but the Kafka publish sat behind a retry loop that gave up after three attempts. Inventory never reserved the stock, and the order showed as paid in one system and pending in another. That incident settled the question we had been circling for weeks: order state and emitted events had to commit together, and we needed a mechanism that made that true by construction.
This post is the decision record for how we solved it. We chose a polling transactional outbox over log-based change data capture with Debezium. Both are reasonable. The right answer depended on our write volume, our team size, and what our consumers expected from us.
Context and constraintsLink to this section
The order service runs on PostgreSQL 14, with a single primary and two streaming replicas. Peak write load is about 4,000 transactions per second, and each order transaction typically touches the orders table and one or two child tables. After fan-out, the service emits roughly 5,000 to 6,000 domain events per second at peak.
The hard requirement is atomicity. An order state change and the event describing it must either both commit or both roll back. A dual write, where we update the database and then call Kafka, was already ruled out after the incident above.
Downstream, we have around a dozen Kafka consumers: inventory, fulfillment, billing, fraud, analytics, and two external-facing webhook relays. Most of them already deduplicate by event ID, but their assumptions about latency and ordering differ. Inventory reservation needs p99 publish lag under one second. Analytics tolerates minutes. Within a single order, consumers expect events in the order the state changed, so OrderPaid must not arrive before OrderPlaced.
The platform team that would own this is three engineers, and they also run the Kafka cluster, the service mesh, and on-call for about forty other services. Any component that adds a new stateful system to operate carries real weight for us. Operating a Kafka Connect cluster with Debezium is not impossible, but it is another thing to upgrade, monitor, and debug at 2 a.m.
Options consideredLink to this section
Polling transactional outboxLink to this section
The design is simple. Every transaction that changes order state also inserts a row into an outbox table. A separate relay process reads batches of unpublished rows, publishes them to Kafka, and deletes them once the broker acknowledges.
The table is intentionally small and explicit:
CREATE TABLE outbox (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
event_id uuid NOT NULL,
aggregate_type text NOT NULL,
aggregate_id text NOT NULL,
event_type text NOT NULL,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE INDEX outbox_created_at_idx ON outbox (created_at);The pros are concrete. The event shape is whatever we decide to write, so OrderPaid carries curated fields rather than a raw row diff. This matches the point made in [1] that CDC events describe row changes, and consumers then have to reconstruct business meaning. The schema of orders can change without changing the event contract. Failure modes are visible in our own code and our own tables. There is no replication slot holding WAL on the primary.
The cons are real. Polling adds latency, bounded by the poll interval. It adds read and delete load on the primary. Multiple relay instances can race, so you need a coordination strategy. And the relay is our code, which means we own its bugs. [4] notes the race conditions when scaling horizontally and the DB overhead, and I agree with both. Its rough guidance that polling suits low volume is consistent with [1], which places outbox-plus-polling at the default choice for services up to about 10,000 events per second, with the caveat that the team has to accept the polling cost.
Debezium log-based CDCLink to this section
The alternative is to skip polling entirely. Debezium connects to PostgreSQL through a logical replication slot, reads the WAL using pgoutput, and streams changes to Kafka through Kafka Connect. Teams often combine this with an outbox table, so the application writes a clean event row and Debezium streams it, as described in [2] and [3].
The pros are latency and load. Debezium does not poll, so it adds no periodic query load to the primary, and end-to-end lag is typically sub-second once tuned. The outbox pattern works well with it, and [3] lists the reduced maintenance burden and the transactional guarantees as major advantages.
The cons hit us harder. A logical replication slot pins WAL on the primary until the consumer confirms it. If the connector stalls, WAL accumulates. At 4,000 writes per second, that can fill a disk in hours, and a stuck slot on the primary is an outage risk for the whole database, not just for events. PostgreSQL 13 and later support max_slot_wal_keep_size to cap this, but capping it means the slot gets invalidated and you lose the stream, which requires a resnapshot. Schema changes are another concern: altering the outbox or orders tables requires coordinating connector behavior, and the connector and its converters become a versioned dependency we would need to track. And [3] itself notes the added latency from reading the log and going through a broker, which requires tuning.
The decision and technical rationaleLink to this section
We chose the polling outbox. The reasoning came down to three constraints.
First, the operational blast radius. A stalled Debezium slot threatens the primary's disk, and that is a failure we would rather not be able to cause from the events pipeline. A stalled relay, by contrast, only delays events. The database keeps working.
Second, the team. Three people owning Kafka, the mesh, and this pipeline meant we valued a component whose failure modes we could read in a single file. The relay is about 400 lines of Go.
Third, the latency requirement was achievable. Inventory needed p99 under one second, and polling at 100 ms intervals with batching gives us that with room to spare. We were not operating in the range where polling latency becomes the bottleneck.
What we traded away is worth naming. We accepted extra read and delete load on the primary, and we accepted that the relay's correctness is our responsibility. We gave up the ability to capture changes that bypass the application, such as manual SQL fixes, unless they also write outbox rows. We also accepted that event delivery is at-least-once, so consumers must deduplicate. [4] makes the same point, and it applies here.
The relay uses FOR UPDATE SKIP LOCKED to claim rows, and it deletes them only after the Kafka acknowledgement:
BEGIN;
SELECT id, event_id, aggregate_id, event_type, payload
FROM outbox
ORDER BY id
LIMIT 500
FOR UPDATE SKIP LOCKED;
-- publish the batch to Kafka with acks=all, wait for all acks
DELETE FROM outbox WHERE id = ANY($1); -- $1 = ids published
COMMIT;We run one active relay at a time. A second instance waits on a PostgreSQL advisory lock and takes over if the first one dies. SKIP LOCKED keeps the query safe if a handover overlaps briefly, and it lets us scale out later by partitioning on aggregate_id. Deleting on publish, rather than marking rows, means the table stays small and has no dead tuples piling up from status updates. There is no cursor either. Because the relay always selects the lowest remaining IDs, a row that commits late, after a higher ID was already published, is still picked up on the next poll. Per-order ordering holds because all events for one order are written under the order row lock, and only one relay publishes them.
The poll settings are a batch size of 500 and a 100 ms interval when the last poll returned rows, falling back to 250 ms when it came back empty. Those numbers came from load testing, not theory.
The rejected option would become the better choice under different conditions: sustained volume past roughly 10,000 events per second, a team with a dedicated data platform function that already runs Kafka Connect, a need to capture changes that bypass the application, or a requirement for sub-100 ms lag that polling cannot meet.
Results after shippingLink to this section
Measured from commit timestamp to Kafka broker acknowledgement over the first four months at production load:
- p50 publish lag: 38 ms
- p99 publish lag: 210 ms
- Sustained relay throughput: about 6,100 events per second at peak, and roughly 14,000 per second when draining a backlog
- Primary CPU impact: about 3 to 4 percentage points, mostly from the outbox insert and the relay's select and delete. The unindexed scan we initially had was the largest single cost, and the
created_atindex plus the delete-on-publish design brought it down. - Duplicate event rate: about 0.02 percent of events were redelivered, almost all after relay restarts between publish and commit. Every one was absorbed by the consumer dedupe on
event_id.
IncidentsLink to this section
We had two notable ones. In month two, a relay pod was OOM-killed while holding a batch of 500 rows. A new upstream change had added line-item arrays to some payloads, and a few batches exceeded 40 MB. The backlog reached about 1.4 million rows and the oldest unpublished row was nine minutes old before we noticed. We fixed it by capping batches by total payload bytes as well as row count.
In month three, autovacuum fell behind on the outbox table. Deletes create dead tuples, and the default thresholds were too lax for a table that churns at this rate. We set per-table storage parameters, autovacuum_vacuum_scale_factor = 0.01 and a higher autovacuum_vacuum_cost_limit, and the problem went away.
The stalled relay failure modeLink to this section
The failure we planned for most is a stalled relay. If the relay stops, rows accumulate. The table grows, autovacuum has more work, and the created_at index grows with it. Consumers see lag rise, but the order service keeps accepting writes, which is the behavior we wanted. The risk is that a long stall turns into a slow database problem, so we alert on the age of the oldest row rather than on row count:
SELECT COALESCE(EXTRACT(EPOCH FROM now() - min(created_at)), 0) AS backlog_age_seconds
FROM outbox;We page at 60 seconds of backlog age and warn at 15 seconds. We also alert on the relay's own heartbeat, because a relay that is alive but stuck in a Kafka retry will not show up in backlog age until the backlog is already large.
TakeawaysLink to this section
- Pick the outbox pattern based on who will operate it. A relay you can read in one file is easier to own than a connector whose failure can fill the primary's disk.
- Treat replication slots as a database liability. If you choose CDC, cap WAL retention and alert on slot lag before you ship.
- Alert on the age of the oldest unpublished row, not on row count. Age tells you how stale consumers are, and count depends on your write volume.
- Delete on publish instead of marking rows. It keeps the outbox small and avoids dead-tuple churn, but you must tune autovacuum for that table.
- Make consumers idempotent from day one. At-least-once delivery is the price of atomic commits, and deduplicating on
event_idcosts very little compared to the alternative.
SourcesLink to this section
- Outbox: Transactional Outbox vs Change Data Capture (CDC) — when each?
- The Transactional Outbox Pattern: Transforming Real-Time Data Distribution at SeatGeek
- Distributed Data for Microservices — Event Sourcing vs. Change Data Capture
- Transactional Outbox: PostgreSQL Transactions, Routing and Kafka Replay