Change Data Capture at Scale: Building Real-Time Pipelines Without the Chaos
Batch ETL still matters, but modern products need fresher data. Change data capture (CDC) turns database writes into a stream of events, allowing analytics, search indexes, caches, and microservices to react in near real time. The promise is simple: read the transaction log, publish every insert, update, and delete, and let downstream systems stay current. The reality is harder. Production CDC is an operational discipline that touches replication, schema evolution, ordering, exactly-once semantics, security, and cost. This article is a practical guide to building CDC pipelines that survive real traffic.
What CDC Really Does
CDC captures row-level changes from a source database. The best implementations read the database transaction log instead of polling tables. PostgreSQL exposes logical decoding, MySQL has the binary log, MongoDB offers change streams, SQL Server provides CDC tables, and Oracle uses LogMiner or GoldenGate. Each has different operational constraints, but the goal is the same: emit a durable, ordered stream of changes with enough metadata to reconstruct the current state.
- Log-based CDC: low impact on the source, captures deletes, and preserves commit order per partition. It requires database configuration and careful replication slot or binlog retention management.
- Trigger-based CDC: works on more databases and can capture custom logic, but adds write overhead and schema coupling. Triggers can become a performance and maintenance burden.
- Query-based CDC: polls tables using timestamps or version columns. It is simple, but it misses hard deletes, can race with transactions, and becomes expensive at scale.
Choose log-based CDC when the database supports it and the use case demands low latency. Use query-based CDC only for low-volume, tolerant workloads.
Reference Architecture
A reliable CDC system separates capture, transport, processing, and serving. The source database produces changes. A connector captures them and writes to a durable log such as Apache Kafka, Apache Pulsar, or Amazon Kinesis. Stream processors transform, enrich, and route events. Sinks materialize data into warehouses, search engines, caches, or service-owned read models. A schema registry and control plane sit across the pipeline to enforce contracts and manage connectors.
- Ingest: Debezium, AWS Database Migration Service, Google Datastream, Fivetran, or a custom log reader.
- Transport: Kafka topics with partitioning by primary key. Use compacted topics for latest-state snapshots and standard topics for event history.
- Process: Kafka Streams, Apache Flink, Spark Structured Streaming, or a stream database such as Materialize.
- Serve: Snowflake, BigQuery, Redshift, Elasticsearch, Redis, PostgreSQL read models, or feature stores.
Keep the pipeline idempotent at every hop. Assume events can be duplicated, delayed, or replayed. Design sinks to handle upserts and deletes by primary key.
Designing the Event Envelope
Every CDC event needs a consistent envelope. Include the operation type (insert, update, delete), before and after row images, source metadata, transaction identifier, log sequence number, timestamp, and schema version. The key should be the table and primary key so that all changes for one entity land in the same partition. This preserves per-entity order and enables stateful processing.
A practical envelope contains fields such as: op=u, before={...}, after={...}, source.table=orders, source.lsn=123456, ts_ms=1710000000000, and schema_version=2. For deletes, emit a tombstone or a delete event with the key and no after image. Tombstones are essential for compacted topics because they remove stale state.
Partitioning matters. If you partition by table, a large table can create a hot partition. If you partition by entity key, you get better parallelism and per-entity order. For multi-tenant systems, include tenant id in the key but avoid unbounded key cardinality that breaks compaction. Monitor partition skew and rebalance when necessary.
Schema Evolution and Data Contracts
Schema drift is the most common cause of CDC pipeline failures. A source team adds a column, changes a type, or drops a field, and downstream jobs break. A schema registry with Avro, Protobuf, or JSON Schema helps, but it is not enough. You need explicit data contracts between producers and consumers.
- Add new fields as optional with sensible defaults.
- Never rename or change the type of an existing field without a major version and a migration plan.
- Treat enum changes as breaking unless consumers can ignore unknown values.
- Version the schema and include the version in every event.
- Run compatibility checks in CI before deploying connector or processor changes.
Data contracts should define ownership, freshness expectations, allowed values, and retention. Treat a contract violation like a production incident. Use a dead-letter queue for malformed events, but alert on DLQ growth because it usually indicates a hidden schema or logic problem.
Exactly-Once, Idempotency, and Ordering
Exactly-once end-to-end is difficult and often unnecessary. At-least-once delivery with idempotent sinks is the pragmatic default. Deduplicate using a stable event id such as source system, table, primary key, and log sequence number. Store the deduplication key in a fast store with a TTL, or rely on the sink to perform an idempotent upsert.
- Kafka transactions: useful for read-process-write pipelines that consume from and produce to Kafka.
- Flink checkpoints: provide exactly-once state within the stream processor, but sinks must still be idempotent.
- Warehouse MERGE: upsert by primary key, but be careful with late-arriving events and delete markers.
- Search indexes: use the primary key as the document id and index with external versioning.
Ordering is per partition only. Do not assume global order across a table or database. If a business process requires order for one entity, key by that entity. For cross-entity consistency, use event-time processing with watermarks, or design commands that tolerate reordering. Out-of-order events are normal in distributed systems, especially during snapshots or connector restarts.
The Outbox Pattern and Transactional Boundaries
Microservices often need to update a database and publish an event. Writing to the database and then publishing to Kafka is a dual write: one can succeed while the other fails. CDC can help, but only if the event is part of the same database transaction. The transactional outbox pattern solves this. The service writes the business row and an outbox row in one transaction. CDC then streams the outbox table to the message log. Downstream consumers process the event and acknowledge it.
This pattern keeps the database as the source of truth and avoids distributed transactions. It also makes event publishing replayable. The trade-off is an extra table and the need to clean up delivered outbox rows. Some teams use CDC directly on domain tables, but that couples consumers to the internal schema. Use the outbox pattern when you need a stable public event contract.
Backpressure, Throughput, and Cost
CDC can overwhelm a pipeline during initial snapshots, bulk backfills, or traffic spikes. Plan for backpressure instead of hoping it never happens. Use batching and compression, scale partitions and consumers, and tune connector fetch sizes. For initial snapshots, use incremental snapshotting with chunking and rate limits. Read from a replica when possible, but remember that replicas can lag and have their own log retention limits.
- Monitor source log retention and replication slot lag. A stalled connector can cause unbounded disk growth on the source.
- Set topic retention based on replay needs, not habit. Longer retention costs more but enables recovery.
- Use tiered storage for Kafka to reduce cost while keeping replayability.
- Compact topics for latest-state views; use standard retention for full history.
- Autoscale stream processors based on consumer lag, not just CPU.
Security and Privacy
CDC streams often contain sensitive data. Encrypt data in transit and at rest. Apply least-privilege access to the source database, Kafka topics, schema registry, and sinks. Use field-level encryption or tokenization for PII before it reaches broad-access topics. Mask columns in the connector using single message transforms or in the stream processor. Never put PII in topic names, consumer group names, or tags.
Privacy regulations add requirements. A delete request must propagate to every derived store. Tombstones and delete events are part of the answer, but you also need to verify that backups, caches, search indexes, and analytics tables honor the deletion. Maintain a data map that shows where CDC data flows. Audit access to CDC topics because they are often a high-value target.
Operations and Observability
You cannot operate CDC with database dashboards alone. Build end-to-end observability around the event lifecycle. Track connector status, source log position, replication slot lag, consumer lag, end-to-end latency, error rate, schema errors, DLQ size, and sink freshness. Alert on lag trends, not just static thresholds.
- Use distributed tracing with the event id to follow a change from database to sink.
- Run data quality checks for null rates, referential integrity, row counts, and freshness.
- Create runbooks for snapshot restart, schema change, poison pill, connector failover, and source database failover.
- Chaos test by killing connectors, brokers, and sinks. Verify recovery and replay.
One critical metric is replication slot lag. In PostgreSQL, an inactive slot retains WAL and can fill the disk. In MySQL, binlog retention may be too short for connector recovery. In MongoDB, the oplog window may expire. Know the retention limits of your source and alert before you hit them.
Implementation Blueprint with Debezium and Kafka
A common stack is Debezium for capture, Kafka for transport, and Flink or Kafka Streams for processing. On PostgreSQL, enable logical replication, create a publication, and let Debezium create a replication slot. Configure the connector with table include lists, topic routing, and single message transforms for masking. Register schemas in a registry. Produce one topic per table or one topic per aggregate, depending on consumer needs.
Then build a stream processor that validates schemas, deduplicates, and upserts into the target. Send malformed events to a DLQ with enough context to replay. Monitor the replication slot and Kafka consumer lag. Test snapshot and incremental phases separately. Start with one low-risk table, prove the operational model, then expand.
For MySQL, ensure binlog format is row-based and retention is longer than your recovery window. For MongoDB, size the oplog for peak write volume and monitor its window. For SQL Server, understand the CDC cleanup job and its impact on retention.
Common Pitfalls
- Ignoring replication slot lag and causing source disk exhaustion.
- Skipping schema compatibility checks until a deploy breaks consumers.
- Assuming global event ordering.
- Running an initial snapshot without throttling and saturating the source.
- Forgetting deletes and tombstones in downstream sinks.
- Leaking PII into broad-access topics.
- Having no DLQ, replay strategy, or dead-letter runbook.
- Treating CDC as zero-impact on the source database.
- Coupling consumers directly to internal source schemas without a contract.
When Not to Use CDC
CDC is not free. It adds connectors, topics, schemas, processing jobs, monitoring, and on-call responsibility. If daily batch analytics is enough, keep the simpler pipeline. If the source database cannot enable logical replication, or if write volume is so high that log retention is a constant risk, consider alternatives such as application-level events, periodic exports, or a purpose-built operational data store. Use CDC when low latency, deletes, and event-driven integration justify the complexity.
Conclusion
Change data capture at scale is less about a single tool and more about disciplined system design. Define a stable event envelope. Enforce schema contracts. Make every sink idempotent. Preserve per-entity order. Watch source log retention like a production SLO. Secure sensitive data end to end. Start small, measure end-to-end latency, and expand only after the operational model is proven. Done well, CDC turns database changes into a reliable real-time nervous system for your data platform.

