Java SDK 通過 Stream API 消費資料表的插入、更新和刪除等增量變更資料。
前提條件
-
安裝 Tablestore Java SDK並初始化用戶端。
-
資料表已開啟 Stream 功能。建立資料表時通過
StreamSpecification啟用,具體操作請參見建立資料表。
功能說明
Stream 將資料表的增量變更按 Shard 組織,消費按照 list → describe → getIterator → 迴圈 getStreamRecord 的步驟進行:
-
listStream(ListStreamRequest)列出執行個體下已開啟 Stream 的所有表的streamId。 -
describeStream(DescribeStreamRequest)查詢 Stream 的描述資訊(建立時間、到期時間、目前狀態)以及包含的Shard列表。 -
getShardIterator(GetShardIteratorRequest)擷取指定Shard的讀取迭代值(shardIterator),作為後續拉取增量資料的起點。 -
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 的唯一標識,由 |
|
inclusiveStartShardId(可選) |
String |
返回的 Shard 列表的起始 |
|
shardLimit(可選) |
int |
本次返回的 Shard 數量上限。 |
擷取 Shard 迭代值請求
GetShardIteratorRequest 包含以下參數。
|
名稱 |
類型 |
說明 |
|
streamId(必選) |
String |
Stream 的唯一標識,由 |
|
shardId(必選) |
String |
Shard 的唯一標識,由 |
|
timestamp(可選) |
long |
指定迭代起點的時間戳記(微秒),用於從指定時間開始讀取。不指定時從 Shard 起始位置讀取。 |
讀取增量資料請求
GetStreamRecordRequest 包含以下參數。
|
名稱 |
類型 |
說明 |
|
shardIterator(必選) |
String |
讀取迭代值,由 |
|
limit(可選) |
int |
本次返回的 |
|
tableName(可選) |
String |
目標 Shard 所屬的資料表名稱。 |
傳回值
Stream 列表
ListStreamResponse 包含以下業務欄位。
|
名稱 |
類型 |
說明 |
|
streams |
List<Stream> |
Stream 資訊列表。每個元素包含資料表名稱、Stream ID 和到期時間等資訊。通過 |
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。為 |
|
timeseriesDataTable |
boolean |
資料表是否為時序資料表。通過 |
Shard 迭代值
GetShardIteratorResponse 包含以下業務欄位。
|
名稱 |
類型 |
說明 |
|
shardIterator |
String |
指定 Shard 的讀取迭代值,用於首次調用 |
增量資料
GetStreamRecordResponse 包含以下業務欄位。
|
名稱 |
類型 |
說明 |
|
records |
List<StreamRecord> |
本次返回的增量記錄列表。 |
|
nextShardIterator |
String |
下一次讀取使用的迭代值。為 |
|
mayMoreRecord |
Boolean |
當前 Shard 是否可能還有更多記錄。 |
情境樣本
分頁擷取 Shard 列表
Stream 包含的 Shard 數量較多時,通過 inclusiveStartShardId 和 shardLimit 分批擷取。nextShardId 為 null 表示已遍曆完所有 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 的增量資料。nextShardIterator 為 null 表示當前 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);