消费通道数据

更新时间:
复制 MD 格式

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

客户端自定义标识,用于生成客户端 ID 和区分不同的 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

当前拉取请求的追踪 ID,通过 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);