消费增量数据

更新时间:
复制 MD 格式

Java SDK 通过 Stream API 消费数据表的插入、更新和删除等增量变更数据。

前提条件

  • 安装 Tablestore Java SDK并初始化客户端。

  • 数据表已开启 Stream 功能。创建数据表时通过 StreamSpecification 启用,具体操作请参见创建数据表

功能说明

Stream 将数据表的增量变更按 Shard 组织,消费按照 list → describe → getIterator → 循环 getStreamRecord 的步骤进行:

  1. listStream(ListStreamRequest) 列出实例下已开启 Stream 的所有表的 streamId

  2. describeStream(DescribeStreamRequest) 查询 Stream 的描述信息(创建时间、过期时间、当前状态)以及包含的 Shard 列表。

  3. getShardIterator(GetShardIteratorRequest) 获取指定 Shard 的读取迭代值(shardIterator),作为后续拉取增量数据的起点。

  4. getStreamRecord(GetStreamRecordRequest) 通过 shardIterator 拉取一批增量记录(StreamRecord 列表),并返回 nextShardIterator 用于继续拉取。

public ListStreamResponse listStream(ListStreamRequest request) throws TableStoreException, ClientException
public DescribeStreamResponse describeStream(DescribeStreamRequest request) throws TableStoreException, ClientException
public GetShardIteratorResponse getShardIterator(GetShardIteratorRequest request) throws TableStoreException, ClientException
public GetStreamRecordResponse getStreamRecord(GetStreamRecordRequest request) throws TableStoreException, ClientException

以下示例端到端消费数据表 stream_test_demo 的 Stream,输出每条增量记录的类型和主键。

String demoTable = "stream_test_demo";

// 1. 列出实例下开启 Stream 的所有表,找到目标表的 streamId
ListStreamRequest listRequest = new ListStreamRequest(demoTable);
ListStreamResponse listResponse = client.listStream(listRequest);

String targetStreamId = null;
for (Stream stream : listResponse.getStreams()) {
    if (demoTable.equals(stream.getTableName())) {
        targetStreamId = stream.getStreamId();
        break;
    }
}
System.out.println("Stream ID: " + targetStreamId);

// 2. 查询 Stream 的所有 Shard
DescribeStreamRequest describeRequest = new DescribeStreamRequest(targetStreamId);
DescribeStreamResponse describeResponse = client.describeStream(describeRequest);
List<StreamShard> shards = describeResponse.getShards();
System.out.println("Shard count: " + shards.size());

if (!shards.isEmpty()) {
    String shardId = shards.get(0).getShardId();

    // 3. 拿 Shard 的初始读取迭代值
    GetShardIteratorRequest iterRequest =
            new GetShardIteratorRequest(targetStreamId, shardId);
    GetShardIteratorResponse iterResponse = client.getShardIterator(iterRequest);
    String shardIterator = iterResponse.getShardIterator();

    // 4. 用迭代值拉取 Shard 的增量记录
    GetStreamRecordRequest recordRequest = new GetStreamRecordRequest(shardIterator);
    recordRequest.setLimit(100);
    GetStreamRecordResponse recordResponse = client.getStreamRecord(recordRequest);

    List<StreamRecord> records = recordResponse.getRecords();
    System.out.println("Records fetched: " + records.size());
    for (StreamRecord record : records) {
        System.out.println("RecordType: " + record.getRecordType()
                + ", PK: " + record.getPrimaryKey());
    }

    // nextShardIterator 用于继续拉取后续增量
    System.out.println("Next iterator: "
            + (recordResponse.getNextShardIterator() != null ? "yes" : "no"));
}

参数说明

列出 Stream 请求

ListStreamRequest 包含以下参数。

名称

类型

说明

tableName(可选)

String

数据表名称。不指定时返回当前实例下所有开启 Stream 的表的 Stream 信息;指定时仅返回该表的 Stream 信息。

查询 Stream 请求

DescribeStreamRequest 包含以下参数。

名称

类型

说明

streamId(必选)

String

Stream 的唯一标识,由 listStream 返回。

