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

Mục lục bài viết(14 mục)
Hướng dẫn này trình bày chi tiết việc triển khai một pipeline Change Data Capture (CDC) mạnh mẽ cho PostgreSQL, tận dụng Debezium, Kafka Connect và mô hình Transactional Outbox. Mục tiêu là xây dựng một kiến trúc microservice hướng sự kiện, đảm bảo không có sự không nhất quán ghi kép và truyền dữ liệu đáng tin cậy.
Cấu hình PostgreSQL cho Logical Replication
Debezium dựa vào tính năng giải mã logic của PostgreSQL. Điều này đòi hỏi cấu hình máy chủ cụ thể.
postgresql.conf Điều chỉnh
Tham số wal_level phải được đặt thành logical để bật giải mã logic. max_replication_slots và max_wal_senders nên được cấu hình để phù hợp với slot replication của Debezium và các tiến trình WAL sender.
# 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
Sau khi sửa đổi postgresql.conf, việc khởi động lại PostgreSQL là bắt buộc để các thay đổi này có hiệu lực.
Tạo Replication Slot
Debezium yêu cầu một slot replication logic để theo dõi các thay đổi. Slot này đảm bảo rằng PostgreSQL giữ lại các phân đoạn Write-Ahead Log (WAL) cần thiết cho đến khi Debezium xử lý chúng, ngăn ngừa mất dữ liệu. Plugin pgoutput là tiêu chuẩn cho giải mã logic.
-- Connect as a superuser or a user with REPLICATION privileges
SELECT * FROM pg_create_logical_replication_slot('debezium_slot', 'pgoutput');
Xác minh sự tồn tại và trạng thái của slot:
SELECT slot_name, plugin, active, wal_status, restart_lsn FROM pg_replication_slots;
restart_lsn cho biết vị trí WAL mà từ đó slot sẽ bắt đầu truyền các thay đổi. wal_status lý tưởng nhất là reserved hoặc active. Nếu active là f, Debezium hiện không được kết nối.
Thiết lập Debezium Kafka Connect
Debezium là một nền tảng phân tán biến các cơ sở dữ liệu thành các luồng sự kiện. Nó tích hợp với Kafka Connect để truyền các thay đổi từ PostgreSQL vào các Kafka topic.
Triển khai Kafka Connect
Kafka Connect có thể được triển khai ở chế độ độc lập hoặc phân tán. Đối với môi trường sản xuất, chế độ phân tán được ưu tiên để đảm bảo khả năng chịu lỗi và khả năng mở rộng.
# 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
Đảm bảo JAR của connector Debezium PostgreSQL được đặt trong một thư mục được chỉ định bởi plugin.path.
Cấu hình Debezium PostgreSQL Connector
Cấu hình connector Debezium định nghĩa cách nó kết nối với PostgreSQL và dữ liệu nào nó sẽ truyền.
{
"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"
}
}
Triển khai cấu hình này thông qua Kafka Connect REST API:
curl -X POST -H "Content-Type: application/json" --data @debezium-connector.json http://kafka-connect:8083/connectors
Tiến hóa Schema với Confluent Schema Registry và Avro
Debezium, khi được cấu hình với các bộ chuyển đổi Avro, sẽ tự động đăng ký schema với Confluent Schema Registry. Điều này rất quan trọng để quản lý sự tiến hóa schema trong các luồng sự kiện. Khi một schema bảng thay đổi (ví dụ: thêm một cột), Debezium sẽ phát hiện điều này, cập nhật schema Avro trong Schema Registry và các tin nhắn tiếp theo sẽ phản ánh schema mới. Các consumer sau đó có thể sử dụng Schema Registry để giải mã tin nhắn một cách chính xác, ngay cả khi phiên bản schema của chúng khác nhau.
Mô hình Transactional Outbox
Mô hình Transactional Outbox giải quyết vấn đề ghi kép: đảm bảo rằng một giao dịch cơ sở dữ liệu và việc xuất bản một tin nhắn đi là nguyên tử. Thay vì trực tiếp xuất bản lên Kafka, các sự kiện được ghi vào một bảng outbox trong cùng một giao dịch cơ sở dữ liệu với logic nghiệp vụ. Debezium sau đó sẽ lấy các thay đổi của bảng outbox này và xuất bản chúng lên Kafka.
Schema Bảng Outbox
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
);
Triển khai Mô hình
Hãy xem xét một UserService tạo một người dùng mới và cần xuất bản một sự kiện UserCreated.
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, được cấu hình để giám sát bảng outbox, sẽ nắm bắt hoạt động INSERT và xuất bản nó dưới dạng một tin nhắn Kafka.
Khóa Idempotency cho Event Consumer
Khi tiêu thụ các sự kiện từ Kafka, điều quan trọng là phải đảm bảo rằng việc xử lý một sự kiện nhiều lần (do thử lại hoặc cân bằng lại consumer) không dẫn đến trạng thái không nhất quán. Khóa idempotency giải quyết vấn đề này.
Trường id của bảng outbox đóng vai trò là khóa idempotency tự nhiên. Các consumer nên lưu trữ id của các sự kiện đã xử lý và từ chối các bản sao.
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);
So sánh Kiến trúc: CDC so với Polling so với Dual-Write
| Tính năng | CDC (Debezium + Outbox) | Polling | Dual-Write (Kafka trực tiếp) |
|---|---|---|---|
| Tính nguyên tử | Đảm bảo (giao dịch DB đơn lẻ) | N/A (truy vấn DB & gửi tin nhắn riêng biệt) | Không (ghi DB & gửi Kafka là các hoạt động riêng biệt) |
| Độ trễ | Gần thời gian thực | Định kỳ (phút đến giờ) | Gần thời gian thực |
| Sử dụng tài nguyên | Thấp trên DB (dựa trên WAL), vừa phải trên Connect | Cao trên DB (truy vấn lặp lại) | Thấp trên DB, vừa phải trên dịch vụ ứng dụng |
| Rủi ro mất dữ liệu | Tối thiểu (WAL + replication slot) | Cao (bỏ lỡ các thay đổi giữa các lần thăm dò) | Cao (nếu ứng dụng gặp sự cố giữa ghi DB & gửi Kafka) |
| Độ phức tạp | Vừa phải (Debezium, Kafka Connect, Schema Registry) | Thấp (truy vấn đơn giản) | Thấp-Vừa phải (tích hợp client Kafka) |
| Tiến hóa Schema | Tuyệt vời (Avro + Schema Registry) | Thủ công (yêu cầu cập nhật truy vấn/consumer cẩn thận) | Thủ công (yêu cầu cập nhật producer/consumer cẩn thận) |
| Khả năng mở rộng | Cao (Kafka Connect phân tán) | Hạn chế (tải DB tăng theo tần suất thăm dò) | Cao (Kafka mở rộng tốt) |
| Trường hợp sử dụng | Microservice hướng sự kiện, kho dữ liệu, kiểm toán | Tích hợp đơn giản, khối lượng thấp | Không khuyến nghị cho tính nhất quán dữ liệu quan trọng |
Những vấn đề và cách khắc phục trong môi trường sản xuất
-
Phình to Replication Slot:
- Chế độ lỗi:
pg_replication_slotshiển thịwal_statuslàunreservedhoặcretainedtrong thời gian dài,restart_lsnkhông tiến triển, và dung lượng đĩa đầy do các phân đoạn WAL. - Nguyên nhân: Debezium connector bị dừng, tạm dừng hoặc bị kẹt, không thể tiêu thụ từ replication slot. PostgreSQL giữ lại WAL vô thời hạn cho slot không hoạt động.
- Cách khắc phục:
- Khởi động lại Debezium connector.
- Kiểm tra log của Kafka Connect để tìm lỗi (ví dụ: sự cố mạng, Kafka broker không khả dụng, sự cố Schema Registry).
- Nếu connector không thể phục hồi nhanh chóng, hãy cân nhắc xóa replication slot (là biện pháp cuối cùng, điều này sẽ làm mất các thay đổi kể từ khi slot ngừng tiến triển) và tạo lại nó.
SELECT pg_drop_replication_slot('debezium_slot'); - Giám sát
pg_replication_slotsvàpg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)để biết độ trễ của slot.
- Chế độ lỗi:
-
Không tương thích Schema Registry:
- Chế độ lỗi: Debezium connector không khởi động được hoặc tạo ra các tin nhắn không đọc được trong Kafka, với các lỗi như "Schema not found" hoặc "Incompatible schema."
- Nguyên nhân: Schema Registry không khả dụng, hoặc có một thay đổi schema gây lỗi (ví dụ: thay đổi kiểu cột từ
VARCHARsangINTEGERmà không có cài đặt tương thích Avro phù hợp). - Cách khắc phục:
- Xác minh tính khả dụng của Schema Registry và kết nối mạng.
- Kiểm tra log của Schema Registry.
- Đảm bảo
value.converter.schema.registry.urllà chính xác trong cấu hình connector. - Đối với các thay đổi gây lỗi, hãy cân nhắc một topic mới, một connector mới, hoặc một chiến lược để xử lý tiến hóa schema (ví dụ: sử dụng SMT tùy chỉnh để chuyển đổi dữ liệu trước khi tuần tự hóa Avro). Các chế độ tương thích của Confluent Schema Registry (ví dụ:
BACKWARD,FORWARD,FULL) rất quan trọng ở đây.
-
Lỗi Kafka Connect Worker:
- Chế độ lỗi: Các connector ngừng xử lý, các tác vụ được đánh dấu là
FAILED, hoặc toàn bộ cụm Kafka Connect trở nên không phản hồi. - Nguyên nhân: Lỗi hết bộ nhớ, phân vùng mạng, plugin cấu hình sai, hoặc các vấn đề với Kafka/ZooKeeper cơ bản.
- Cách khắc phục:
- Giám sát log của Kafka Connect worker để tìm lỗi OOM hoặc các ngoại lệ khác. Điều chỉnh kích thước heap JVM nếu cần.
- Kiểm tra tình trạng cụm Kafka và ZooKeeper.
- Đảm bảo
offset.storage.topic,config.storage.topic,status.storage.topickhỏe mạnh và có thể truy cập được. - Ở chế độ phân tán, Kafka Connect sẽ tự động cân bằng lại các tác vụ, nhưng các lỗi dai dẳng cho thấy một vấn đề sâu sắc hơn.
- Chế độ lỗi: Các connector ngừng xử lý, các tác vụ được đánh dấu là
-
Độ trễ xử lý bảng Outbox:
- Chế độ lỗi: Các sự kiện được ghi vào bảng
outboxnhưng không xuất hiện trong các Kafka topic kịp thời. - Nguyên nhân: Debezium connector chậm, tạm dừng, hoặc đang gặp áp lực ngược từ Kafka. Bảng
outboxcó thể đang phát triển quá mức. - Cách khắc phục:
- Kiểm tra trạng thái và log của Debezium connector.
- Giám sát tình trạng Kafka broker và độ trễ của consumer group cho topic Debezium.
- Đảm bảo bảng
outboxđược bao gồm trongtable.include.listvàpublication.autocreate.modeđược đặt chính xác. - Cân nhắc thêm một chỉ mục vào
createdattrên bảngoutboxnếu bạn định cắt bỏ các sự kiện cũ (mặc dù Debezium không yêu cầu điều này cho hoạt động của nó).
- Chế độ lỗi: Các sự kiện được ghi vào bảng
Các câu hỏi thường gặp
-
Tại sao lại sử dụng mô hình Transactional Outbox thay vì chỉ để Debezium theo dõi các bảng nghiệp vụ của tôi? Việc theo dõi trực tiếp các bảng nghiệp vụ hoạt động tốt cho các trường hợp đơn giản, nhưng mô hình Transactional Outbox cung cấp một luồng sự kiện chuyên dụng, tách rời. Nó cho phép bạn định nghĩa các hợp đồng sự kiện rõ ràng, làm giàu các sự kiện với siêu dữ liệu bổ sung và đảm bảo rằng tải trọng sự kiện chính xác là những gì bạn muốn xuất bản, thay vì chỉ là một thay đổi hàng cơ sở dữ liệu thô. Nó cũng đơn giản hóa logic consumer bằng cách cung cấp một cấu trúc sự kiện nhất quán.
-
Debezium xử lý các thay đổi schema trong các bảng PostgreSQL như thế nào? Khi một schema thay đổi (ví dụ: thêm một cột, thay đổi kiểu), Debezium sẽ phát hiện điều này. Nếu được cấu hình với các bộ chuyển đổi Avro và Schema Registry, nó sẽ đăng ký một phiên bản mới của schema Avro cho topic đó. Các consumer sau đó sử dụng Schema Registry để lấy phiên bản schema chính xác cho mỗi tin nhắn, cho phép tương thích ngược và tương thích tiến tùy thuộc vào cài đặt tương thích của Schema Registry.
-
Điều gì xảy ra nếu Debezium connector bị dừng? Tôi có bị mất dữ liệu không? Không, bạn sẽ không bị mất dữ liệu. Logical replication slot của PostgreSQL đảm bảo rằng các phân đoạn WAL được giữ lại cho đến khi Debezium xử lý chúng thành công. Khi Debezium khởi động lại, nó sẽ tiếp tục từ LSN (Log Sequence Number) cuối cùng đã commit trong bộ lưu trữ offset của nó, tiếp tục chính xác từ nơi nó đã dừng lại. Tuy nhiên, thời gian ngừng hoạt động kéo dài có thể dẫn đến phình to WAL trên máy chủ PostgreSQL.
-
Làm cách nào để dọn dẹp các sự kiện cũ khỏi bảng
outbox? Debezium chỉ đọc từ bảngoutbox; nó không xóa bản ghi. Bạn cần một tiến trình riêng biệt (ví dụ: một công việc theo lịch trình) để định kỳ xóa các sự kiện cũ khỏi bảngoutbox. Công việc này chỉ nên xóa các sự kiện đã đủ cũ và được xác nhận là đã được Debezium xử lý (ví dụ: bằng cách kiểm tra offset của Debezium hoặc một dấu thời gian). Hãy cẩn thận không xóa các sự kiện mà Debezium chưa truyền. Một chiến lược phổ biến là xóa các sự kiện cũ hơn một ngưỡng nhất định (ví dụ: 7 ngày). -
Tôi có thể sử dụng Debezium với các plugin giải mã logic khác ngoài
pgoutputkhông? Có, Debezium hỗ trợ các plugin khác nhưwal2json. Tuy nhiên,pgoutputlà plugin giải mã logic gốc được giới thiệu trong PostgreSQL 10 và thường được khuyến nghị vì hiệu quả và tích hợp chặt chẽ hơn với giao thức replication của PostgreSQL.wal2jsoncung cấp đầu ra JSON dễ đọc hơn, có thể hữu ích cho việc gỡ lỗi hoặc các trường hợp sử dụng cụ thể, nhưngpgoutputlà tiêu chuẩn cho CDC sản xuất với Debezium.
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

Tối ưu hóa truy vấn PostgreSQL 17: Kế hoạch thực thi, điều chỉnh bộ nhớ & EXPLAIN ANALYZE
Hướng dẫn toàn diện về tối ưu hóa truy vấn PostgreSQL 17: kế hoạch thực thi, điều chỉnh bộ nhớ & explain analyze với kiến trúc cấp độ sản xuất và ví dụ mã.
Read morePostgreSQL Vacuum & Bloat Index: Phát hiện, Giảm thiểu và Tinh chỉnh Tự động
Chẩn đoán và loại bỏ tình trạng phình (bloat) bảng và index trong PostgreSQL. Nắm vững các công thức tinh chỉnh autovacuum, nén dữ liệu không downtime với pg_repack, và cơ chế visibility map của MVCC.
Read more
Khóa phân tán trong hệ thống phân tán: Redlock, PostgreSQL Advisory Locks & etcd Leases
Hướng dẫn toàn diện về khóa phân tán trong hệ thống phân tán: redlock, postgresql advisory locks & etcd leases với kiến trúc cấp độ sản xuất và ví dụ mã.
Read more