Cân bằng lại nhóm consumer là thứ được vẽ nhiều mà đo ít. Bài này bắt từng sự kiện kèm dấu thời gian, và kết quả không khớp với lời khuyên phổ biến.

Ba lần đo cân bằng lại ở hai chiến lược, và thời gian phát hiện consumer chết

Luật cơ bản

Nhóm consumer dùng chung một group.id. Broker chia partition cho các thành viên theo hai luật cứng:

  1. Một partition chỉ thuộc về một consumer trong nhóm tại một thời điểm.
  2. Một consumer có thể giữ nhiều partition.

Hệ quả trực tiếp: số consumer > số partition thì phần thừa ngồi không. Topic 6 partition thì consumer thứ 7 không bao giờ nhận được gì.

Chạy một consumer trên topic 6 partition:

GROUP  TOPIC  PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
g1     grp    0          101062          101062          0
g1     grp    1           98024           98024          0
...    (cả sáu, cùng một CONSUMER-ID)

Một consumer ôm cả sáu.

Thêm consumer: chiến lược mặc định

Tôi bật kafka-verifiable-consumer.sh — nó in ra JSON mỗi lần nhận hoặc nhả partition, kèm dấu thời gian mili giây. Mốc 0 là lúc consumer thứ hai khởi động.

+3468 ms   C1 NHẢ    [0,1,2,3,4,5]
+3472 ms   C1 NHẬN   [3,4,5]
+3474 ms   C2 NHẬN   [0,1,2]

Đo ba lần: cửa sổ nhả–nhận là 8 ms, 4 ms, 4 ms. Nhanh hơn tôi tưởng nhiều.

Nhưng chú ý dòng đầu: C1 nhả cả sáu partition, kể cả ba cái nó sẽ nhận lại ngay sau đó. Đây là kiểu "nhả hết" — mọi thành viên buông hết, broker chia lại từ đầu, ai nấy nhận phần mới. Trong 4 ms đó toàn bộ nhóm đứng hình.

Chiến lược mặc định trong Kafka 3.9 là [RangeAssignor, CooperativeStickyAssignor] — danh sách hai phần tử, và phần tử đầu thắng. Nên mặc định vẫn là nhả hết.

Kiểu hợp tác

CooperativeStickyAssignor chỉ nhả phần phải nhường:

+3411 ms   C1 NHẢ    [3,4,5]
+3412 ms   C1 NHẬN   []
+3412 ms   C2 NHẬN   []          <- vòng một: chưa ai được gì
+6415 ms   C2 NHẬN   [3,4,5]     <- vòng hai

Partition 0, 1, 2 không xuất hiện trong lệnh nhả nào — chúng chạy liên tục suốt quá trình. Đúng như quảng cáo.

Nhưng partition 3, 4, 5 mất 3.004 ms mới về tay chủ mới, thay vì 4 ms.

Ba lần đo cho khoảng cách 3.001 / 3.004 / 3.003 ms. Con số ổn định đó chính là heartbeat.interval.ms mặc định — kiểu hợp tác cần hai vòng, và vòng thứ hai phải đợi nhịp tim kế tiếp.

Đổi lại không phải cái người ta hay kể

Cộng tổng thời gian partition ngừng chạy:

Phép tính Tổng
Nhả hết 6 partition × 4 ms 24 partition-ms
Hợp tác 3 partition × 3.004 ms 9.012 partition-ms

Hợp tác tốn hơn 375 lần theo phép cộng này.

Lời khuyên phổ biến là "luôn dùng cooperative-sticky". Với nhóm consumer đơn giản — đọc tin, ghi vào cơ sở dữ liệu, không giữ trạng thái gì — phép đo trên nói ngược lại: nhả hết xong trong 4 ms, và 4 ms thì không ai để ý.

