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.
Mọi thao tác vận hành đều qua REST API của worker (mặc định cổng 8083):
| Method & Endpoint | Tác dụng |
|---|---|
GET /connectors | Liệt kê tất cả connector |
GET /connectors/{name}/status | Xem trạng thái connector & từng task |
POST /connectors | Tạo connector mới (body là JSON config) |
PUT /connectors/{name}/config | Cập nhật cấu hình connector |
POST /connectors/{name}/restart | Khởi động lại connector |
PUT /connectors/{name}/pause | Tạm dừng |
PUT /connectors/{name}/resume | Chạ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"
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 |
FAILED | Gặp lỗi và đã dừng — xem trace trong status để biết nguyên nhân |
PAUSED | Bị tạm dừng chủ động qua API |
UNASSIGNED | Chưa được gán cho worker nào (đang rebalance) |
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"
}
none).| Lỗi | Cá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 DB | Task 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 replay | Xó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 task | Bật errors.tolerance=all + DLQ để cô lập message lỗi |
localhost:8080 để xem topic, message và consumer lag; gọi REST API tại localhost:8083 để kiểm tra status connector.