•12 min read

PostgreSQL Change Data Capture (CDC): Debezium, Kafka Connect & Transactional Outbox

PostgreSQL Change Data Capture (CDC): Debezium, Kafka Connect & Transactional Outbox

This guide details the implementation of a robust Change Data Capture (CDC) pipeline for PostgreSQL, leveraging Debezium, Kafka Connect, and the Transactional Outbox pattern. The objective is to construct an event-driven microservice architecture ensuring zero dual-write inconsistencies and reliable data propagation.

Audio Briefing
0:00 / 0:00

PostgreSQL Configuration for Logical Replication

Debezium relies on PostgreSQL's logical decoding feature. This necessitates specific server configuration.

postgresql.conf Adjustments

The wal_level parameter must be set to logical to enable logical decoding. max_replication_slots and max_wal_senders should be configured to accommodate Debezium's replication slot and WAL sender processes.

# postgresql.conf
wal_level = logical
max_replication_slots = 10 # Adjust based on number of Debezium connectors
max_wal_senders = 10       # Adjust based on number of Debezium connectors

After modifying postgresql.conf, a PostgreSQL restart is mandatory for these changes to take effect.

Creating a Replication Slot

Debezium requires a logical replication slot to track changes. This slot ensures that PostgreSQL retains the necessary Write-Ahead Log (WAL) segments until Debezium has processed them, preventing data loss. The pgoutput plugin is the standard for logical decoding.

-- Connect as a superuser or a user with REPLICATION privileges
SELECT * FROM pg_create_logical_replication_slot('debezium_slot', 'pgoutput');

Verify the slot's existence and state:

SELECT slot_name, plugin, active, wal_status, restart_lsn FROM pg_replication_slots;

The restart_lsn indicates the WAL location from which the slot will start streaming changes. wal_status should ideally be reserved or active. If active is f, Debezium is not currently connected.

Advertisement

Debezium Kafka Connect Setup

Debezium is a distributed platform that turns databases into event streams. It integrates with Kafka Connect to stream changes from PostgreSQL into Kafka topics.

Kafka Connect Deployment

Kafka Connect can be deployed in standalone or distributed mode. For production, distributed mode is preferred for fault tolerance and scalability.

# Example: Distributed Kafka Connect worker configuration (connect-distributed.properties)
bootstrap.servers=kafka-broker-1:9092,kafka-broker-2:9092
group.id=connect-cluster
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-statuses
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://schema-registry:8081
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://schema-registry:8081
internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false
plugin.path=/usr/share/java,/usr/share/confluent-hub-components

Ensure the Debezium PostgreSQL connector JAR is placed in a directory specified by plugin.path.

Debezium PostgreSQL Connector Configuration

The Debezium connector configuration defines how it connects to PostgreSQL and what data it streams.

{
  "name": "postgres-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "tasks.max": "1",
    "database.hostname": "postgres-db",
    "database.port": "5432",
    "database.user": "debezium_user",
    "database.password": "debezium_password",
    "database.dbname": "your_database",
    "database.server.name": "your_logical_server_name",
    "schema.include.list": "public",
    "table.include.list": "public.users,public.orders",
    "slot.name": "debezium_slot",
    "publication.name": "debezium_publication",
    "publication.autocreate.mode": "all_tables",
    "plugin.name": "pgoutput",
    "topic.prefix": "dbserver",
    "heartbeat.interval.ms": "5000",
    "snapshot.mode": "initial",
    "decimal.handling.mode": "double",
    "time.precision.mode": "connect",
    "hstore.handling.mode": "json",
    "converters": "dateConverter",
    "dateConverter.type": "org.apache.kafka.connect.data.Timestamp",
    "dateConverter.format": "yyyy-MM-dd HH:mm:ss.SSS",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "key.converter.schema.registry.url": "http://schema-registry:8081",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "http://schema-registry:8081"
  }
}

Deploy this configuration via the Kafka Connect REST API:

curl -X POST -H "Content-Type: application/json" --data @debezium-connector.json http://kafka-connect:8083/connectors

Schema Evolution with Confluent Schema Registry and Avro

