•15 min read

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

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

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.

Audio Briefing
0:00 / 0:00

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.

Advertisement

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ăngCDC (Debezium + Outbox)PollingDual-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ênThấp trên DB (dựa trên WAL), vừa phải trên ConnectCao 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ệuTố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ạpVừ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 SchemaTuyệ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ộngCao (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ụngMicroservice hướng sự kiện, kho dữ liệu, kiểm toánTích hợp đơn giản, khối lượng thấpKhông khuyến nghị cho tính nhất quán dữ liệu quan trọng
Advertisement

Những vấn đề và cách khắc phục trong môi trường sản xuất

  1. Phình to Replication Slot:

    • Chế độ lỗi: pg_replication_slots hiển thị wal_status là unreserved hoặc retained trong thời gian dài, restart_lsn khô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_slots và pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) để biết độ trễ của slot.
  2. 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ừ VARCHAR sang INTEGER mà 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.url là 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.
  3. 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.topic khỏ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.
  4. Độ trễ xử lý bảng Outbox:

    • Chế độ lỗi: Các sự kiện được ghi vào bảng outbox như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 outbox có 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 trong table.include.list và publication.autocreate.mode được đặt chính xác.
      • Cân nhắc thêm một chỉ mục vào createdat trên bảng outbox nế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ó).

Các câu hỏi thường gặp

  1. 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.

  2. 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.

  3. Đ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.

  4. 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ảng outbox; 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ảng outbox. 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).

  5. Tôi có thể sử dụng Debezium với các plugin giải mã logic khác ngoài pgoutput không? Có, Debezium hỗ trợ các plugin khác như wal2json. Tuy nhiên, pgoutput là 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. wal2json cung 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ưng pgoutput là tiêu chuẩn cho CDC sản xuất với 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