Hợp tác thắng khi việc dừng một partition đắt:

  • Kafka Streams có state store: nhả partition là mất kho trạng thái cục bộ, nhận lại là phải dựng lại từ changelog topic — có thể mất hàng phút.
  • Consumer giữ kết nối tới hệ thống bên ngoài, mở lại tốn giây.
  • Consumer giữ bộ nhớ đệm lớn phải nạp lại.

Với những trường hợp đó, 3 giây chờ thêm cho một phần ba partition rẻ hơn nhiều so với việc dựng lại state cho toàn bộ partition.

Cách chọn thực dụng: dùng Kafka Streams hay giữ trạng thái thì bật hợp tác. Consumer không trạng thái và nhóm ổn định thì mặc định đã đủ, và đơn giản hơn.

Một lưu ý khi chuyển: đổi sang cooperative-sticky phải triển khai hai lần — lần một thêm nó vào cuối danh sách chiến lược, lần hai bỏ chiến lược cũ đi. Đổi thẳng một lần khiến các thành viên không thống nhất được chiến lược và nhóm không hình thành.

Consumer chết đột ngột

Đóng đàng hoàng thì consumer gửi lệnh rời nhóm và việc chia lại xảy ra ngay. Bị SIGKILL thì không.

Tôi giết một consumer bằng kill -9:

+29.767 ms   C1 NHẢ    [3,4,5]
+29.770 ms   C1 NHẬN   [0,1,2,3,4,5]

Gần 30 giây mới phát hiện. Đó là session.timeout.ms của công cụ đo (30.000 ms). Mặc định của chính thư viện Kafka là 45.000 ms — kiểm bằng cách đọc thẳng ConsumerConfig.configDef().defaultValues().

Nghĩa là: một pod bị OOMKilled khiến partition của nó ngừng chạy tới 45 giây, và không có gì báo động trong khoảng đó ngoài lag tăng dần.

Hai điều rút ra:

  • Luôn gọi consumer.close() khi tắt. Trong Spring Boot là để container tự lo; trong ứng dụng viết tay là Runtime.addShutdownHook. Chênh lệch giữa đóng sạch và bị giết là 45 giây.
  • Hạ session.timeout.ms chỉ giúp một phần và có giá của nó: quá thấp thì một lần GC dài hoặc một nhịp mạng chậm cũng bị coi là chết, và nhóm cân bằng lại liên tục. Sàn của broker là group.min.session.timeout.ms (mặc định 6.000).

max.poll.interval.ms — thứ hay gây "cân bằng lại vô cớ"

Có một thời gian chờ thứ hai, tách hẳn khỏi nhịp tim: max.poll.interval.ms, mặc định 300.000 ms. Nếu vòng lặp của bạn xử lý một lô lâu hơn con số đó mà chưa gọi poll() lại, broker coi consumer đã chết và chia lại partition — dù nhịp tim vẫn đều, vì nhịp tim chạy ở luồng riêng.

Đây là nguồn của phần lớn các ca "nhóm cân bằng lại liên tục mà không rõ vì sao": max.poll.records mặc định 500, mỗi bản ghi gọi một API chậm 1 giây, thế là 500 giây cho một lô — vượt trần.

Chữa bằng cách hạ max.poll.records, không phải nâng max.poll.interval.ms. Lô nhỏ thì vòng lặp quay lại poll() thường xuyên hơn, và consumer chết thật cũng được phát hiện sớm hơn.

Thử ba mươi giây

Xem nhóm của bạn đang chia thế nào và ai đang ngồi không:

kafka-consumer-groups.sh --bootstrap-server kf:9092 --describe --group nhom-cua-ban

Đếm số CONSUMER-ID khác nhau. Nếu nó nhỏ hơn số tiến trình bạn đang chạy, phần chênh là những consumer không nhận được partition nào — và bạn đang trả tiền cho chúng mà không nhận lại gì.

Phần sau đo việc chốt offset: cùng một lần sập, ba cách chốt cho ra 1.150 tin trùng, 50 tin trùng, hoặc 50 tin mất hẳn.