Change Data Capture (CDC) at Scale: Streaming PostgreSQL WAL Changes to Redis & Kafka with Debezium and pgoutput

Application-level dual writes fail silently, while table polling destroys database IOPS. Learn how to stream PostgreSQL write-ahead log (WAL) mutations directly to Kafka and Redis using native pgoutput logical replication and Debezium.

Dual-Write Inconsistencies & The High Cost of Table Polling

In modern decoupled architectures, keeping secondary data stores—such as Redis cache layers, Elasticsearch clusters, and analytics data warehouses—synchronized with your primary PostgreSQL database is an ongoing engineering challenge. Two common architectural patterns are typically deployed to solve this, and both carry fatal production flaws:

  1. Application Dual-Writes: The application code updates PostgreSQL and immediately attempts an update in Redis or Kafka within the same request lifecycle. When network blips, process crashes, or timeout exceptions occur between the primary commit and the secondary write, the stores permanently diverge. The application is left serving stale or corrupted state with zero automated recovery path.
  2. Batch Table Polling: A background worker repeatedly queries SELECT * FROM orders WHERE updated_at > :last_sync every 10 seconds. As table sizes expand into millions of records, polling imposes severe read lock contention, burns database IOPS, pollutes buffer caches, and fails completely to capture record deletions.

The definitive architectural solution is Change Data Capture (CDC). By tailing the database's internal transaction log directly, CDC extracts every insert, update, and delete mutation with zero application overhead and mathematical consistency.

PostgreSQL Logical Decoding Internals: wal_level = logical and pgoutput

PostgreSQL achieves durability by writing all mutations to its Write-Ahead Log (WAL) before persisting changes to disk data pages. Historically, WAL streaming was binary and reserved strictly for physical read replicas. With the introduction of Logical Decoding, PostgreSQL can parse raw WAL segments into structured logical tuples using its native pgoutput replication plugin.

To enable logical decoding, the PostgreSQL server configuration must specify:

# postgresql.conf
wal_level = logical
max_replication_slots = 10
max_wal_senders = 10
track_commit_timestamp = on

A logical publication defines the exact scope of tables streamed to external subscribers:

-- Create logical publication for core domain tables
CREATE PUBLICATION cdc_publication FOR TABLE orders, customer_profiles, inventory_items;

-- Verify replication slot creation
SELECT * FROM pg_create_logical_replication_slot('debezium_cdc_slot', 'pgoutput');

Production Debezium Connector Configuration

Rather than authoring low-level replication parsers from scratch, Debezium serves as the industry-standard CDC connector. It attaches to the PostgreSQL replication slot, acts as a virtual standby server, and streams formatted JSON or Avro events directly into Kafka topics:

{
  "name": "postgres-inventory-cdc",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "tasks.max": "1",
    "plugin.name": "pgoutput",
    "database.hostname": "postgres-primary.internal",
    "database.port": "5432",
    "database.user": "debezium_replicator",
    "database.password": "${file:/secrets/db.properties:cdc_password}",
    "database.dbname": "production_app",
    "database.server.name": "pg_prod",
    "publication.name": "cdc_publication",
    "slot.name": "debezium_cdc_slot",
    "slot.drop.on.stop": "false",
    "decimal.handling.mode": "double",
    "tombstones.on.delete": "true",
    "time.precision.mode": "connect"
  }
}

Taming WAL Bloat: Disk Ceilings and Slot Safeguards

While logical decoding is remarkably fast, it introduces a severe operational risk: WAL accumulation. PostgreSQL will retain all WAL segments on disk indefinitely until every active replication slot confirms it has consumed and acknowledged those LSN (Log Sequence Number) positions. If the Debezium connector or Kafka broker goes offline for several hours, PostgreSQL cannot purge old WAL files. Unmonitored, this behavior will completely fill the database volume, forcing PostgreSQL into an emergency panic shutdown.

To prevent catastrophic disk saturation, configure strict WAL retention ceilings in PostgreSQL 13+:

# Enforce absolute upper bound on WAL retained by replication slots
# If Debezium falls behind by more than 50GB, the slot is automatically invalidated to save the server
max_slot_wal_keep_size = 51200MB

Additionally, configure continuous alerting against pg_replication_slots:

SELECT 
    slot_name,
    active,
    wal_status,
    pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal_bytes
FROM pg_replication_slots;

Consuming CDC Events Idempotently in Python

Downstream consumers reading from CDC topics must handle duplicates caused by network retries. Below is an idempotent Python consumer synchronizing inventory state directly into Redis:

import json
import redis
from confluent_kafka import Consumer

r = redis.Redis(host='redis-cache.internal', port=6379, db=0)
consumer = Consumer({
    'bootstrap.servers': 'kafka-cluster.internal:9092',
    'group.id': 'redis-cache-sync-group',
    'auto.offset.reset': 'earliest',
    'enable.auto.commit': False
})

consumer.subscribe(['pg_prod.public.inventory_items'])

while True:
    msg = consumer.poll(1.0)
    if msg is None or msg.error():
        continue

    payload = json.loads(msg.value().decode('utf-8'))
    op = payload.get('op') # 'c'=create, 'u'=update, 'd'=delete
    before = payload.get('before')
    after = payload.get('after')
    tx_timestamp = payload.get('ts_ms', 0)

    item_id = after['id'] if after else before['id']
    redis_key = f"cache:inventory:{item_id}"

    # Use Redis Lua script to enforce monotonic timestamp ordering
    # Prevents out-of-order event delivery from overwriting newer state
    lua_upsert = (
        "local current_ts = tonumber(redis.call('HGET', KEYS[1], '_ts') or 0); "
        "local incoming_ts = tonumber(ARGV[1]); "
        "if incoming_ts >= current_ts then "
        "  redis.call('HMSET', KEYS[1], 'stock', ARGV[2], '_ts', ARGV[1]); "
        "  return 1; "
        "end; "
        "return 0;"
    )

    if op in ('c', 'u'):
        r.eval(lua_upsert, 1, redis_key, tx_timestamp, after['available_stock'])
    elif op == 'd':
        r.delete(redis_key)

    consumer.commit(msg, asynchronous=True)
Architectural Continuity & Deep Dives

For related production architectures and system implementations, explore these companion guides:

Operating System Tuning: Sustained high-throughput WAL generation from CDC replication pipelines can rapidly fill Linux page caches, leading to catastrophic synchronous write stalls during checkpoints. See how to tune kernel flusher threads in Linux Kernel Dirty Page Writeback & I/O Stalls: Tuning vm.dirty_ratio for Heavy PostgreSQL WAL and Logging Workloads.

Production Takeaway

Change Data Capture fundamentally decouples primary transaction processing from secondary indexing, caching, and stream analytics. By extracting mutations directly from the PostgreSQL WAL using pgoutput and Debezium, you achieve sub-100ms synchronization across caches and search engines with zero dual-write bugs and negligible primary database overhead.

All Insights
Chat on WhatsApp