Tablestore Java SDK 可通過通道持續消費資料,使用消費回調處理各批記錄,並配置心跳、Checkpoint、線程池和消費並行度。
注意事項
-
增量日誌的保留時間與資料表的 Stream 日誌到期時間一致,最長為 7 天。使用全量加增量類型的通道時,如果全量資料未能在增量日誌保留時間內消費完成,開始消費增量資料時會返回
OTSTunnelExpired錯誤,無法繼續消費增量資料。 -
消費增量資料的進度落後於增量日誌保留時間時,通道可能從當前仍可用的最新資料開始消費,導致部分資料未被消費。
-
通道到期後可能被禁用。通道連續處于禁用狀態超過 30 天后會被刪除,刪除後無法恢複。
前提條件
功能說明
TunnelWorker 根據通道 ID 串連通道,通過心跳擷取分配給當前用戶端的 Channel,持續拉取資料,並將每批記錄傳入 IChannelProcessor。多個 TunnelWorker 消費同一通道時,服務端會在各用戶端之間分配 Channel。
消費通道資料包括以下步驟:
-
實現
IChannelProcessor介面。process方法處理每批記錄,shutdown方法釋放消費回調使用的資源。 -
建立
TunnelWorkerConfig,配置資料處理回調和消費參數。 -
使用通道 ID、
TunnelClient和TunnelWorkerConfig建立TunnelWorker。 -
調用
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 |
已初始化的 |
|
workerConfig(必選) |
TunnelWorkerConfig |
消費回調和消費行為配置。 |
消費配置
workerConfig 的類型為 TunnelWorkerConfig,包含以下參數。
|
名稱 |
類型 |
說明 |
|
channelProcessor(必選) |
IChannelProcessor |
資料處理回調。使用 |
|
heartbeatTimeoutInSec(可選) |
long |
心跳逾時時間,單位為秒。預設值為 |
|
heartbeatIntervalInSec(可選) |
long |
心跳間隔,單位為秒。預設值為 |
|
checkpointIntervalInMillis(可選) |
long |
向服務端記錄消費位點的間隔,單位為毫秒。預設值為 |
|
clientTag(可選) |
String |
用戶端自訂標識,用於產生用戶端識別碼 和區分不同的 |
|
readRecordsExecutor(可選) |
ThreadPoolExecutor |
拉取資料的線程池。預設線程池的核心線程數為 |
|
processRecordsExecutor(可選) |
ThreadPoolExecutor |
處理資料的線程池。預設配置與 |
|
maxChannelParallel(可選) |
int |
同時拉取和處理資料的最大 Channel 數量,用於限制記憶體使用量。預設值為 |
|
channelHelperExecutor(可選) |
ThreadPoolExecutor |
初始化 Channel、調度流水線和處理執行階段錯誤的輔助線程池。未設定時使用緩衝線程池。 |
|
maxRetryIntervalInMillis(可選) |
int |
增量資料拉取的指數退避最大基礎間隔,單位為毫秒。預設值為 |
|
readMaxTimesPerRound(可選) |
int |
單輪流水線最多調用 |
|
readMaxBytesPerRound(可選) |
int |
單輪流水線最多拉取的資料量,單位為位元組。預設值為 |
|
enableClosingChannelDetect(可選) |
boolean |
是否即時檢測處於 |
在同一台機器上啟動多個 TunnelWorker 時,可以複用一個 TunnelWorkerConfig 以共用讀取和處理線程池。停止所有工作器後,只調用一次 config.shutdown()。
回調資料
process 方法接收 ProcessRecordsInput 對象,包含以下欄位。
|
欄位 |
類型 |
說明 |
|
records |
|
當前批次拉取到的記錄列表,通過 |
|
nextToken |
String |
下一批資料的分頁憑證,通過 |
|
traceId |
String |
當前提取要求的追蹤識別碼,通過 |
|
channelId |
String |
當前批次所屬的 Channel 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);