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 |
客户端自定义标识,用于生成客户端 ID 和区分不同的 |
|
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 |
当前拉取请求的追踪 ID,通过 |
|
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);