全部產品
Search
文件中心

Tablestore:消費通道資料

更新時間:Aug 01, 2026

Tablestore Java SDK 可通過通道持續消費資料,使用消費回調處理各批記錄,並配置心跳、Checkpoint、線程池和消費並行度。

注意事項

  • 增量日誌的保留時間與資料表的 Stream 日誌到期時間一致,最長為 7 天。使用全量加增量類型的通道時,如果全量資料未能在增量日誌保留時間內消費完成,開始消費增量資料時會返回 OTSTunnelExpired 錯誤,無法繼續消費增量資料。

  • 消費增量資料的進度落後於增量日誌保留時間時,通道可能從當前仍可用的最新資料開始消費,導致部分資料未被消費。

  • 通道到期後可能被禁用。通道連續處于禁用狀態超過 30 天后會被刪除,刪除後無法恢複。

前提條件

安裝Tablestore Java SDK,並初始化 TunnelClient

功能說明

TunnelWorker 根據通道 ID 串連通道,通過心跳擷取分配給當前用戶端的 Channel,持續拉取資料,並將每批記錄傳入 IChannelProcessor。多個 TunnelWorker 消費同一通道時,服務端會在各用戶端之間分配 Channel。

消費通道資料包括以下步驟:

  1. 實現 IChannelProcessor 介面。process 方法處理每批記錄,shutdown 方法釋放消費回調使用的資源。

  2. 建立 TunnelWorkerConfig,配置資料處理回調和消費參數。

  3. 使用通道 ID、TunnelClientTunnelWorkerConfig 建立 TunnelWorker

  4. 調用 connectAndWorking 方法啟動消費。

    void process(ProcessRecordsInput input);
    void shutdown();

以下樣本列印通道拉取到的每條記錄,然後啟動消費。

private static class SimpleProcessor implements IChannelProcessor {
    @Override
    public void process(ProcessRecordsInput input) {
        for (StreamRecord record : input.getRecords()) {
            System.out.println(record);
        }
    }

    @Override
    public void shutdown() {
        // 釋放消費回調使用的資源。
    }
}

String tunnelId = "example_tunnel_id";
TunnelWorkerConfig config =
        new TunnelWorkerConfig(new SimpleProcessor());
TunnelWorker worker =
        new TunnelWorker(tunnelId, tunnelClient, config);
worker.connectAndWorking();
重要

connectAndWorking 啟動後台消費任務後會返回,應用進程需要保持運行。停止消費時,依次調用 worker.shutdown()config.shutdown()tunnelClient.shutdown()worker.shutdown() 會關閉通道串連並調用消費回調的 shutdown 方法,config.shutdown() 會關閉讀取、處理和輔助線程池。TunnelWorker 會註冊 JVM 關閉鉤子以嘗試關閉工作器,但應用仍應顯式釋放上述資源。

參數說明

消費工作器

TunnelWorker 的構造方法包含以下參數。

名稱

類型

說明

tunnelId(必選)

String

通道 ID。可以通過建立、列出或查詢通道擷取。

client(必選)

TunnelClientInterface

已初始化的 TunnelClient

workerConfig(必選)

TunnelWorkerConfig

消費回調和消費行為配置。

消費配置

workerConfig 的類型為 TunnelWorkerConfig,包含以下參數。

名稱

類型

說明

channelProcessor(必選)

IChannelProcessor

資料處理回調。使用 TunnelWorker 的三參數構造方法時必須設定。

heartbeatTimeoutInSec(可選)

long

心跳逾時時間,單位為秒。預設值為 300,且必須大於 heartbeatIntervalInSec。發生心跳逾時後,服務端將當前用戶端視為不可用,用戶端會重新串連通道。

heartbeatIntervalInSec(可選)

long

心跳間隔,單位為秒。預設值為 30,最小值為 5。心跳用於擷取活躍 Channel、更新 Channel 狀態和初始化資料處理任務,因此也會影響 TunnelWorker 的預熱時間。

checkpointIntervalInMillis(可選)

long

