全部產品
Search
文件中心

Tablestore:消費增量資料

更新時間:Aug 06, 2026

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);