Kafka Consumer Groups thực sự hoạt động như thế nào?
Partition được chia cho consumer ra sao, offset được lưu ở đâu, và vì sao rebalance có thể làm cả hệ thống đứng hình vài giây.
Mục lục
Hầu hết tutorial về Kafka dừng lại ở câu “consumer group giúp scale việc đọc message”. Câu đó đúng, nhưng không giải thích được những gì bạn sẽ gặp ở production: consumer đứng im dù topic còn đầy message, message bị xử lý hai lần sau khi deploy, hay cả group “đóng băng” vài giây mỗi khi một pod khởi động lại.
Bài này đi từ mô hình cơ bản đến cơ chế rebalance — nguồn gốc của phần lớn các sự cố đó.
Mô hình cơ bản: partition là đơn vị song song
Một topic được chia thành nhiều partition. Trong một consumer group, mỗi partition được giao cho đúng một consumer tại một thời điểm, còn một consumer có thể nhận nhiều partition[1].
Hệ quả trực tiếp: số consumer làm việc thực sự không bao giờ vượt quá số partition.
| Số partition | Số consumer trong group | Phân bổ | Ghi chú |
|---|---|---|---|
| 6 | 2 | mỗi consumer 3 partition | Bình thường |
| 6 | 4 | 2 consumer nhận 2, 2 consumer nhận 1 | Tải không đều |
| 6 | 6 | mỗi consumer 1 partition | Song song tối đa |
| 6 | 8 | 6 consumer làm việc, 2 consumer rảnh | Thêm consumer không giúp gì |
Offset: consumer nhớ mình đã đọc đến đâu
Kafka không xoá message khi consumer đọc xong. Thay vào đó, mỗi group lưu lại offset — vị trí đã xử lý — cho từng partition vào một topic nội bộ tên là __consumer_offsets.
Thời điểm commit offset quyết định ngữ nghĩa giao nhận:
- Commit trước khi xử lý → nếu consumer chết giữa chừng, message bị mất (at-most-once).
- Commit sau khi xử lý → nếu consumer chết sau khi xử lý nhưng trước khi commit, message bị xử lý lại (at-least-once).
Trong thực tế, at-least-once kết hợp với xử lý idempotent là lựa chọn phổ biến nhất:
Properties props = new Properties();props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-service");// Tắt auto commit: chỉ commit sau khi xử lý thành côngprops.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) { consumer.subscribe(List.of("orders")); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500)); for (ConsumerRecord<String, String> record : records) { orderService.handle(record.key(), record.value()); // phải idempotent } consumer.commitSync(); }}Rebalance: khi group thay đổi thành viên
Rebalance là quá trình chia lại partition cho các consumer. Nó xảy ra khi:
- Một consumer mới tham gia group (scale up, deploy phiên bản mới).
- Một consumer rời group — chủ động (shutdown) hoặc bị coi là đã chết.
- Số partition của topic thay đổi, hoặc tập topic mà group subscribe thay đổi.
Một broker đóng vai trò group coordinator theo dõi thành viên của group thông qua heartbeat. Có hai ngưỡng thời gian quan trọng, rất hay bị nhầm lẫn:
| Cấu hình | Ý nghĩa | Khi vượt ngưỡng |
|---|---|---|
session.timeout.ms |
Thời gian tối đa coordinator chờ heartbeat | Consumer bị coi là chết → rebalance |
heartbeat.interval.ms |
Chu kỳ gửi heartbeat (nên nhỏ hơn nhiều so với session timeout) | — |
max.poll.interval.ms |
Khoảng cách tối đa giữa hai lần gọi poll() |
Consumer tự rời group → rebalance |
Eager vs. cooperative rebalance
Với giao thức rebalance truyền thống (eager), mọi consumer phải thu hồi toàn bộ partition trước khi nhận phân bổ mới. Trong khoảng thời gian đó, cả group ngừng xử lý — hiện tượng thường được gọi là “stop-the-world”.
Giao thức incremental cooperative rebalancing (KIP-429[3]) chỉ thu hồi những partition thực sự cần chuyển chủ, các partition còn lại tiếp tục được xử lý. Với Java client, bạn bật nó bằng assignor tương ứng:
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignorXa hơn nữa, KIP-848[4] chuyển logic phân bổ partition từ phía client về broker, giúp rebalance nhẹ hơn và dễ vận hành hơn. Hãy kiểm tra phiên bản Kafka và client bạn đang dùng trước khi bật giao thức mới.
Quan sát consumer group trong thực tế
Công cụ dòng lệnh đi kèm Kafka cho biết partition nào đang thuộc consumer nào và lag — số message chưa xử lý — của từng partition:
kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --describe --group order-serviceTóm tắt
- Partition là đơn vị song song; số consumer hữu ích bị giới hạn bởi số partition.
- Offset được lưu theo group, và thời điểm commit quyết định at-most-once hay at-least-once.
- Rebalance là cái giá phải trả cho tính đàn hồi; hiểu các timeout và dùng cooperative rebalancing giúp giảm đáng kể thời gian gián đoạn[2].
Tài liệu tham khảo
- [1]Apache Kafka Documentation. Apache Software Foundation.
- [2]Martin Kleppmann. Designing Data-Intensive Applications. O'Reilly Media, 2017. Chương 11 — Stream Processing.
- [3]KIP-429: Kafka Consumer Incremental Rebalance Protocol. Apache Kafka Improvement Proposals.
- [4]KIP-848: The Next Generation of the Consumer Rebalance Protocol. Apache Kafka Improvement Proposals.