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

Table of Contents(14 sections)
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.
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.
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
| Feature | CDC (Debezium + Outbox) | Polling | Dual-Write (Direct Kafka) |
|---|---|---|---|
| Atomicity | Guaranteed (single DB transaction) | N/A (separate DB query & message send) | No (DB write & Kafka send are separate operations) |
| Latency | Near real-time | Periodic (minutes to hours) | Near real-time |
| Resource Usage | Low on DB (WAL-based), moderate on Connect | High on DB (repeated queries) | Low on DB, moderate on app service |
| Data Loss Risk | Minimal (WAL + replication slot) | High (missed changes between polls) | High (if app crashes between DB write & Kafka send) |
| Complexity | Moderate (Debezium, Kafka Connect, Schema Registry) | Low (simple queries) | Low-Moderate (Kafka client integration) |
| Schema Evolution | Excellent (Avro + Schema Registry) | Manual (requires careful query/consumer updates) | Manual (requires careful producer/consumer updates) |
| Scalability | High (Kafka Connect distributed) | Limited (DB load increases with polling frequency) | High (Kafka scales well) |
| Use Case | Event-driven microservices, data warehousing, auditing | Simple, low-volume integrations | Not recommended for critical data consistency |
Production Gotchas & Troubleshooting
-
Replication Slot Bloat:
- Failure Mode:
pg_replication_slotsshowswal_statusasunreservedorretainedfor extended periods,restart_lsnnot 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_slotsandpg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)for slot lag.
- Failure Mode:
-
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
VARCHARtoINTEGERwithout proper Avro compatibility settings). - Fix:
- Verify Schema Registry availability and network connectivity.
- Check Schema Registry logs.
- Ensure
value.converter.schema.registry.urlis 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.
-
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.topicare healthy and accessible. - In distributed mode, Kafka Connect should automatically rebalance tasks, but persistent failures indicate a deeper issue.
- Failure Mode: Connectors stop processing, tasks are marked as
-
Outbox Table Processing Lag:
- Failure Mode: Events are written to the
outboxtable but don't appear in Kafka topics promptly. - Cause: Debezium connector is slow, paused, or experiencing backpressure from Kafka. The
outboxtable 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
outboxtable is included intable.include.listandpublication.autocreate.modeis set correctly. - Consider adding an index to
createdaton theoutboxtable if you plan to prune old events (though Debezium doesn't require this for its operation).
- Failure Mode: Events are written to the
Frequently Asked Questions
-
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.
-
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.
-
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.
-
How do I clean up old events from the
outboxtable? Debezium only reads from theoutboxtable; it does not delete records. You need a separate process (e.g., a scheduled job) to periodically delete old events from theoutboxtable. 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). -
Can I use Debezium with other logical decoding plugins besides
pgoutput? Yes, Debezium supports other plugins likewal2json. However,pgoutputis the native logical decoding plugin introduced in PostgreSQL 10 and is generally recommended for its efficiency and tighter integration with PostgreSQL's replication protocol.wal2jsonprovides a more human-readable JSON output, which might be useful for debugging or specific use cases, butpgoutputis the standard for production CDC with Debezium.
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

PostgreSQL 17 Query Optimization: Execution Plans, Memory Tuning & EXPLAIN ANALYZE
Comprehensive guide covering postgresql 17 query optimization: execution plans, memory tuning & explain analyze with production-grade architecture and code examples.
Read more
Distributed Locking in Distributed Systems: Redlock, PostgreSQL Advisory Locks & etcd Leases
Comprehensive guide covering distributed locking in distributed systems: redlock, postgresql advisory locks & etcd leases with production-grade architecture and code examples.
Read more
PgBouncer Architecture & Tuning: Transaction Pooling, Prepared Statements & Session Overhead
Comprehensive guide covering pgbouncer architecture & tuning: transaction pooling, prepared statements & session overhead with production-grade architecture and code examples.
Read more