Kafka Connect là một framework mã nguồn mở đi kèm Apache Kafka, chuyên để di chuyển dữ liệu giữa Kafka và các hệ thống bên ngoài một cách ổn định, có thể mở rộng. Bài này giải thích kiến trúc và các khái niệm bạn buộc phải nắm trước khi viết bất kỳ connector nào.
Khi muốn đưa dữ liệu từ một database vào Kafka, hoặc lấy dữ liệu từ Kafka ghi xuống một database khác, cách thủ công là tự viết ứng dụng Kafka Producer/Consumer. Việc này lặp đi lặp lại, dễ sai, và phải tự lo offset, retry, scale, chịu lỗi.
Kafka Connect giải quyết đúng bài toán đó: nó là một dịch vụ chạy riêng (worker process) mà bạn chỉ cần khai báo cấu hình JSON là dữ liệu chạy. Connect lo hết phần khó: theo dõi offset, song song hóa, tự khởi động lại khi lỗi, rebalance khi thêm/bớt worker.
Có hai loại connector: Source (kéo dữ liệu từ ngoài VÀO Kafka) và Sink (đẩy dữ liệu từ Kafka RA hệ thống ngoài). Cộng đồng và Confluent cung cấp sẵn hàng trăm connector (Debezium, JDBC, S3, Elasticsearch...).
Connector — đơn vị logic định nghĩa "lấy dữ liệu từ đâu, đưa đi đâu". Bản thân connector không di chuyển dữ liệu; nó chịu trách nhiệm chia công việc thành các Task. Bạn tạo connector bằng một file cấu hình JSON.
Task — đơn vị thực sự di chuyển dữ liệu. Một connector có thể sinh ra nhiều task để xử lý song song (vd mỗi task đọc một nhóm bảng / một nhóm partition). Số task tối đa khai báo qua tasks.max.
Worker — tiến trình JVM chạy connector và task. Có thể chạy một worker (standalone) hoặc nhiều worker thành cluster (distributed). Worker là thứ bạn thật sự khởi động.
Converter — bộ chuyển đổi giữa dữ liệu nội bộ của Connect và định dạng byte trên Kafka. Quyết định message được serialize thành JSON, Avro hay Protobuf. Source dùng converter để ghi, Sink dùng để đọc — hai bên phải khớp nhau.
Transform (SMT – Single Message Transform) — biến đổi từng message ngay trong luồng, trước khi ghi vào Kafka (source) hoặc trước khi ghi ra ngoài (sink). Ví dụ: bỏ field, đổi tên topic, unwrap cấu trúc Debezium. Cấu hình bằng JSON, không cần code.
JsonConverter (dễ đọc, không cần hạ tầng thêm) và AvroConverter (gọn hơn, có ràng buộc schema nhưng cần Schema Registry). Tách riêng converter cho key và value qua key.converter / value.converter. Với JSON thường tắt schema bằng value.converter.schemas.enable=false cho gọn.| Tiêu chí | Standalone | Distributed |
|---|---|---|
| Cách chạy | Một tiến trình duy nhất, cấu hình trong file .properties | Nhiều worker chung một group.id, tạo connector qua REST API |
| Chịu lỗi | Không — worker chết là dừng hết | Có — task được rebalance sang worker còn sống |
| Scale | Không scale ngang được | Thêm worker là tự chia lại task |
| Lưu offset | Trong file cục bộ trên đĩa | Trong các Kafka topic nội bộ (offset/config/status) |
| Dùng khi nào | Dev, test, demo nhỏ một máy | Mọi môi trường production |
Ở môi trường thật gần như luôn dùng distributed mode, kể cả khi chỉ có một worker — vì khi đó mọi cấu hình nằm trong Kafka, dễ vận hành và sẵn sàng scale.
| Source Connector | Sink Connector | |
|---|---|---|
| Hướng | Hệ thống ngoài → Kafka | Kafka → Hệ thống ngoài |
| Vai trò Kafka | Là Producer | Là Consumer |
| Ví dụ | Debezium (CDC từ DB), FileStreamSource | JDBC Sink, Elasticsearch, S3 |
| Cấu hình đặc trưng | table.include.list, topic.prefix | topics, connection.url, insert.mode |
Mỗi connector là một bộ JAR riêng (Debezium, JDBC...). Worker tìm chúng trong thư mục khai báo bởi plugin.path. Mỗi plugin nên nằm trong một thư mục con riêng để tránh xung đột class.
# trong worker config (vd connect-distributed.properties)
plugin.path=/usr/share/java,/opt/connectors
Vì sao dùng Kafka Connect thay vì tự viết Producer/Consumer?