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 đó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.
{
"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.mode | insert chỉ thêm mới · upsert thêm hoặc cập nhật theo khóa · update chỉ cập nhật |
pk.mode | Lấy khóa chính từ đâu: record_key (key của message), record_value (field trong value), none |
pk.fields | Tên (các) field làm khóa chính, vd id |
auto.create | Tự tạo bảng đích nếu chưa có |
auto.evolve | Tự ALTER thêm cột khi schema message có cột mới |
table.name.format | Mẫ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.
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.
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.