MJ
Manish Joshi
ServicesPortfolioFree AI ToolsBlogContact
Start Project →
MJ
Manish Joshi
ServicesPortfolioFree AI ToolsBlogContact
Start Your App →💬 Chat on WhatsApp (+91 95489 50280)
MJ
Manish Joshi

AI-Powered Mobile App Developer. Building production Flutter iOS & Android apps with integrated GenAI, LLMs, computer vision, and scalable ML backends.

Services

  • AI Mobile App Dev
  • Custom Flutter Apps
  • Add AI to Existing Apps
  • AI & ML Infrastructure

Work

  • Case Studies
  • Dliva Delivery
  • SnapQuote AI
  • About & Credentials

Resources

  • Free AI Developer Tools
  • Start Project
  • WhatsApp: +91 95489 50280
  • Privacy Policy

Built with by Manish Joshi

© 2026 manishjoshi.online · All rights reserved

Back to all articles
AI Sep 27, 2026 8 min read

Real‑Time Change Data Capture Pipelines with Debezium, Kafka, and the Outbox Pattern: Architecture, Exactly‑Once Guarantees, and Scaling Strategies

Debezium change data capture streams database WAL events into Kafka, enabling microservices to react in real time. By pairing Debezium with the outbox pattern, you can achieve exactly‑once delivery and simplify transactional consistency. This guide walks through the architecture, scaling techniques, and production‑grade monitoring.
MJ
Manish JoshiAuthor
AI Mobile App Developer & Systems Engineer
AIAI & GENAI PIPELINES

Real‑Time Change Data Capture Pipelines with Debezium, Kafka, and the Outbox Pattern: Architecture, Exactly‑Once Guarantees, and Scaling Strategies

Production InsightsManish Joshi

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:

  1. Latency – sub‑second propagation from OLTP to downstream services.
  2. Consistency – exactly‑once semantics despite retries, crashes, or network partitions.
  3. 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

GuaranteeWhy it mattersHow CDC + Outbox achieves it
Exactly‑once deliveryPrevent duplicate billing or inventory countsKafka’s idempotent producer + Debezium’s transactional source
Order preservation per entityGuarantees correct state reconstructionWAL ordering + partition key = entity ID
Fault‑tolerant replayEnables rapid disaster recoveryKafka retention + compacted outbox topic
Schema evolution without downtimeSupports agile microservice developmentAvro + 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 (c for create, u for update, d for 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?

  1. Transactional write – Service updates its domain tables and inserts a row into outbox_events in the same DB transaction.
  2. Debezium captures – The outbox table appears as another source; its rows become Kafka messages.
  3. 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

  1. Source DB – MySQL 8.x with GTID enabled.
  2. Kafka Cluster – 3‑node KRaft mode, replication factor 3, min‑insync replicas = 2.
  3. Kafka Connect – Debezium MySQL connector, distributed mode, offset storage in Kafka.
  4. Outbox Table – outbox_events(id UUID PK, aggregate_id VARCHAR, type VARCHAR, payload JSONB, created_at TIMESTAMP).
  5. Kafka Streams Application – Enriches events, applies business rules, writes to domain‑specific topics.
  6. Consumer Services – Spring Boot, Quarkus, or Go microservices that read from their topic with exactly‑once processing enabled.
  7. Observability Stack – Prometheus + Grafana dashboards for connector lag, stream processing throughput, and consumer lag.
  8. Schema Registry – Confluent Schema Registry storing Avro schemas for all topics.

Data flow

mermaidUTF-8
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 end

Scaling strategy

  • Connector scaling – Deploy multiple tasks; each task handles a subset of tables or partitions.
  • Kafka Streams parallelism – Set num.stream.threads to match pod CPU limits; use stateless processors for linear scaling.
  • Kubernetes autoscaling – Horizontal Pod Autoscaler (HPA) reacts to kafka.connect.task.max.poll.records and stream.processing.rate metrics from Prometheus.
  • Back‑pressure handling – Enable consumer.max.poll.interval.ms and producer.delivery.timeout.ms to avoid OOM during spikes.

Architecture Patterns Comparison

PatternLatency (ms)Exactly‑once?Operational complexityTypical use‑case
Polling + REST500‑2000No (at‑least‑once)Low – simple cron jobsLow‑volume sync
Dual‑write + Message Queue100‑300No (window of inconsistency)Medium – transaction managementEvent‑driven legacy
CDC + Direct Topic20‑80Yes (Kafka idempotence)High – connector ops, schema mgmtHigh‑throughput microservices
CDC + Outbox + Streams10‑50Yes (transactional + idempotent)Highest – needs DB schema, streams codeMission‑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.

sqlUTF-8
-- 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.

yamlUTF-8
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: - kafka

Register the connector via REST API. Note the transforms setting. We strip the schema to keep payloads lightweight.

bashUTF-8
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/connectors

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

pythonUTF-8
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=engine

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

PitfallSymptomMitigation
Memory leak in the connectorJVM heap climbs steadily, GC pauses > 5 sSet 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 KafkaConsumer lag spikes, fetch.max.bytes exhaustedTune 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 contentionINSERTs block each other, latency > 200 msAdopt 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 DBConnector logs “Too many requests”Reduce max.queue.size to 4096, and set snapshot.delay.ms=1000 to spread load.
Concurrent schema evolutionDuplicate events, missing columnsEnable database.history.kafka.recovery.poll.interval.ms=5000 and run a single schema change per deployment window.

Tuning the Debezium Connector

plainUTF-8
# 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.max modest; more tasks increase parallelism but also raise MySQL connection count.
  • max.batch.size controls how many rows Debezium pulls per poll. Larger batches reduce round‑trips but raise memory pressure.

Optimizing the Kafka Producer for Exactly‑Once

plainUTF-8
# 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=32768
  • enable.idempotence guarantees no duplicate writes.
  • transactional.id lets the connector commit offsets atomically with the data payload.
  • max.in.flight.requests.per.connection=5 balances throughput with ordering guarantees.

Outbox Batch Writer

javaUTF-8
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 batchUpdate call to avoid per‑row round‑trips.
  • The created_at timestamp lets downstream consumers dedupe safely.

Monitoring the Pipeline

  • JVM metrics – Export jvm.memory.used and jvm.gc.pause to Prometheus.
  • Kafka lag – Track consumer_lag per 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:

  1. Connector configuration – Small batches, heartbeats, and explicit snapshot control prevent memory bloat.
  2. Kafka producer settings – Idempotence and transactions lock down duplicates.
  3. Outbox design – Composite keys, batch inserts, and a processed flag keep the write path fast and safe.
  4. 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.

Get in touch with Manish Joshi →

MJ
Written by Manish Joshi

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.

Start Your App Project