Part 08 · Kafka · 8.1.08

Consumer groups, offsets và rebalancing

Consumer group phân partitions cho members để scale một logical application. Committed offset là recovery checkpoint, không phải bằng chứng business side effect đã thành công hoặc đúng một lần.

Ownership invariant: tại một thời điểm ổn định, mỗi partition được một member trong group sở hữu. Khi ownership đổi, old work, offset commit và new owner phải được phối hợp để tránh loss hoặc duplicate ngoài dự kiến.

1. Position, committed offset và failure window

Khái niệmÝ nghĩaKhông chứng minh
Fetch/current positionRecord tiếp theo client sẽ trả trong sessionBusiness processing đã commit
Committed offsetCheckpoint restart/rebalance của groupMọi record trước đó xử lý đúng
Log end offsetCuối log hiện biếtRecord đã replicated theo SLA
Consumer lagKhoảng cách giữa end và position/checkpointMessage age hoặc handler health

Commit trước side effect tạo loss window: crash sau commit làm group bỏ qua work chưa hoàn tất. Process rồi commit tạo duplicate window: crash sau DB commit nhưng trước offset commit sẽ đọc lại. Manual commit chỉ cho quyền kiểm soát window; idempotency/transaction strategy mới xử lý duplicate.

2. Group lifecycle và assignment

Consumer subscribe/join group, coordinator quản membership, assignor phân partitions, member poll/fetch/process và heartbeat/commit. Member join/leave, subscription/topic partition thay đổi, session timeout hoặc max poll violation có thể kích hoạt rebalance.

Eager assignment revoke toàn bộ trước assign lại; cooperative assignment chuyển dần partitions để giảm stop-the-world nhưng application vẫn phải xử lý callbacks cho revoked/lost/assigned. Static membership có thể giảm churn khi restart ngắn nhưng không xóa failure detection hay ownership fencing.

Rebalance callback là correctness boundary. Dừng nhận work mới cho partition revoked, drain hoặc cancel bounded work, commit safe offsets nếu còn ownership, và không để worker cũ commit sau khi partition đã chuyển.

3. Poll loop, heartbeats và long processing

session.timeout.ms liên quan liveness/heartbeat; max.poll.interval.ms giới hạn thời gian giữa các poll để phát hiện application không tiến triển. Nếu processing vượt max poll interval, partition có thể chuyển sang member khác trong khi old task vẫn chạy, tạo concurrent duplicate.

Với work dài, giảm max.poll.records, pause partitions, tiếp tục poll đúng protocol và chuyển work vào bounded executor. Resume khi capacity trở lại. Tăng timeout vô hạn chỉ che under-capacity và làm failover chậm.

4. Parallel processing và commit theo partition

Ordering đến từ partition; chạy records cùng partition song song có thể hoàn tất out of order. Chỉ commit offset liên tục cao nhất mà mọi record trước nó đã hoàn tất. Commit offset 105 khi record 103 còn pending sẽ làm mất 103 sau crash.

5. Lag, age và capacity

Lag có thể tính theo current hoặc committed offset; hiểu metric exporter đang dùng. Cần xem thêm record timestamp/business age, ingress/processing rate, partition skew và retry/DLQ. Lag bằng 0 vẫn có thể che data loss nếu commit quá sớm hoặc handler nuốt lỗi.

Một group không dùng hữu ích nhiều active consumers hơn số partitions; extra members idle. Scale consumer chỉ hiệu quả nếu partitions và downstream capacity cho phép. Hot partition có thể giữ lag cao dù group tổng thể còn idle.

6. Recovery tests và observability

Test kill tại trước/sau side effect và trước/sau commit; inject rebalance khi work đang chạy; test poison record, long handler, coordinator/network interruption và offset out-of-range. Reconciliation theo business ID chứng minh duplicate/loss, không chỉ nhìn lag.

Theo dõi assigned partitions, rebalance count/duration/reason, commit latency/errors, poll interval, records rate, per-partition lag/age, paused duration, worker queue depth và processing outcome.

Review checklist: commit timing có failure-window doc; handler idempotent; revoke/lost callbacks fence old work; poll loop không block quá lâu; executor bounded; per-partition ordering/commit watermark đúng; lag đi cùng age và business reconciliation.
Tài liệu: Consumer API · Consumer Configurations · Offset Tracking · ConsumerRebalanceListener · max.poll.interval.ms