Phần trước đo cân bằng lại khi thêm consumer. Phần này đo hai thứ khác: chiến lược chia có thật sự chia đều không, và chuyện gì xảy ra khi bạn triển khai bản mới.

Ba chiến lược chia và tác dụng của thành viên tĩnh

Chiến lược mặc định không chia đều

Hai topic, mỗi topic 3 partition, hai consumer cùng đăng ký cả hai. Sáu partition, hai consumer — chờ đợi hiển nhiên là 3 và 3.

RangeAssignor — mặc định của Kafka:

C1   [ta-2, tb-2]                     2 partition
C2   [ta-0, ta-1, tb-0, tb-1]         4 partition

2 so với 4. Một consumer làm gấp đôi.

Lý do nằm ở chữ "Range": nó chia riêng từng topic. Với topic ta, 3 partition cho 2 consumer thành 2 và 1. Với topic tb, cũng 2 và 1 — và lệch về cùng một phía, vì thứ tự consumer không đổi giữa các topic.

Hai topic thì lệch 4:2. Năm topic thì lệch 10:5. Càng nhiều topic, càng lệch, và lệch luôn cùng một chiều.

Hai chiến lược còn lại nhìn toàn cục:

RoundRobinAssignor
   C1   [ta-0, ta-2, tb-1]            3
   C2   [ta-1, tb-0, tb-2]            3

StickyAssignor
   C1   [ta-0, ta-1, tb-0]            3
   C2   [ta-2, tb-1, tb-2]            3

Đều tuyệt đối.

Khi nào chuyện này quan trọng: consumer đăng ký nhiều topic và số partition mỗi topic không chia hết cho số consumer. Nhóm chỉ đọc một topic thì Range chia đều như mọi chiến lược khác — đó là lý do phần lớn người dùng không bao giờ gặp.

Nếu bạn đang có một consumer luôn nóng hơn hẳn những cái khác mà không rõ vì sao, chạy kafka-consumer-groups.sh --describe và đếm số dòng mỗi CONSUMER-ID. Chênh lệch 2:4 hiện ra ngay.

Khởi động lại một consumer: hai lần cân bằng

Đây là chuyện xảy ra mỗi lần bạn triển khai bản mới. Tôi cho C2 thoát rồi bật lại ngay:

+0 ms       C2  NHẢ    [ta-2, tb-2]
+1.925 ms   C1  NHẢ    [ta-0, ta-1, tb-0, tb-1]     <- C1 bị lôi vào
+1.931 ms   C1  NHẬN   tất cả 6
+3.938 ms   C1  NHẢ    tất cả 6                      <- lần thứ hai
+3.945 ms   C1  NHẬN   [ta-2, tb-2]
+3.949 ms   C2b NHẬN   [ta-0, ta-1, tb-0, tb-1]

Hai lần cân bằng cho một lần khởi động lại. Một lần khi C2 rời, một lần khi nó quay lại. C1 chẳng liên quan gì mà phải dừng hai lượt và bị đổi partition hai lần.

Nhân con số này với số pod trong một lần triển khai cuốn chiếu: 10 pod là 20 lần cân bằng, và mỗi lần đều là toàn nhóm dừng. Với consumer giữ trạng thái, đó là 20 lần dựng lại state store.

Chú ý thêm: C1 nhận lại [ta-2, tb-2] chứ không phải [ta-0, ta-1, tb-0, tb-1] như trước. Partition nhảy lung tung giữa các consumer, và mọi bộ nhớ đệm gắn với partition đều mất.

Thành viên tĩnh xoá sạch chuyện đó

Thêm đúng một thuộc tính:

group.instance.id=inst-2

Cùng kịch bản:

+0 ms       C2  NHẢ    [ta-2, tb-2]
+3.194 ms   C2b NHẬN   [ta-2, tb-2]      <- đúng partition cũ

C1 không có sự kiện nào. Nó chạy liên tục suốt quá trình. Không lần cân bằng nào xảy ra.

Cơ chế: group.instance.id biến consumer thành thành viên tĩnh. Khi nó rời đi, broker không chia lại ngay mà giữ chỗ cho nó. Consumer mới khai cùng group.instance.id được nhận lại đúng phần cũ, và nhóm coi như chưa có gì xảy ra.

Thời gian gián đoạn 3.194 ms kia là thời gian khởi động JVM của consumer, không phải chi phí của Kafka.

Với triển khai cuốn chiếu 10 pod: từ 20 lần cân bằng xuống 0.

Điều kiện, và cái giá

Khởi động lại phải xong trong session.timeout.ms. Quá hạn thì broker coi là chết thật và cân bằng lại như thường.

Nên với thành viên tĩnh, session.timeout.ms phải dài hơn thời gian khởi động pod — thường là 60.000 đến 120.000 thay vì 45.000 mặc định. Trong Kubernetes, group.instance.id lấy từ tên pod của StatefulSet là tự nhiên nhất, vì tên đó ổn định qua các lần khởi động lại.

Cái giá phải trả rõ ràng: một consumer chết thật cũng bị giữ chỗ nguyên chừng đó thời gian. Với session.timeout.ms=120000, một pod bị OOMKilled khiến partition của nó ngừng chạy hai phút. Không có cấu hình nào phân biệt được "đang triển khai" với "vừa chết" — bạn chọn một con số và chấp nhận cả hai mặt.

Cách chọn: đo thời gian khởi động thật của ứng dụng, nhân đôi, và đó là session.timeout.ms. Đừng đặt cao hơn "để chắc ăn" — mỗi giây thừa là một giây partition đứng khi có sự cố thật.

Ba thứ có thể gây cân bằng lại

Xếp theo mức độ khó chẩn đoán:

  1. Thành viên vào hoặc ra. Dễ thấy nhất, có trong log.
  2. Số partition của topic thay đổi. --alter --partitions kích hoạt cân bằng lại trên mọi nhóm đang đọc topic đó.
  3. max.poll.interval.ms bị vượt. Khó nhất, vì không có gì báo ngoài dòng log Member ... sending LeaveGroup request lẫn giữa hàng nghìn dòng khác. Nhịp tim vẫn đều — nó chạy ở luồng riêng — nên mọi số đo về sức khoẻ đều bình thường.

Cách nhận ra trường hợp 3: nhóm cân bằng lại theo chu kỳ đều đặn xấp xỉ max.poll.interval.ms, và luôn cùng một consumer khởi xướng. Chữa bằng cách hạ max.poll.records, không phải nâng thời gian chờ.

Thử ba mươi giây

Đếm số lần cân bằng nhóm của bạn đã trải qua:

docker compose logs consumer 2>&1 | grep -c "Successfully joined group"

Chia cho số ngày ứng dụng đã chạy. Nếu ra hơn vài lần mỗi ngày mà bạn không triển khai bấy nhiêu lần, có gì đó đang kích hoạt cân bằng lại — và ba nguyên nhân ở trên bao gần hết các trường hợp.

Phần sau đo độ trễ đầu-cuối và tách nó ra từng chặng: chặng nào đắt nhất có thể không phải chặng bạn nghĩ.