
線上發版或者業務高峰期有時會遇到一個奇怪的現象消費端的服務進程沒掛JVM 內存正常日志也在正常輸出顯示業務代碼還在繼續處理數據。但 Kafka 服務端卻判定該消費者已失效強行將其踢出消費組并觸發了 Rebalance。這種在進程和日志均正常的情況下依然被判定為失效的現象在技術上通常被稱為假死。為什么有心跳也會被判定死亡Kafka 極早期的版本消費端是單線程的。這個單線程要同時負責兩件事拉取數據并執行寫庫、打日志等業務邏輯還要給服務端發心跳包報告自己還活著。但這個設計在性能上存在缺陷如果業務邏輯卡住了比如查數據庫超時或者 JVM 發生 Full GC主線程被卡住就無法按時發心跳。服務端等不到心跳就會誤認為客戶端掛了立刻觸發 Rebalance。后來 Kafka 把拉數據和發心跳徹底拆成了兩個線程。雙線程架構重構之后啟動一個 Kafka 消費者底層實際上是兩個線程在協作。心跳線程只負責在后臺按照固定頻率給 Broker 發心跳包。只要心跳不斷服務端就認為你的進程還活著。它判斷存活的閾值是session.timeout.ms默認是 45 秒。另一個業務主線程負責在循環里調用poll()拉取數據并執行業務邏輯。它判斷是否卡死的閾值是max.poll.interval.ms默認是 5 分鐘。這里有個最容易發生沖突的地方。假設你的消費者一次性拉了 500 條消息因為下游寫庫慢主業務線程處理這 500 條消息一共花了 6 分鐘。在這 6 分鐘里主線程正在處理數據并輸出日志。由于這 500 條數據還沒處理完它是沒辦法回到循環去調用下一次poll()的。可后臺的心跳線程依然在按時給服務端發心跳。但在服務端看來距離這臺機器上一次調用poll()已經過去 6 分鐘了超過了規定的 5 分鐘死線。服務端就會判定這個消費者雖然還有心跳但它已經失去了消費能力占著分區卻不拉取新數據導致數據積壓。為了不影響整個消費組的吞吐必須判定它為假死強行將其踢出消費組這就引發了 Rebalance。怎么避免假死知道是由于業務處理太慢導致poll()間隔超時其實解決思路非常清晰。1. 減小拉取批次客戶端配置里限制每次poll()拉取的最大消息數max.poll.records: 50把大批次拆成高頻的小批次。即使單條消息處理需要 100 毫秒50 條也只需要 5 秒鐘。主線程處理完能快速回到下一次poll()循環避免超出 5 分鐘超時。2. 調大間隔參數根據最壞的業務場景比如下游服務宕機、網絡抖動等合理調大最大 poll 間隔時間max.poll.interval.ms: 600000不過這通常要配合max.poll.records一起微調不建議設得無限大。因為如果主線程真的徹底死鎖卡死了服務端需要等 10 分鐘才能發現這期間對應的分區就會被其他消費者接管造成數據消費中斷。3. 異步多線程消費如果單條業務邏輯確實極其耗時無論怎么調小 records 都不行那就不要在 Kafka 的消費主線程里直接干活。可以只讓主線程負責拉數據拉完立刻丟進自定義的 Java 線程池里去并行處理讓主線程秒級回到下一次poll()。不過這個方案會引入消息丟失和重復消費的風險自動提交可能導致消息丟失如果開啟了自動提交主線程下一次poll()就會把上一批數據對應的 offset 提交。但此時線程池里可能還有一堆任務在排隊。一旦服務此時重啟或崩潰這些還沒來得及處理的消息就徹底丟失了。手動提交的亂序與重復消費要知道線程池里的并發線程是亂序完成的如果在子線程里直接提交 offset會發生 offset 覆蓋導致重復消費。要解決這個問題需要自己設計一套滑窗提交機制或者基于每個 Partition 單獨配內存隊列開發成本非常高。說在最后分布式系統對活著其實有兩種定義一個是進程存活代表進程在端口通心跳還在跳另一個是業務存活代表業務還能響應外部輸入推動數據往下走。寫業務消費邏輯時不能僅憑進程在日志在刷就認為消費者運行正常需要合理評估max.poll.records和max.poll.interval.ms的配置避免因為單次處理耗時過長導致不必要的 Rebalance。