inclusiveStartShardId(可选)

String

返回的 Shard 列表的起始 shardId,用于分页拿取大量 Shard。

shardLimit(可选)

int

本次返回的 Shard 数量上限。

获取 Shard 迭代值请求

GetShardIteratorRequest 包含以下参数。

名称

类型

说明

streamId(必选)

String

Stream 的唯一标识,由 describeStream 返回。

shardId(必选)

String

Shard 的唯一标识,由 describeStream 返回的 StreamShard 中获取。

timestamp(可选)

long

指定迭代起点的时间戳(微秒),用于从指定时间开始读取。不指定时从 Shard 起始位置读取。

读取增量数据请求

GetStreamRecordRequest 包含以下参数。

名称

类型

说明

shardIterator(必选)

String

读取迭代值,由 getShardIterator 或上一次 getStreamRecordnextShardIterator 返回。

limit(可选)

int

本次返回的 StreamRecord 数量上限。

tableName(可选)

String

目标 Shard 所属的数据表名称。

返回值

Stream 列表

ListStreamResponse 包含以下业务字段。

名称

类型

说明

streams

List<Stream>

Stream 信息列表。每个元素包含数据表名称、Stream ID 和过期时间等信息。通过 getStreams() 获取。

Stream 信息

DescribeStreamResponse 包含以下业务字段。

名称

类型

说明

streamId

String

Stream ID。

tableName

String

数据表名称。

creationTime

long

Stream 创建时间。

expirationTime

int

Stream 过期时间。

status

StreamStatus

Stream 状态。

shards

List<StreamShard>

本次返回的 Shard 列表。

nextShardId

String

下一页的起始 Shard ID。为 null 时表示 Shard 已全部返回。

timeseriesDataTable

boolean

数据表是否为时序数据表。通过 isTimeseriesDataTable() 获取。

Shard 迭代值

GetShardIteratorResponse 包含以下业务字段。

名称

类型

说明

shardIterator

String

指定 Shard 的读取迭代值,用于首次调用 getStreamRecord()

增量数据

GetStreamRecordResponse 包含以下业务字段。

名称

类型

说明

records

List<StreamRecord>

本次返回的增量记录列表。

nextShardIterator

String

下一次读取使用的迭代值。为 null 时表示当前 Shard 已读完。

mayMoreRecord

Boolean

当前 Shard 是否可能还有更多记录。

场景示例

分页获取 Shard 列表

Stream 包含的 Shard 数量较多时,通过 inclusiveStartShardIdshardLimit 分批获取。nextShardIdnull 表示已遍历完所有 Shard。

String currentStreamId = "<your-stream-id>";
String startShardId = null;
int totalShards = 0;

while (true) {
    DescribeStreamRequest request = new DescribeStreamRequest(currentStreamId);
    if (startShardId != null) {
        request.setInclusiveStartShardId(startShardId);
    }
    request.setShardLimit(50);

    DescribeStreamResponse response = client.describeStream(request);
    totalShards += response.getShards().size();

    // nextShardId 为 null 表示已遍历完所有 Shard
    if (response.getNextShardId() == null) {
        break;
    }
    startShardId = response.getNextShardId();
}
System.out.println("Total shards: " + totalShards);

持续轮询增量数据

nextShardIterator 循环调用 getStreamRecord 持续拉取一个 Shard 的增量数据。nextShardIteratornull 表示当前 Shard 已读完。

String currentStreamId = "<your-stream-id>";
String shardId = "<your-shard-id>";

GetShardIteratorRequest iterRequest =
        new GetShardIteratorRequest(currentStreamId, shardId);
String shardIterator = client.getShardIterator(iterRequest).getShardIterator();

int totalRecords = 0;
while (shardIterator != null) {
    GetStreamRecordRequest recordRequest = new GetStreamRecordRequest(shardIterator);
    recordRequest.setLimit(100);
    GetStreamRecordResponse response = client.getStreamRecord(recordRequest);

    totalRecords += response.getRecords().size();
    shardIterator = response.getNextShardIterator();
}
System.out.println("Polling total records: " + totalRecords);