Kafka Connect đưa dữ liệu vào và ra Kafka bằng JSON cấu hình thay vì mã. Bài này dựng một đường ống PostgreSQL → Kafka, đo nó, và tìm ra một cách hỏng im lặng hoàn toàn.

Cấu hình và số đo, cái bẫy của mode incrementing, và bốn chế độ

Dựng bằng một lời gọi HTTP

curl -X POST http://cn:8083/connectors -H "Content-Type: application/json" -d '{
 "name":"pg-source",
 "config":{
  "connector.class":"io.confluent.connect.jdbc.JdbcSourceConnector",
  "connection.url":"jdbc:postgresql://pg:5432/shop",
  "connection.user":"postgres","connection.password":"pw",
  "table.whitelist":"don_hang",
  "mode":"incrementing","incrementing.column.name":"id",
  "topic.prefix":"pg-","poll.interval.ms":"1000"}}'

Không có mã nào. Connect đọc bảng, biến mỗi dòng thành một bản ghi Kafka, và nhớ vị trí đã đọc tới.

trạng thái: RUNNING | task: ['RUNNING']
{"id":1,"ma":"ORD-1","khach":"CUST-1","tien":137,"trang_thai":"PENDING",...}

Số đo

thời gian Connect khởi động 30,1 giây
bộ nhớ Connect 875 MB → 1,1 GB khi chạy
nạp 100.000 dòng mới 1,0 giây = 104.394 dòng/giây
độ trễ một dòng 899 – 1.920 ms

Hai điều đáng chú ý.

Connect nặng hơn broker. Một broker Kafka tốn 339 MB (phần 21); Connect tốn hơn một gigabyte. Nó là một ứng dụng Java đầy đủ với hệ thống nạp plugin, REST API và cơ chế phân tán riêng.

Thông lượng thì rất tốt — hơn một trăm nghìn dòng mỗi giây từ PostgreSQL. Nút thắt không nằm ở Connect.

Độ trễ bị chặn bởi poll.interval.ms. Đặt 1.000 nghĩa là trung bình chờ nửa giây trước khi Connect nhìn lại bảng. (Con số đo được còn cộng thêm khoảng 0,9 giây thời gian khởi động JVM của công cụ tôi dùng để kiểm — phần 2 đã đo. Độ trễ thật của Connect thấp hơn con số trong bảng.)

Cái bẫy: mode=incrementing mù với UPDATE và DELETE

Đây là phát hiện quan trọng nhất của bài.

Tôi chạy trong PostgreSQL:

update don_hang set trang_thai='DA_HUY', tien=0 where id <= 50000;   -- 50.000 dòng
delete from don_hang where id between 50001 and 60000;                -- 10.000 dòng

Rồi đo Kafka:

offset trước:   200.003
offset sau:     200.003        chênh 0

Không một bản ghi nào được sinh ra.

Trong PostgreSQL, 50.000 dòng đã là DA_HUY. Trong Kafka, đọc lại bản ghi id=1:

id 1   trang_thai = PENDING   tien = 137

Dữ liệu cũ, và nó sẽ ở đó vĩnh viễn.

Điều tệ nhất là không có lỗi, không cảnh báo, không dòng log nào. Connector vẫn RUNNING. Bảng theo dõi vẫn xanh. Mọi thứ trông hoàn hảo, và hạ nguồn đang tính toán trên dữ liệu sai.

Cơ chế thì rõ ràng khi nghĩ lại: mode=incrementing chạy truy vấn select * from don_hang where id > <đã_đọc_tới>. UPDATE không đổi id, nên dòng đó không bao giờ lọt vào lần đọc sau. DELETE thì càng không — dòng đã biến mất khỏi bảng.

Bốn chế độ, và cái gì chúng thấy