向服務端記錄消費位點的間隔,單位為毫秒。預設值為 5000。通道服務至少投遞一次資料並保持記錄順序;處理任務重啟後從最近一次位點繼續消費,因此部分資料可能被重複處理。縮短間隔可減少重複處理的資料,但記錄位點過於頻繁會影響輸送量。

clientTag(可選)

String

用戶端自訂標識,用於產生用戶端識別碼 和區分不同的 TunnelWorker。預設值為 Java 系統屬性 os.name

readRecordsExecutor(可選)

ThreadPoolExecutor

拉取資料的線程池。預設線程池的核心線程數為 32、最大線程數為 1000、隊列容量為 16,線程空閑 60 秒後可回收。

processRecordsExecutor(可選)

ThreadPoolExecutor

處理資料的線程池。預設配置與 readRecordsExecutor 相同。自訂線程池時,可根據通道的 Channel 數量配置線程數。

maxChannelParallel(可選)

int

同時拉取和處理資料的最大 Channel 數量,用於限制記憶體使用量。預設值為 -1,表示不限制。Tablestore Java SDK 5.10.0 及以上版本支援該參數。

channelHelperExecutor(可選)

ThreadPoolExecutor

初始化 Channel、調度流水線和處理執行階段錯誤的輔助線程池。未設定時使用緩衝線程池。

maxRetryIntervalInMillis(可選)

int

增量資料拉取的指數退避最大基礎間隔,單位為毫秒。預設值為 2000,最小值為 200。當一批資料不超過 500 條且不超過 900 KB 時,用戶端逐步增加退避間隔,實際間隔會在當前基礎間隔的 75%~125% 範圍內隨機取值。Tablestore Java SDK 5.4.0 及以上版本支援該參數。

readMaxTimesPerRound(可選)

int

單輪流水線最多調用 ReadRecords 的次數。預設值為 1

readMaxBytesPerRound(可選)

int

單輪流水線最多拉取的資料量,單位為位元組。預設值為 4194304,即 4 MiB。達到該值或 readMaxTimesPerRound 後停止本輪拉取。

enableClosingChannelDetect(可選)

boolean

是否即時檢測處於 CLOSING 狀態的 Channel。CLOSING 表示 Channel 正在從一個用戶端遷移到另一個用戶端。Tablestore Java SDK 5.13.13 及以上版本支援該參數;5.17.0 及以上版本的預設值為 true。關閉檢測後,如果 Channel 較多但用戶端資源不足,Channel 遷移可能受阻並導致消費中斷。

在同一台機器上啟動多個 TunnelWorker 時,可以複用一個 TunnelWorkerConfig 以共用讀取和處理線程池。停止所有工作器後,只調用一次 config.shutdown()

回調資料

process 方法接收 ProcessRecordsInput 對象,包含以下欄位。

欄位

類型

說明

records

List<StreamRecord>

當前批次拉取到的記錄列表,通過 getRecords() 擷取。

nextToken

String

下一批資料的分頁憑證,通過 getNextToken() 擷取。TunnelWorker 會自動使用該值繼續拉取資料並記錄消費位點。

traceId

String

當前提取要求的追蹤識別碼,通過 getTraceId() 擷取。

channelId

String

當前批次所屬的 Channel ID,通過 getChannelId() 擷取。可以通過 getPartitionId() 從 Channel ID 中擷取分區 ID。

情境樣本

調整消費配置

消費輸送量或記憶體佔用不符合預期時,可以同時調整心跳和位點間隔、Channel 並行度、單輪拉取次數與資料量,以及增量拉取的退避間隔。

TunnelWorkerConfig config =
        new TunnelWorkerConfig(new SimpleProcessor());
config.setHeartbeatIntervalInSec(10);
config.setHeartbeatTimeoutInSec(60);
config.setCheckpointIntervalInMillis(10_000);
config.setMaxChannelParallel(16);
config.setReadMaxTimesPerRound(4);
config.setReadMaxBytesPerRound(8 * 1024 * 1024);
config.setMaxRetryIntervalInMillis(3_000);