重設消費位點是指改變訂閱者當前的消費位置。當消費者出現故障或者消費錯誤資料時,您可通過重設消費位點將消費位置復原到之前的某個位點或者指定分區位點,重新開始消費。您也可以將消費位置移動至最新位點,暫時不處理堆積的訊息。
前提條件
已停止所有Consumer用戶端(雲訊息佇列 Kafka 版不支援線上重設消費位點)。
在停止Consumer用戶端後,需要經過ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG配置的時間(預設10000 ms),服務端才認為Consumer真正下線。
背景資訊
雲訊息佇列 Kafka 版支援以下重設消費位點方式:
從最新位點開始消費:不再消費Broker上堆積的訊息,將消費位點重設到最新的位置。
指定時間點開始消費:將消費位點重設到過去的某個時間點(該時間點以Topic的訊息儲存時間為準)。只要訊息仍在 Kafka 的訊息保留周期內(預設 3 天),選擇此方式即可重新消費該時間點之後的所有未消費訊息,不會遺漏。
說明執行重設前,需先停止所有 Consumer 用戶端,並等待會話逾時時間(預設 10 秒)後再執行重設,以確保重設生效。詳見前提條件。
按分區消費位點進行重設:如果只有少量分區產生訊息堆積,可以僅重設指定分區的消費位點,避免重複消費其他分區中已正確處理的訊息。
堆積的訊息本身並不會因此被刪除,改變的只是消費位點。
操作步驟
在概览頁面的资源分布地區,選擇地區。
在实例列表頁面,單擊目標執行個體名稱。
在左側導覽列,單擊Group 管理。
在Group 管理頁面,單擊目標Group ID。
在重設Group的消費位點面板,瞭解其前提条件,設定重設策略。
設定重置所有 Topic。
單擊是,重設所有Topic的消費位點。
單擊否,在Topic文字框輸入需要重設Topic的名稱。
設定重置方式。
單擊从最新位点开始消费,將消費位點指定到最新的位置,單擊确定。
單擊从指定时间点的位点开始消费,在时间点文字框,單擊
,從指定時間點的位點開啟消費功能,單擊确定。單擊按分区消费位点进行重置,在目標資料分割所在行, 消费位点文字框輸入開始消費位點值,單擊确定。
在提示對話方塊,確認提示資訊,單擊确定。
常見問題
重設消費位點能否解決部分 Topic 消費失敗問題?
重設消費位點可能緩解部分 Topic 消費失敗的情況,但建議先排查根本原因。
若消費失敗由頻繁 Rebalance 或分區分配異常導致,可按以下步驟排查:
檢查用戶端版本,確認用戶端與服務端版本相容。
調整用戶端的
session.timeout.ms和max.poll.interval.ms參數,避免因心跳逾時或訊息拉取逾時頻繁觸發 Rebalance。重設消費位點後,確認新消費組已正確訂閱目標 Topic,再重新啟動消費。
更換 Kafka 消費組 ID 後是否能繼續消費?
更換消費組 ID 後可以繼續消費,但起始消費位置取決於是否初次開機及 Topic 狀態。建議先排查導致更換消費組 ID 的根本原因(如網路異常或心跳逾時),避免新消費組後續也出現類似問題。
刪除 Kafka 消費組是否會同時刪除其訂閱的 Topic?
刪除 Kafka 消費組不會刪除其訂閱的 Topic。消費組和 Topic 是獨立的資源,刪除消費組僅會移除該消費組及其消費位點,Topic 中的資料和配置不受影響。
Spark 或 Lindorm-Spark 類型的 ConsumerGroup 如何提交消費位點?
根據 ConsumerGroup 類型,提交方式如下:
spark-kafka-source 類型:支援向 Kafka 提交消費位點。可通過
enable.auto.commit參數控制提交方式:設為
true時自動認可消費位點。設為
false時,需在消費邏輯完成後手動調用commit(offsets)函數提交位點。
Lindorm-Spark 類型:建議手動提交消費位點,以確保訊息處理完成後再提交,防止因自動認可導致位點跳躍或監控誤判:
將
enable.auto.commit設為false。在消費邏輯完成後調用
commit(offsets)函數手動提交位點。
相關文檔
如果您希望通過API來重設消費位點,請參見重設消費者組的消費位點。
重設完成後,您可以通過查看消費狀態來擷取最新的消費位點資訊。