Kafka Connect

Sink & JDBC

Sink connector làm chiều ngược lại với Source: nó đọc message từ Kafka topic rồi ghi ra hệ thống ngoài. JDBC Sink connector ghi xuống các database quan hệ như Oracle, MySQL, PostgreSQL — đây là mảnh ghép biến CDC thành một pipeline replication hoàn chỉnh.

Sink connector & JDBC
Từ Kafka ghi xuống database

Sink connector đóng vai trò Consumer: nó subscribe các topic, lấy message và ghi ra đích. JDBC Sink connector (của Confluent, class io.confluent.connect.jdbc.JdbcSinkConnector) ghi xuống bất kỳ database nào có JDBC driver.

Nó có thể tự tạo bảng, tự thêm cột khi schema thay đổi, và hỗ trợ nhiều chế độ ghi (insert thuần, hoặc upsert theo khóa chính). Rất hợp để đồng bộ dữ liệu CDC từ Debezium xuống một DB đích.

Cấu hình JDBC Sink
Ghi xuống Oracle/MySQL
{
  "name": "jdbc-sink-orders",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics": "mysqlsrv.shop.orders",
    "connection.url": "jdbc:oracle:thin:@oracle:1521/XEPDB1",
    "connection.user": "app",
    "connection.password": "app",
    "insert.mode": "upsert",
    "pk.mode": "record_key",
    "pk.fields": "id",
    "auto.create": "true",
    "auto.evolve": "true",
    "table.name.format": "ORDERS",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter"
  }
}

Các option quan trọng:

OptionÝ nghĩa
insert.modeinsert chỉ thêm mới · upsert thêm hoặc cập nhật theo khóa · update chỉ cập nhật
pk.modeLấy khóa chính từ đâu: record_key (key của message), record_value (field trong value), none
pk.fieldsTên (các) field làm khóa chính, vd id
auto.createTự tạo bảng đích nếu chưa có
auto.evolveTự ALTER thêm cột khi schema message có cột mới
table.name.formatMẫu tên bảng đích, vd ${topic} hoặc tên cố định

Muốn upsert thì bắt buộc khai báo pk.mode + pk.fields, vì connector cần biết theo khóa nào để quyết định insert hay update.

SMT – Single Message Transform
Unwrap Debezium & đổi tên topic

Message từ Debezium được gói trong cấu trúc before/after/op. Nếu ghi thẳng xuống DB sẽ ra một bảng có cột "after" lồng nhau — không phải thứ ta muốn. Cần SMT để biến đổi từng message trước khi ghi.

1) ExtractNewRecordState — unwrap, chỉ lấy phần after làm bản ghi phẳng (với delete sẽ phát tombstone hoặc bỏ tùy cấu hình).
2) RegexRouter — đổi tên topic mysqlsrv.shop.orders thành tên bảng orders.

{
  "config": {
    // ... các config JDBC sink ở trên ...
    "transforms": "unwrap,route",

    "transforms.unwrap.type":
        "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "false",
    "transforms.unwrap.delete.handling.mode": "rewrite",

    "transforms.route.type":
        "org.apache.kafka.connect.transforms.RegexRouter",
    "transforms.route.regex": "mysqlsrv\\.shop\\.(.*)",
    "transforms.route.replacement": "$1"
  }
}

transforms liệt kê các SMT chạy theo thứ tự (unwrap trước, route sau). Mỗi SMT có một .type riêng và các tham số của nó. Với regex trên, topic mysqlsrv.shop.orders sẽ được route về bảng orders.

Converter phải khớp giữa source và sink: nếu source ghi bằng JsonConverter thì sink cũng phải đọc bằng JsonConverter; nếu source dùng Avro thì sink phải dùng AvroConverter + cùng Schema Registry. Lệch converter là nguyên nhân lỗi deserialize phổ biến nhất.
← Quay lại
Bài 2: Source & Debezium