Debezium, when configured with Avro converters, automatically registers schemas with Confluent Schema Registry. This is crucial for managing schema evolution in event streams. When a table schema changes (e.g., adding a column), Debezium detects this, updates the Avro schema in the Schema Registry, and subsequent messages reflect the new schema. Consumers can then use the Schema Registry to deserialize messages correctly, even if their schema versions differ.

Transactional Outbox Pattern

The Transactional Outbox pattern addresses the dual-write problem: ensuring that a database transaction and an outgoing message publication are atomic. Instead of directly publishing to Kafka, events are written to an outbox table within the same database transaction as the business logic. Debezium then picks up these outbox table changes and publishes them to Kafka.

Outbox Table Schema

CREATE TABLE outbox (
    id UUID PRIMARY KEY,
    aggregatetype VARCHAR(255) NOT NULL,
    aggregateid UUID NOT NULL,
    type VARCHAR(255) NOT NULL,
    payload JSONB NOT NULL,
    createdat TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

Implementing the Pattern

Consider a UserService that creates a new user and needs to publish a UserCreated event.

import { Pool } from 'pg';
import { v4 as uuidv4 } from 'uuid';

interface User {
  id: string;
  name: string;
  email: string;
}

interface OutboxEvent {
  id: string;
  aggregateType: string;
  aggregateId: string;
  type: string;
  payload: any;
}

class UserService {
  private pool: Pool;

  constructor(pool: Pool) {
    this.pool = pool;
  }

  public async createUser(name: string, email: string): Promise<User> {
    const client = await this.pool.connect();
    try {
      await client.query('BEGIN');

      const userId = uuidv4();
      const user: User = { id: userId, name, email };

      // 1. Persist business entity
      await client.query(
        'INSERT INTO users (id, name, email) VALUES ($1, $2, $3)',
        [user.id, user.name, user.email]
      );

      // 2. Create outbox event
      const event: OutboxEvent = {
        id: uuidv4(),
        aggregateType: 'User',
        aggregateId: user.id,
        type: 'UserCreated',
        payload: { userId: user.id, name: user.name, email: user.email, timestamp: new Date().toISOString() },
      };

      // 3. Persist outbox event in the same transaction
      await client.query(
        'INSERT INTO outbox (id, aggregatetype, aggregateid, type, payload) VALUES ($1, $2, $3, $4, $5)',
        [event.id, event.aggregateType, event.aggregateId, event.type, JSON.stringify(event.payload)]
      );

      await client.query('COMMIT');
      return user;
    } catch (error) {
      await client.query('ROLLBACK');
      console.error('Failed to create user and publish event:', error);
      throw error;
    } finally {
      client.release();
    }
  }
}

// Example usage (assuming 'pool' is a configured pg.Pool instance)
// const userService = new UserService(pool);
// userService.createUser('John Doe', 'john.doe@example.com')
//   .then(user => console.log('User created:', user))
//   .catch(err => console.error('Error:', err));

Debezium, configured to monitor the outbox table, will capture the INSERT operation and publish it as a Kafka message.

Idempotency Keys for Event Consumers

When consuming events from Kafka, it's critical to ensure that processing an event multiple times (due to retries or consumer rebalancing) does not lead to inconsistent state. Idempotency keys solve this.

The id field of the outbox table serves as a natural idempotency key. Consumers should store the id of processed events and reject duplicates.

import { Kafka, Consumer } from 'kafkajs';
import { Pool } from 'pg';

interface UserCreatedEvent {
  userId: string;
  name: string;
  email: string;
  timestamp: string;
}

class UserEventHandler {
  private consumer: Consumer;
  private pool: Pool;

  constructor(kafka: Kafka, pool: Pool) {
    this.consumer = kafka.consumer({ groupId: 'user-service-consumer-group' });
    this.pool = pool;
  }

  public async start(): Promise<void> {
    await this.consumer.connect();
    await this.consumer.subscribe({ topic: 'dbserver.public.outbox', fromBeginning: false });

    await this.consumer.run({
      eachMessage: async ({ topic, partition, message }) => {
        if (!message.value) return;

        const event = JSON.parse(message.value.toString());
        const outboxEventId = event.payload.id; // The ID from the outbox table, our idempotency key
        const operation = event.op; // 'c' for create, 'u' for update, 'd' for delete

        if (operation !== 'c') {
          // We only care about new outbox events, not updates/deletes to the outbox table itself
          return;
        }

        const payload = event.after.payload; // Debezium 'after' field contains the new row
        const eventType = event.after.type;

        if (eventType === 'UserCreated') {
          const userCreatedEvent: UserCreatedEvent = payload;
          await this.processUserCreatedEvent(outboxEventId, userCreatedEvent);
        }
        // Handle other event types
      },
    });
  }

  private async processUserCreatedEvent(outboxEventId: string, event: UserCreatedEvent): Promise<void> {
    const client = await this.pool.connect();
    try {
      await client.query('BEGIN');

      // Check for idempotency: has this event ID been processed before?
      const checkResult = await client.query(
        'SELECT 1 FROM processed_events WHERE event_id = $1',
        [outboxEventId]
      );

      if (checkResult.rows.length > 0) {
        console.log(`Event ${outboxEventId} already processed. Skipping.`);
        await client.query('ROLLBACK'); // Rollback the transaction as nothing new was done
        return;
      }

      // Process the event (e.g., create a user in a read model, send a welcome email)
      console.log(`Processing UserCreated event for user ${event.userId}:`, event);
      // Example: Insert into a read-model table
      await client.query(
        'INSERT INTO read_model_users (id, name, email) VALUES ($1, $2, $3) ON CONFLICT (id) DO NOTHING',
        [event.userId, event.name, event.email]
      );

      // Record the event ID as processed
      await client.query(
        'INSERT INTO processed_events (event_id, processed_at) VALUES ($1, NOW())',
        [outboxEventId]
      );

      await client.query('COMMIT');
      console.log(`Successfully processed event ${outboxEventId}`);
    } catch (error) {
      await client.query('ROLLBACK');
      console.error(`Error processing event ${outboxEventId}:`, error);
      throw error;
    } finally {
      client.release();
    }
  }
}

// Ensure 'processed_events' table exists for idempotency tracking
// CREATE TABLE processed_events (
//     event_id UUID PRIMARY KEY,
//     processed_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
// );

// Example usage
// const kafka = new Kafka({ brokers: ['kafka-broker-1:9092'] });
// const userEventHandler = new UserEventHandler(kafka, pool);
// userEventHandler.start().catch(console.error);

Architecture Comparison: CDC vs. Polling vs. Dual-Write

FeatureCDC (Debezium + Outbox)PollingDual-Write (Direct Kafka)
AtomicityGuaranteed (single DB transaction)N/A (separate DB query & message send)No (DB write & Kafka send are separate operations)
LatencyNear real-timePeriodic (minutes to hours)Near real-time
Resource UsageLow on DB (WAL-based), moderate on ConnectHigh on DB (repeated queries)Low on DB, moderate on app service
Data Loss RiskMinimal (WAL + replication slot)High (missed changes between polls)High (if app crashes between DB write & Kafka send)
ComplexityModerate (Debezium, Kafka Connect, Schema Registry)Low (simple queries)Low-Moderate (Kafka client integration)
Schema EvolutionExcellent (Avro + Schema Registry)Manual (requires careful query/consumer updates)Manual (requires careful producer/consumer updates)
ScalabilityHigh (Kafka Connect distributed)Limited (DB load increases with polling frequency)High (Kafka scales well)
Use CaseEvent-driven microservices, data warehousing, auditingSimple, low-volume integrationsNot recommended for critical data consistency
Advertisement

Production Gotchas & Troubleshooting

  1. Replication Slot Bloat:

    • Failure Mode: pg_replication_slots shows wal_status as unreserved or retained for extended periods, restart_lsn not advancing, and disk space filling up due to WAL segments.
    • Cause: Debezium connector is down, paused, or stuck, failing to consume from the replication slot. PostgreSQL retains WALs indefinitely for the inactive slot.
    • Fix:
      • Restart the Debezium connector.
      • Check Kafka Connect logs for errors (e.g., network issues, Kafka broker unavailability, Schema Registry issues).
      • If the connector cannot be recovered quickly, consider dropping the replication slot (as a last resort, this will lose changes since the slot stopped advancing) and recreating it. SELECT pg_drop_replication_slot('debezium_slot');
      • Monitor pg_replication_slots and pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) for slot lag.
  2. Schema Registry Incompatibility:

    • Failure Mode: Debezium connector fails to start or produces unreadable messages in Kafka, with errors like "Schema not found" or "Incompatible schema."
    • Cause: Schema Registry is unavailable, or there's a breaking schema change (e.g., changing a column type from VARCHAR to INTEGER without proper Avro compatibility settings).
    • Fix:
      • Verify Schema Registry availability and network connectivity.
      • Check Schema Registry logs.
      • Ensure value.converter.schema.registry.url is correct in the connector config.
      • For breaking changes, consider a new topic, a new connector, or a strategy to handle schema evolution (e.g., using a custom SMT to transform data before Avro serialization). Confluent Schema Registry's compatibility modes (e.g., BACKWARD, FORWARD, FULL) are crucial here.
  3. Kafka Connect Worker Failures:

    • Failure Mode: Connectors stop processing, tasks are marked as FAILED, or the entire Kafka Connect cluster becomes unresponsive.
    • Cause: Out-of-memory errors, network partitions, misconfigured plugins, or issues with underlying Kafka/ZooKeeper.
    • Fix:
      • Monitor Kafka Connect worker logs for OOM errors or other exceptions. Adjust JVM heap size if necessary.
      • Check Kafka and ZooKeeper cluster health.
      • Ensure offset.storage.topic, config.storage.topic, status.storage.topic are healthy and accessible.
      • In distributed mode, Kafka Connect should automatically rebalance tasks, but persistent failures indicate a deeper issue.
  4. Outbox Table Processing Lag:

    • Failure Mode: Events are written to the outbox table but don't appear in Kafka topics promptly.
    • Cause: Debezium connector is slow, paused, or experiencing backpressure from Kafka. The outbox table might be growing excessively.
    • Fix:
      • Check Debezium connector status and logs.
      • Monitor Kafka broker health and consumer group lag for the Debezium topic.
      • Ensure the outbox table is included in table.include.list and publication.autocreate.mode is set correctly.
      • Consider adding an index to createdat on the outbox table if you plan to prune old events (though Debezium doesn't require this for its operation).

Frequently Asked Questions

  1. Why use the Transactional Outbox pattern instead of just letting Debezium watch my business tables? Watching business tables directly works for simple cases, but the Transactional Outbox pattern provides a dedicated, decoupled event stream. It allows you to define explicit event contracts, enrich events with additional metadata, and ensures that the event payload is exactly what you intend to publish, rather than just a raw database row change. It also simplifies consumer logic by providing a consistent event structure.

  2. How does Debezium handle schema changes in PostgreSQL tables? When a schema changes (e.g., adding a column, changing a type), Debezium detects this. If configured with Avro converters and Schema Registry, it registers a new version of the Avro schema for that topic. Consumers then use the Schema Registry to fetch the correct schema version for each message, allowing for backward and forward compatibility depending on the Schema Registry's compatibility settings.

  3. What happens if the Debezium connector goes down? Will I lose data? No, you will not lose data. PostgreSQL's logical replication slot ensures that WAL segments are retained until Debezium has successfully processed them. When Debezium restarts, it will resume from the last committed LSN (Log Sequence Number) in its offset storage, picking up exactly where it left off. However, prolonged downtime can lead to WAL bloat on the PostgreSQL server.

  4. How do I clean up old events from the outbox table? Debezium only reads from the outbox table; it does not delete records. You need a separate process (e.g., a scheduled job) to periodically delete old events from the outbox table. This job should only delete events that are sufficiently old and confirmed to have been processed by Debezium (e.g., by checking Debezium's offset or a timestamp). Be cautious not to delete events that Debezium has not yet streamed. A common strategy is to delete events older than a certain threshold (e.g., 7 days).

  5. Can I use Debezium with other logical decoding plugins besides pgoutput? Yes, Debezium supports other plugins like wal2json. However, pgoutput is the native logical decoding plugin introduced in PostgreSQL 10 and is generally recommended for its efficiency and tighter integration with PostgreSQL's replication protocol. wal2json provides a more human-readable JSON output, which might be useful for debugging or specific use cases, but pgoutput is the standard for production CDC with Debezium.

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement