Someone subscribes to your website, so the application executes this SQL against the database:

UPDATE users
SET plan = 'pro'
WHERE id = 42;

Now your derived data stores (such as Redis or OpenSearch) may be out of sync, showing free instead of pro.

Change Data Capture ("CDC") is a pattern that you can use to reconcile the differences in near-real time by propagating database changes to downstream systems in a reliable way when eventual consistency is tolerable.

A committed update reaches Redis

Postgres commits plan pro before Redis applies the CDC record. Redis briefly remains on plan free, then matches Postgres.

CDC has a few steps:

  1. Postgres commits an UPDATE, recording it in the write-ahead log.
  2. Debezium reads the logical-decoding stream and writes a change record to a table-specific Kafka topic.
  3. A consumer applies the record to a derived data store.

Debezium reads the WAL in order

Debezium reads committed Postgres WAL records in increasing LSN order. LSN 24023128 updates user 42 from plan free to pro and maps to Kafka offset zero. LSN 24023144 deletes the row for user 7 and maps to offset one. LSN 24023160 updates user 9 from plan free to team and maps to offset two. Consumers can follow this ordered topic without polling Postgres tables.

A key insight is that downstream systems do not need to poll source tables directly; they consume an ordered stream of committed changes and track their position in that stream. In this example, clients can consume from the app.public.users topic while keeping track of their offset in the stream of events.

A simplified Debezium-style update record might look like this:

{
  "topic": "app.public.users",
  "key": {
    "id": 42
  },
  "value": {
    "before": {
      "id": 42,
      "plan": "free"
    },
    "after": {
      "id": 42,
      "plan": "pro"
    },
    "source": {
      "schema": "public",
      "table": "users",
      "lsn": 24023128
    },
    "op": "u"
  }
}

Example: Data Retention

Let's say you need a data retention pipeline: whenever a parent entity is deleted, every service holding related data should purge its own copy. This is useful to enforce retention policies and tenant offboarding. One approach is to choreograph the operations across microservices, letting each service know when a parent entity has been deleted so it can delete the entities it owns. For example, when a customer is deleted, the order service should delete all orders for that customer.

Instead of implementing a transactional outbox in each microservice to reliably publish "Purge Events", you can let the deletion itself be the event. A Flink job listens to the change streams, and for every delete (op: "d") it publishes a purge event to a unified Kafka topic. Each service subscribes to that topic and deletes any related data on its end. Because those deletes produce their own change records, the cascade is self-propagating: deleting a customer emits a purge event, the order service deletes its orders, and those order deletes flow back through CDC as the next round of purge events.

Event-Based Delete Cascade

Customer #91 is deleted first. That deletion creates separate purge events for Order #7012 and Order #7013. Each order is deleted only after its event reaches the Order service, and each order deletion creates a distinct purge event for its matching shipment. Shipment #5012 and Shipment #5013 are then deleted. The complete cascade deletes one customer, two orders, and two shipments.

At LinkedIn

Change Data Capture at LinkedIn is powered by Brooklin, which helps stream database updates with low latency. A major benefit of this system is isolation between applications and online databases. Applications can scale independently from the database, which avoids the risk of bringing down the database.

Incremental ETL

For questions such as...

  • "What were the top 5 products purchased by English-speaking customers this week?"
  • "Which marketing campaigns drove paid upgrades from customers who first signed up on a mobile device?"

... you typically need to join tables that are not located in the same database (e.g., users, orders, product_catalog)

To decouple analytical queries from transaction processing, companies typically maintain a data lake. The data lake helps avoid resource contention with the online databases that are performing work on behalf of a customer.

At LinkedIn, we store offline data in HDFS. We use Trino, Opal, and OpenHouse to query and manage the datasets. The datasets are initialized based on a snapshot of the database (such as mysqldump) and kept in sync through change streams with Brooklin, Kafka, and Gobblin.

Online → Offline Data Flow

Follow one users row from the online commit to its offline HDFS copy.

The MySQL users record with id 42 and the Opal row on HDFS both start with plan free. The application submits an update, MySQL commits plan pro, and Brooklin captures the committed MySQL change and publishes a CDC event to Kafka. Kafka accepts the event and records it in the app.public.users topic. Gobblin reads that topic and writes the record into HDFS. Opal applies plan pro to the offline row. The final state has one online update synchronized to one offline row.

Try it yourself

Run the companion demo repository to exercise the same CDC pipeline locally: wcygan/change-data-capture-demo.