Real‑Time Change Data Capture Pipelines with Debezium, Kafka, and the Outbox Pattern: Architecture, Exactly‑Once Guarantees, and Scaling Strategies
Real‑Time Change Data Capture Pipelines with Debezium, Kafka, and the Outbox Pattern: Architecture, Exactly‑Once Guarantees, and Scaling Strategies
Debezium change data capture: Real‑Time Pipelines with Kafka, Streams, and the Outbox Pattern
Debezium streams database WAL events into Kafka, letting microservices react instantly.
Combined with the outbox pattern, it guarantees exactly‑once delivery across service boundaries.
Debezium change data capture provides sub‑second, exactly‑once data pipelines by streaming WAL events to Kafka.
Debezium change data capture – Introduction
Real‑time inventory updates and order state sync have become non‑negotiable after Google’s AI‑driven shopping test on Flipkart.
Insurers losing 942 M because stale data feeds inaccurate risk models illustrate the cost of lag.
At the same time, Crusoe’s 1.25 B AI data‑center pull‑back shows the industry’s appetite for energy‑efficient pipelines.
In this context, a CDC‑first architecture solves three problems:
- Latency – sub‑second propagation from OLTP to downstream services.
- Consistency – exactly‑once semantics despite retries, crashes, or network partitions.
- Scalability – horizontal growth on Kubernetes without manual sharding. Debezium change data capture taps the transaction log of MySQL, PostgreSQL, or MongoDB, converting each commit into a Kafka record.
Kafka retains the order, provides durable storage, and lets Kafka Streams enrich or filter events before they reach consumers.
The outbox pattern sits on the producer side: business services write domain events to an “outbox” table in the same transaction that updates their core tables. Debezium then treats the outbox as just another source, guaranteeing atomicity between state change and event emission.
Why traditional approaches fall short
- Polling introduces seconds‑to‑minutes of drift and adds load on the source DB.
- Application‑level publishing risks lost events when the DB commit succeeds but the message send fails.
- Dual‑writes break atomicity, leading to eventual consistency bugs that are hard to debug.
Desired guarantees
| Guarantee | Why it matters | How CDC + Outbox achieves it |
|---|---|---|
| Exactly‑once delivery | Prevent duplicate billing or inventory counts | Kafka’s idempotent producer + Debezium’s transactional source |
| Order preservation per entity | Guarantees correct state reconstruction | WAL ordering + partition key = entity ID |
| Fault‑tolerant replay | Enables rapid disaster recovery | Kafka retention + compacted outbox topic |
| Schema evolution without downtime | Supports agile microservice development | Avro + Confluent Schema Registry |
What is Debezium change data capture?
Debezium is an open‑source connector suite for Apache Kafka that reads database change logs (WAL, redo logs, etc.) and emits a stream of change events.
Each event contains:
- before and after snapshots of the row.
- source metadata (transaction ID, timestamp, binlog position).
- op code (
cfor create,ufor update,dfor delete). The connector runs as a Kafka Connect worker, handling offsets, schema registration, and fault‑tolerant restarts automatically.
How does the outbox pattern complement CDC?
- Transactional write – Service updates its domain tables and inserts a row into
outbox_eventsin the same DB transaction. - Debezium captures – The outbox table appears as another source; its rows become Kafka messages.
- Consumer reads – Downstream services consume from the outbox topic, process idempotently, and acknowledge. Because the outbox write shares the same transaction as the state change, there is no window where one succeeds and the other fails.
Problem Statement & System Architecture
Enterprises need a pipeline that:
- Propagates every OLTP change to any number of microservices.
- Guarantees exactly‑once semantics end‑to‑end.
- Scales to tens of thousands of TPS on Kubernetes.
- Offers observability for latency, lag, and schema drift.
Core components
- Source DB – MySQL 8.x with GTID enabled.
- Kafka Cluster – 3‑node KRaft mode, replication factor 3, min‑insync replicas = 2.
- Kafka Connect – Debezium MySQL connector, distributed mode, offset storage in Kafka.
- Outbox Table –
outbox_events(id UUID PK, aggregate_id VARCHAR, type VARCHAR, payload JSONB, created_at TIMESTAMP). - Kafka Streams Application – Enriches events, applies business rules, writes to domain‑specific topics.
- Consumer Services – Spring Boot, Quarkus, or Go microservices that read from their topic with exactly‑once processing enabled.
- Observability Stack – Prometheus + Grafana dashboards for connector lag, stream processing throughput, and consumer lag.
- Schema Registry – Confluent Schema Registry storing Avro schemas for all topics.
Data flow
flowchart TD
A[Source DB] -->|WAL| B[Debezium Connector]
B -->|Kafka Connect| C[Kafka Topics]
C -->|Enrichment| D[Kafka Streams]
D -->|Domain Topics| E[Microservice Consumers]
subgraph Outbox
A -->|Transactional INSERT| O[outbox_events]
O -->|Debezium reads| C
endScaling strategy
- Connector scaling – Deploy multiple tasks; each task handles a subset of tables or partitions.
- Kafka Streams parallelism – Set
num.stream.threadsto match pod CPU limits; use stateless processors for linear scaling. - Kubernetes autoscaling – Horizontal Pod Autoscaler (HPA) reacts to
kafka.connect.task.max.poll.recordsandstream.processing.ratemetrics from Prometheus. - Back‑pressure handling – Enable
consumer.max.poll.interval.msandproducer.delivery.timeout.msto avoid OOM during spikes.
Architecture Patterns Comparison
| Pattern | Latency (ms) | Exactly‑once? | Operational complexity | Typical use‑case |
|---|---|---|---|---|
| Polling + REST | 500‑2000 | No (at‑least‑once) | Low – simple cron jobs | Low‑volume sync |
| Dual‑write + Message Queue | 100‑300 | No (window of inconsistency) | Medium – transaction management | Event‑driven legacy |
| CDC + Direct Topic | 20‑80 | Yes (Kafka idempotence) | High – connector ops, schema mgmt | High‑throughput microservices |
| CDC + Outbox + Streams | 10‑50 | Yes (transactional + idempotent) | Highest – needs DB schema, streams code | Mission‑critical, multi‑service pipelines |
About the Author
I'm a senior backend engineer with a focus on real‑time data pipelines and Kubernetes‑native deployments.
Feel free to reach out with questions or consulting inquiries: Contact Me.
Step-by-Step Implementation Guide
You've got the theory down. Now let's build the actual pipeline. We start with the database, move through Kafka, and end with the consumer that guarantees consistency.
1. Configure the Database for Logical Decoding
Debezium doesn't poll tables. It reads the Write-Ahead Log (WAL). You need to enable logical decoding in your PostgreSQL instance.
-- In postgresql.conf
wal_level = logical
max_replication_slots = 10
max_wal_senders = 10
-- Create a publication
CREATE PUBLICATION dbz_publication
FOR TABLE orders, customers, inventory;
-- Create a user with replication privileges
CREATE USER debezium WITH REPLICATION LOGIN PASSWORD 'secure_pwd';The wal_level = logical setting is non-negotiable. Physical replication uses a different format. Debezium needs the logical format to extract row-level changes.
Set max_replication_slots to match your expected number of Debezium instances. Too low, and you'll hit connection errors under load.
2. Deploy the Debezium Connector
Use the Kafka Connect REST API or a containerized deployment. Here’s a Docker Compose snippet for a single-node setup.
version: '3.8'
services:
postgres:
image: postgres:14
environment:
POSTGRES_PASSWORD: postgres
command: >
postgres
-c wal_level=logical
-c max_replication_slots=10
kafka:
image: confluentinc/cp-kafka:7.4.0
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
connect:
image: debezium/connect:2.4
ports:
- "8083:8083"
environment:
BOOTSTRAP_SERVERS: kafka:9092
GROUP_ID: 1
CONFIG_STORAGE_TOPIC: connect_configs
OFFSET_STORAGE_TOPIC: connect_offsets
STATUS_STORAGE_TOPIC: connect_statuses
depends_on:
- kafkaRegister the connector via REST API. Note the transforms setting. We strip the schema to keep payloads lightweight.
curl -X POST -H "Content-Type: application/json" \
-d '{
"name": "inventory-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "debezium",
"database.password": "secure_pwd",
"database.dbname": "inventory_db",
"database.server.name": "inventory",
"plugin.name": "pgoutput",
"publication.name": "dbz_publication",
"table.include.list": "inventory_db.orders",
"topic.prefix": "inventory",
"transforms": "strip",
"transforms.strip.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.strip.drop.tombstones": "true"
}
}' http://localhost:8083/connectorsThe ExtractNewRecordState transform is critical. Without it, every message contains the full before/after state. That doubles your network bandwidth.
3. Implement the Outbox Pattern in Your Application
Don't write directly to Kafka from your service. Write to a local outbox table in the same transaction.
from fastapi import FastAPI
from sqlalchemy import Column, String, DateTime, create_engine
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
import json
from datetime import datetime
Base = declarative_base()
engine = create_engine("postgresql://user:pass@db/outbox_db")
SessionLocal = sessionmaker(bind=engineProduction Pitfalls & Performance Optimization
Debezium change data capture works great until a hidden bottleneck surfaces. Below are the most common edge cases and how to tame them.
| Pitfall | Symptom | Mitigation |
|---|---|---|
| Memory leak in the connector | JVM heap climbs steadily, GC pauses > 5 s | Set snapshot.mode=initial only once, enable connector.class=io.debezium.connector.mysql.MySqlConnector max.batch.size=2048, and add heartbeat.interval.ms=30000. Restart the connector after a clean snapshot. |
| Back‑pressure on Kafka | Consumer lag spikes, fetch.max.bytes exhausted | Tune consumer.max.poll.records=500 and increase fetch.max.wait.ms=500. Use linger.ms=5 on the producer to let batches coalesce. |
| Outbox table lock contention | INSERTs block each other, latency > 200 ms | Adopt a composite primary key (aggregate_id, event_id) and enable innodb_lock_wait_timeout=5. Batch inserts in groups of 100 using INSERT … VALUES … , …. |
| Rate‑limit throttling by the source DB | Connector logs “Too many requests” | Reduce max.queue.size to 4096, and set snapshot.delay.ms=1000 to spread load. |
| Concurrent schema evolution | Duplicate events, missing columns | Enable database.history.kafka.recovery.poll.interval.ms=5000 and run a single schema change per deployment window. |
Tuning the Debezium Connector
# src/main/resources/debezium-mysql.properties
name=inventory-connector
connector.class=io.debezium.connector.mysql.MySqlConnector
tasks.max=2
database.hostname=db
database.port=3306
database.user=debezium
database.password=debezium_pwd
database.server.id=85744
database.server.name=inventory
database.history.kafka.bootstrap.servers=kafka:9092
database.history.kafka.topic=schema-changes.inventory
snapshot.mode=when_needed
max.batch.size=2048
heartbeat.interval.ms=30000- Keep
tasks.maxmodest; more tasks increase parallelism but also raise MySQL connection count. max.batch.sizecontrols how many rows Debezium pulls per poll. Larger batches reduce round‑trips but raise memory pressure.
Optimizing the Kafka Producer for Exactly‑Once
# producer.properties
bootstrap.servers=kafka:9092
acks=all
enable.idempotence=true
transactional.id=debezium-tx
max.in.flight.requests.per.connection=5
linger.ms=5
batch.size=32768enable.idempotenceguarantees no duplicate writes.transactional.idlets the connector commit offsets atomically with the data payload.max.in.flight.requests.per.connection=5balances throughput with ordering guarantees.
Outbox Batch Writer
public void writeEvents(List<DomainEvent> events) {
String sql = "INSERT INTO outbox (aggregate_id, event_id, payload, created_at) VALUES ";
List<Object> args = new ArrayList<>();
for (DomainEvent e : events) {
sql += "(?, ?, ?, ?),";
args.add(e.getAggregateId());
args.add(UUID.randomUUID().toString());
args.add(e.getPayload());
args.add(Instant.now());
}
sql = sql.replaceAll(",$", ""); // strip trailing comma
jdbcTemplate.batchUpdate(sql, args);
}- Group events in batches of 100–200.
- Use a single
batchUpdatecall to avoid per‑row round‑trips. - The
created_attimestamp lets downstream consumers dedupe safely.
Monitoring the Pipeline
- JVM metrics – Export
jvm.memory.usedandjvm.gc.pauseto Prometheus. - Kafka lag – Track
consumer_lagper partition; alert when > 5000. - Outbox backlog – Count rows where
processed = false. If any metric crosses its threshold, automatically scale the connector task count or spin up an extra consumer group.
Final Summary & Key Takeaways
Debezium change data capture gives you a reliable, low‑latency feed from the source database. Pair it with Kafka’s exactly‑once semantics and an outbox table, and you have a robust event‑driven backbone. The trick is to keep each component tuned:
- Connector configuration – Small batches, heartbeats, and explicit snapshot control prevent memory bloat.
- Kafka producer settings – Idempotence and transactions lock down duplicates.
- Outbox design – Composite keys, batch inserts, and a processed flag keep the write path fast and safe.
- Observability – Export JVM, Kafka, and outbox metrics; set tight alerts. When you respect these boundaries, scaling is as simple as adding more connector tasks or consumer instances. The architecture remains linear, and you retain exactly‑once guarantees without sacrificing throughput.
How do I handle schema changes without breaking the pipeline?
Debezium stores schema evolution in a dedicated Kafka topic (schema-changes.*). Let the connector run in snapshot.mode=when_needed so it captures the new DDL once. Downstream consumers should read the schema topic and apply the latest Avro/Protobuf definition before deserializing events. If you need an immediate cut‑over, pause the connector, apply the DDL, then resume; the connector will emit a DDL event that downstream services can react to.
What’s the safest way to achieve exactly‑once when the outbox table is shared across services?
Wrap the outbox insert and the business transaction in a single DB transaction. After commit, let Debezium capture the row and forward it to Kafka. On the consumer side, enable enable.idempotence and use a transactional consumer (isolation.level=read_committed). The consumer should mark the outbox row as processed = true within the same transaction that publishes the downstream event, guaranteeing atomicity.
When should I consider moving from a single‑node Kafka cluster to a multi‑broker setup?
If you see any of these signals: sustained producer latency > 100 ms, consumer lag consistently above 10 k messages, or disk usage on the broker exceeding 70 %. Multi‑broker clusters also protect against node failure and let you increase partition count for parallelism. Start with three brokers for quorum, then scale partitions per topic to match your consumer parallelism.
If you’re ready to turn these concepts into production‑grade code, let’s talk. Manish Joshi blends deep Flutter UI expertise, AI‑augmented workflows, and battle‑tested FastAPI/Node.js back‑ends. He can help you design a real‑time CDC pipeline that scales, stays consistent, and integrates cleanly with your existing stack.
Building an AI Mobile App or Scalable System?
I engineer production Flutter apps integrated with LLMs, computer vision, LangGraph agents, and high-performance ML backends.