Chế độ Thấy được
bulk Đọc lại toàn bộ bảng mỗi lần poll — không dùng được ở quy mô thật
incrementing Chỉ INSERT
timestamp INSERT và UPDATE, nếu bảng có cột thời gian đáng tin
timestamp+incrementing Như trên, chống trùng tốt hơn khi có nhiều dòng cùng dấu thời gian

Không chế độ nào thấy DELETE.

timestamp có điều kiện riêng của nó: cột thời gian phải được mọi đường ghi cập nhật. Một câu UPDATE do ai đó chạy tay quên đụng tới cột đó là một thay đổi Connect không bao giờ thấy. Trigger giúp được, nhưng khi ấy bạn đã đang sửa lược đồ cơ sở dữ liệu nguồn — và nếu chấp nhận làm vậy thì CDC là lựa chọn tốt hơn.

Kết luận thực dụng: JDBC source chỉ phù hợp với bảng chỉ-thêm. Nhật ký sự kiện, bản ghi audit, dữ liệu đo. Với bảng nghiệp vụ có sửa và xoá, nó cho ra dữ liệu sai một cách im lặng — và phần sau đo giải pháp đúng.

Plugin phải cài riêng

Ảnh Connect chỉ có sẵn ba connector MirrorMaker. Mọi thứ khác phải cài thêm:

confluent-hub install --no-prompt confluentinc/kafka-connect-jdbc:10.7.6
# rồi khởi động lại Connect

Bước khởi động lại là chỗ dễ quên nhất khi dựng lần đầu: cài xong, gọi API tạo connector, nhận lỗi "class not found", và tưởng cài hỏng.

Kiểm bằng:

curl -s http://cn:8083/connector-plugins | jq -r '.[].class'

Sink connector: chiều ngược lại

Cùng cơ chế, ngược hướng — Kafka → cơ sở dữ liệu:

{"connector.class":"io.confluent.connect.jdbc.JdbcSinkConnector",
 "topics":"pg-don_hang",
 "connection.url":"jdbc:postgresql://dw:5432/kho",
 "insert.mode":"upsert","pk.mode":"record_key","pk.fields":"id",
 "auto.create":"true"}

insert.mode=upsert là thứ khiến việc đọc trùng (phần 8) không gây hại: đọc lại cùng bản ghi thì ghi đè, không tạo dòng thừa. Đây là ví dụ cụ thể của nguyên tắc "làm phía tiêu thụ chịu được lặp".

Ba thứ nên biết khi vận hành

Connect chạy phân tán được và nên chạy như vậy. Nhiều thực thể cùng group.id chia nhau các task, và một thực thể chết thì task được giao lại. Cấu hình, offset và trạng thái nằm trong ba topic Kafka nội bộ, nên không có trạng thái cục bộ nào cần sao lưu.

tasks.max là mức song song. JDBC source chia theo bảng: một bảng thì tasks.max lớn hơn 1 không giúp gì.

Trạng thái RUNNING của connector không có nghĩa task đang chạy. Phải xem cả hai:

curl -s http://cn:8083/connectors/pg-source/status | jq '{c:.connector.state, t:[.tasks[].state]}'

Connector RUNNING với task FAILED là tình huống có thật, và bảng theo dõi chỉ nhìn trường đầu sẽ báo khoẻ.

Thử ba mươi giây

Kiểm connector JDBC của bạn có đang bỏ sót thay đổi không:

# số dòng trong nguồn
psql -h pg -U postgres -d shop -tAc "select count(*) from don_hang"

# số bản ghi trong Kafka
kafka-get-offsets.sh --bootstrap-server kf:9092 --topic pg-don_hang

Nếu Kafka lớn hơn nguồn, đã có DELETE mà Kafka không biết. Nếu chúng bằng nhau nhưng nội dung khác — như bảng trang_thai ở trên — đã có UPDATE bị bỏ qua. Cả hai đều là dữ liệu sai đang được phục vụ.

Phần sau đo Debezium: cùng bài toán, nhưng đọc từ WAL nên thấy được cả UPDATE và DELETE.