Kafka Connect

Vận hành & Xử lý lỗi

Sau khi connector chạy, công việc là theo dõi và xử lý sự cố. Kafka Connect cung cấp một REST API đầy đủ để tạo, sửa, tạm dừng, khởi động lại và giám sát connector — cộng với cơ chế Dead Letter Queue để không cho một message hỏng làm chết cả pipeline.

Quản lý qua REST API
Cổng 8083

Mọi thao tác vận hành đều qua REST API của worker (mặc định cổng 8083):

Method & EndpointTác dụng
GET /connectorsLiệt kê tất cả connector
GET /connectors/{name}/statusXem trạng thái connector & từng task
POST /connectorsTạo connector mới (body là JSON config)
PUT /connectors/{name}/configCập nhật cấu hình connector
POST /connectors/{name}/restartKhởi động lại connector
PUT /connectors/{name}/pauseTạm dừng
PUT /connectors/{name}/resumeChạy lại sau khi pause
DELETE /connectors/{name}Xóa connector

Ví dụ xem trạng thái và khởi động lại:

# xem trạng thái connector + task
curl http://localhost:8083/connectors/jdbc-sink-orders/status

# khởi động lại connector
curl -X POST http://localhost:8083/connectors/jdbc-sink-orders/restart

# khởi động lại riêng task số 0 (kèm cả connector)
curl -X POST "http://localhost:8083/connectors/jdbc-sink-orders/restart?includeTasks=true&onlyFailed=true"
Trạng thái connector & task
RUNNING · FAILED · PAUSED

Một connector và mỗi task của nó có trạng thái riêng. Điều quan trọng: connector RUNNING nhưng task vẫn có thể FAILED — nên luôn đọc cả phần tasks trong status.

Trạng tháiÝ nghĩa
RUNNINGĐang chạy bình thường
FAILEDGặp lỗi và đã dừng — xem trace trong status để biết nguyên nhân
PAUSEDBị tạm dừng chủ động qua API
UNASSIGNEDChưa được gán cho worker nào (đang rebalance)
Xử lý lỗi & Dead Letter Queue
Đừng để một message hỏng giết pipeline

Mặc định, một message không deserialize/transform được sẽ làm task FAILED. Với sink connector, ta nên cấu hình bỏ qua message lỗi và đẩy chúng vào Dead Letter Queue (DLQ) để điều tra sau:

{
  "errors.tolerance": "all",
  "errors.deadletterqueue.topic.name": "dlq.jdbc-sink",
  "errors.deadletterqueue.topic.replication.factor": "1",
  "errors.deadletterqueue.context.headers.enable": "true",
  "errors.log.enable": "true",
  "errors.log.include.messages": "true"
}
Lưu ý: DLQ chỉ áp dụng cho lỗi converter và SMT của sink connector. Lỗi do bản thân connector (vd mất kết nối DB) không vào DLQ mà sẽ làm task FAILED rồi retry.
Lỗi thường gặp & cách xử lý
Bộ checklist khi pipeline kẹt
LỗiCách xử lý
Converter mismatch (deserialize fail)Kiểm tra source & sink cùng converter (cùng JSON hoặc cùng Avro + Schema Registry); chỉnh schemas.enable cho khớp
Schema thay đổi (thêm/xóa cột)Bật auto.evolve=true ở JDBC sink; với Avro dùng schema tương thích (backward/forward)
Mất kết nối DBTask FAILED rồi tự retry; kiểm tra DB sống, credential, network; restart task khi DB trở lại
Offset lệch / muốn replayXóa offset ở consumer group (sink) hoặc reset offset topic (source); cẩn thận vì có thể ghi trùng
Message hỏng làm chết taskBật errors.tolerance=all + DLQ để cô lập message lỗi
Best practices vận hành
Cho pipeline ổn định lâu dài
Khớp với cdc-lab: mở Kafka UI ở localhost:8080 để xem topic, message và consumer lag; gọi REST API tại localhost:8083 để kiểm tra status connector.
← Quay lại
Bài 3: Sink JDBC