Kafka Rebalance Storm & Các Chiến Lược Rebalance Thực Tế
Tóm tắt: Trên một consumer group động (autoscaling + pod bị lỗi), rebalance xảy ra liên tục. Với giao thức classic + chiến lược eager (mặc định của KafkaJS), mọi rebalance đều khiến gần như toàn bộ consumer dừng xử lý — cái gọi là nguyên tắc "stop-the-world". Hậu quả: lag tin nhắn chất đống, thời gian tiêu thụ tăng không kiểm soát. Bài viết này tái hiện tình huống đó trên một lab tự dựng và so sánh ba chiến lược với số liệu đo thực tế.
1. Vấn đề: Rebalance Storm trên môi trường động
Kịch bản thực tế
Giả sử bạn có một consumer group tiêu thụ một topic 25 partitions với 5 consumer instances (mỗi pod lý tưởng 5 partitions). Đây là một hệ thống CDC đồng bộ gần như real-time với lưu lượng rất cao — chúng ta không muốn dừng xử lý chỉ vì một consumer gặp sự cố.
Trên môi trường production, rebalance được kích hoạt liên tục vì ba lý do chính:
- Một consumer instance gặp sự cố (timeout, server crash) → nó ngừng xử lý và rời khỏi group → rebalance.
- HPA (Horizontal Pod Autoscaler) tạo pod mới → pod mới gia nhập group → rebalance.
- HPA thu hồi pod dư thừa → pod rời group → rebalance.
Kết quả: rebalance xảy ra thường xuyên, mỗi lần đều làm các pod xử lý bị dừng. Rebalance càng nhiều, thời gian tiêu thụ tin nhắn càng tăng → tin nhắn chất đống ngày càng nhiều.
[25 partitions] → [consumer group]
│
┌────────┬──────────┼──────────┬──────────┐
│ Pod 1 │ Pod 2 │ Pod 3 │ Pod 4 │ Pod 5
│ 5 parts│ 5 parts │ 5 parts │ 5 parts │ 5 parts
└────────┴──────────┴──────────┴──────────┘
↑ một pod chết / HPA tạo-thu hồi pod → REBALANCE → dừng toàn bộ
Vì sao nó tệ đến vậy với KafkaJS?
Nhiều đội dùng NestJS với Kafka built-in support, dựa trên KafkaJS — thư viện Kafka phổ biến nhất cho Node.js. Mặc định:
- Giao thức:
classic - Chiến lược:
roundRobin(eager)
Với giao thức classic + chiến lược eager này, bất kỳ thay đổi nào trong thành viên group (consumer tham gia/rời khỏi) hoặc metadata topic đều kích hoạt một lần dừng hoàn toàn:
Stop-the-world: Tất cả consumer thu hồi partitions của mình, một leader tính toán assignment mới, partitions được phân phối lại — và chỉ sau đó việc xử lý mới tiếp tục.
Giao thức classic sử dụng client-side để tính toán partition assignments — tức là việc quyết định nằm ở phía consumer, càng nhiều consumer/partition thì việc đạt được đồng thuận càng lâu.
Thêm nữa — một hạn chế quan trọng:
KafkaJS đã ~2 năm không được bảo trì (coi như đã bị drop). NestJS cũng không có kế hoạch chuyển sang thư viện Kafka ổn định khác cho built-in support.
Điều này gần như chặn con đường cải tiến bằng cách nâng cấp thư viện.
2. Hướng Giải Quyết
Từ vấn đề này, Kafka team đã giới thiệu các chiến lược cải tiến:
- Cooperative-Sticky: Giữ partitions không bị ảnh hưởng, chỉ thu hồi partitions cần tái phân bổ. Consumer vẫn xử lý trong lúc rebalance.
- KIP-848 (Consumer Protocol): Tái thiết kế hoàn toàn tương tác consumer-broker — rebalance được server-side xử lý, nhanh hơn và ổn định hơn.
Rào cản lớn: KafkaJS không hỗ trợ các chiến lược mới này (vì đã ngừng bảo trì). Vì vậy cần chuyển sang thư viện mới do Confluent bảo trì.
3. Chuyển sang Confluent Kafka cho JavaScript
Thư viện confluent-kafka-javascript được xây dựng dựa trên hai thư viện Kafka nổi tiếng:
- KafkaJS
- Node Rdkafka (binding của librdkafka)
Điểm mạnh:
- Được Confluent bảo trì — chính Confluent cũng là nền tảng Kafka Chorus, nên không lo về bảo trì hay cập nhật tính năng mới.
- Hỗ trợ migration dễ dàng từ hai thư viện cũ.
- Hỗ trợ đầy đủ KIP-848 và các chiến lược mới.
4. Tái Hiện Trên Lab & Số Liệu Đo Thực Tế
Tôi dựng một lab nhỏ để tự tay chứng minh câu chuyện trên, thay vì chỉ tin vào lý thuyết.
Thiết lập lab
| Thành phần | Chi tiết |
|---|---|
| Broker | Confluent Kafka 8.3.1 (KRaft, single-node, Docker) |
| Topic | repro-rebalance, 25 partitions, RF 1 |
| Consumers | 5 pod, mỗi pod ~5 partitions |
| Libraries | [email protected] (eager) + @confluentinc/[email protected] (cooperative + KIP-848) |
| Kịch bản | ổn định (20s) → kill -9 1 pod → quan sát (70s) → pod thứ 6 join → quan sát (45s) |
Mỗi chiến lược được chạy 3 lần độc lập với cùng kịch bản. Số liệu dưới đây là kết quả đo thật từ log bằng chứng.
4.1 Kết quả tổng hợp (broker ground truth)
| Chỉ số | Classic Eager (RR) | Classic Cooperative | Consumer (KIP-848) |
|---|---|---|---|
| Partitions đổi chủ khi kill 1 consumer | 19/25 (76%) | 5/25 (20%) | 5/25 (20%) |
| Partitions giữ nguyên chủ | 6/25 | 20/25 | 20/25 |
| Partitions di chuyển khi pod mới join | 21/25 (84%) | 5/25 (20%) | 5/25 (20%) |
| Số lần rebalance khởi động (cả 2 phase) | 8 | 4 | 5 |
| Stop-the-world (pod sống bị ảnh hưởng nhất) | 45 ms | 0 ms | 0 ms |
| Commit gap lớn nhất (pod sống) | 267 ms | 254 ms | 253 ms |
Đọc bảng này: Đây là toàn bộ câu chuyện trong một bảng.
- Với eager, kill một consumer khiến 76% partitions (19/25) đổi chủ, và một pod mới join làm 84% partitions xáo trộn lại. Đây chính là chữ ký của rebalance storm: một sự kiện nhỏ → hỗn loạn toàn cục.
- Với cooperative và KIP-848, chỉ 20% partitions (đúng 5 partitions mà consumer chết đang giữ) phải di chuyển. Các pod còn lại không bị chạm tới → stop-the-world = 0 ms.
4.2 Classic Protocol + Eager Strategy (Round Robin)
'group.protocol': 'classic'
'partition.assignment.strategy': 'roundrobin'
Diễn biến đo được khi một pod rời group:
- Toàn bộ pods dừng ngay lập tức khi có một pod rời/tham gia.
- Partition ownership xáo trộn gần như hoàn toàn: từ steady-state
[pod1: 0,5,10,15,20],[pod2: 1,6,11,16,21],[pod3: 2,7,12,17,22],[pod4: 3,8,13,18,23],[pod5: 4,9,14,19,24]→ sau khi pod3 chết, ownership bị tái phân bổ lại từ đầu (pod1 nhận0,4,8,12,16,20,24, v.v.). - 8 lần rebalance khởi động chỉ cho một lần kill + một lần join (mỗi pod sống thấy 2 lần rebalance ở mỗi phase).
- Client-side xử lý rebalance.
- Rebalance càng lâu khi càng nhiều consumer instances.
Kết luận: Mỗi lần leave/join, tất cả consumer trong group dừng nghe và dừng xử lý tin nhắn. Với HPA động, lag tích lũy không kiểm soát. Rủi ro cao cho production.
4.3 Classic Protocol + Cooperative Strategy
'group.protocol': 'classic'
'partition.assignment.strategy': 'cooperative-sticky'
Diễn biến đo được khi một pod rời group:
- Chỉ 5/25 partitions (đúng những partitions mà consumer chết đang giữ) bị thu hồi và tái phân bổ.
- Các pod còn lại giữ nguyên partitions và tiếp tục xử lý — stop-the-world đo được = 0 ms.
- Mỗi pod sống chỉ thấy 1 lần rebalance ở mỗi phase (so với 2 lần của eager).
- ⚠️ Nhưng: commit offset bị pause (~250 ms) cho đến khi rebalance hoàn tất.
- Vẫn client-side xử lý rebalance.
Kết luận: Đã hạn chế stop-the-world — consumer giữ partitions không bị ảnh hưởng tiếp tục xử lý, chỉ phần bị chạm mới bị ảnh hưởng. Tuy nhiên commit bị pause giữa các phase và vẫn phụ thuộc client-side. An toàn hơn eager rõ rệt.
4.4 Consumer Protocol (KIP-848) + Uniform Strategy
'group.protocol': 'consumer'
'group.remote.assignor': 'uniform'
Diễn biến đo được khi một pod rời group:
- Chỉ 5/25 partitions bị chuyển (giống cooperative), nhưng việc phân bổ được server-side tính toán ngay lập tức.
- Stop-the-world = 0 ms; các pod ngoài partitions bị chạm tiếp tục xử lý.
- Commit gap ~253 ms (tương đương cooperative trong lab này).
- Server-side xử lý rebalance: broker tự tính assignment, không cần consumer đồng thuận.
Kết luận: Ưu tiên cho production. Rebalance server-side → nhanh và ổn định. KIP-848 được thiết kế để trở thành giao thức mặc định trong Apache Kafka tương lai (dự kiến Kafka 5.0).
4.5 Một phát hiện quan trọng: độ trễ phát hiện bị chi phối bởi session timeout, không phải protocol
Đây là điều lab dạy tôi mà tài liệu lý thuyết không nói rõ:
| Chỉ số | Eager | Cooperative | KIP-848 |
|---|---|---|---|
| Độ trễ phát hiện (kill → tín hiệu rebalance đầu tiên) | 29.9 s | 32.5 s | 48.6 s |
| Phục hồi sau kill (partition cuối được gán lại) | 30.0 s | 32.5 s | 48.7 s |
Nghịch lý: KIP-848 rebalance "rẻ" hơn (chỉ 20% partitions di chuyển) nhưng phát hiện consumer chết cứng chậm hơn — vì:
- Classic dùng client-side session timeout (mặc định 30s).
- KIP-848 dùng broker-side
group.consumer.session.timeout.ms(lab để 45s).
Đây là trade-off thật: muốn KIP-848 phát hiện nhanh, phải hạ group.consumer.session.timeout.ms phía broker. Nếu không, một pod chết cứng sẽ mất ~45s mới được phát hiện — lâu hơn cả eager.
Bài học: Đừng chỉ chọn protocol rồi quên timeout. Đo cả thời gian phát hiện, không chỉ thời gian rebalance.
5. Bảng So Sánh Tổng Hợp
| Tiêu chí | Classic Eager (RR) | Classic Cooperative | Consumer (KIP-848) |
|---|---|---|---|
| Khi consumer join/leave | Dừng ngay tất cả | Giữ/xử lý partitions không bị ảnh hưởng | Giữ/xử lý partitions không bị ảnh hưởng |
| Rebalance process | Thu hồi tất cả partitions của mọi consumer | Chỉ thu hồi partitions của consumer bị sự cố | Chỉ thu hồi partitions của consumer bị sự cố |
| Partitions di chuyển (kill 1/5) | 76–84% | 20% | 20% |
| Ai tính assignment | Client (leader) | Client (leader) | Server |
| Rebalance phase | 1 phase | Nhiều phase | Nhiều phase (gọn) |
| Assignment consensus | Cần tất cả consumer đồng thuận | Cần tất cả consumer đồng thuận | Server tự assign |
| Phạm vi ảnh hưởng | Tất cả consumer | Chỉ consumer/partition bị ảnh hưởng | Chỉ consumer/partition bị ảnh hưởng |
| Stop-the-world (đo được) | 45 ms | 0 ms | 0 ms |
| Commit offset | Pause trong lúc rebalance | Pause trong lúc rebalance | Pause rất ngắn |
6. Khuyến Nghị Theo Use Case
Không có chiến lược nào "tốt nhất tuyệt đối" — tùy vào đặc điểm hệ thống:
Case A: CDC high-volume, đồng bộ gần real-time
Lưu lượng 1–2 triệu messages/ngày, sync gần real-time, partitions lớn, nhiều consumer instances, HPA, một số message type xử lý lâu.
- Khuyến nghị:
Consumer Protocol (KIP-848)(ưu tiên) - Vì: không muốn dừng tất cả consumer nếu một consumer lỗi; KIP-848 rebalance server-side → nhanh và chỉ ảnh hưởng phần bị chạm. Nhớ hạ
group.consumer.session.timeout.msđể phát hiện pod chết nhanh.
Case B: Nhiều partitions (>10), xử lý message lâu, gần real-time, HPA
- Khuyến nghị:
CooperativehoặcConsumer Protocol (KIP-848)(ưu tiên) - Vì: cho phép tiếp tục xử lý messages ở partitions không bị ảnh hưởng; chỉ delay message của partitions bị ảnh hưởng.
Case C: Ít partitions (<10), message xử lý nhanh, không cần real-time, HPA
- Khuyến nghị:
Classic with Eager (Round Robin) - Vì: với ít partitions, quá trình thu hồi + tái phân bổ diễn ra nhanh (ít consumer/partition, ít thời gian đồng thuận). Không cần real-time thì có thể chấp nhận dừng tất cả consumer.
Case D: Ít partitions, message nhanh, không real-time, không HPA
- Khuyến nghị:
Classic with Eager (Round Robin) - Vì: không có rebalance thường xuyên, eager đơn giản và đủ tốt.
7. Những Yếu Tố Khác Ảnh Hưởng Rebalance
Phần trên nói về kịch bản lý tưởng. Thời gian rebalance thực tế còn bị ảnh hưởng bởi:
- Pod terminated nhưng không disconnect Kafka connection đúng cách (kết nối mồ côi / orphaned connections) → broker phải chờ timeout mới nhận ra, làm chậm rebalance.
- Consumer group timeout được cấu hình lớn → detection lâu hơn (như phát hiện ở mục 4.5).
- Một vài yếu tố khác...
Hãy kiểm tra những điều này trước khi đổ lỗi cho strategy.
8. Kết Luận
Rebalance storm trên môi trường động (HPA + pod crash) là một rủi ro thật cho hệ thống Kafka tiêu thụ gần real-time, đặc biệt khi dùng giao thức classic + chiến lược eager — vốn là mặc định của KafkaJS.
Số liệu lab đã chứng minh rõ ràng:
- Eager: kill một consumer → 76% partitions xáo trộn, 8 lần rebalance chỉ cho 2 sự kiện.
- Cooperative / KIP-848: chỉ 20% partitions di chuyển, 0 ms stop-the-world cho các pod không bị chạm.
- Nhưng: KIP-848 phát hiện pod chết chậm hơn nếu không hạ
session.timeout— một trade-off cần lưu ý.
Ba công cụ, xếp theo mức độ hiện đại:
- Classic Eager: đơn giản, nhưng stop-the-world — không hợp production HPA.
- Classic Cooperative: giữ partitions không ảnh hưởng, tiếp tục xử lý — hạn chế được stop-the-world nhưng vẫn client-side.
- Consumer Protocol (KIP-848): tương lai của Kafka — rebalance server-side, nhanh, ổn định.
Với KafkaJS đã ngừng bảo trì, chuyển sang thư viện Confluent là bước đi đúng đắn cho production.
Trải nghiệm "goblin" của mình: Đừng chỉ tin benchmark — hãy tự tái hiện kịch bản rebalance trên môi trường test với đúng số partitions/consumer của bạn, đo cả stop-the-world và thời gian phát hiện. Số liệu thực của hệ thống bạn mới là chân lý.
Bài viết dựa trên nghiên cứu tài liệu Apache Kafka KIP-848 và thực nghiệm tái hiện trên một lab Kafka tự dựng (Confluent Kafka 8.3.1, KafkaJS + Confluent Kafka cho JavaScript).