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.
1. Position, committed offset và failure window
| Khái niệm | Ý nghĩa | Không chứng minh |
|---|---|---|
| Fetch/current position | Record tiếp theo client sẽ trả trong session | Business processing đã commit |
| Committed offset | Checkpoint restart/rebalance của group | Mọi record trước đó xử lý đúng |
| Log end offset | Cuối log hiện biết | Record đã replicated theo SLA |
| Consumer lag | Khoảng cách giữa end và position/checkpoint | Message 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.
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.
- Giữ per-partition completion tracker/watermark nếu xử lý song song.
- Không gọi consumer client từ nhiều threads nếu API không thread-safe; wakeup là cơ chế đặc biệt.
- Bound queue giữa poller và workers để backpressure thay vì heap growth.
- Pause/resume theo partition để một hot/slow partition không chiếm toàn executor.
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.