Tablestore SDK for Java のストリーム API は、挿入、更新、削除など、テーブルに対する増分データを消費します。
前提条件
-
Tablestore SDK for Java をインストールし、クライアントを初期化しておきます。
-
テーブルのストリームは、テーブル作成時に
StreamSpecificationを設定することで有効になります。詳細については、「データテーブルを作成する」をご参照ください。
説明
ストリームは、テーブルからの増分データをシャードに整理します。ストリームを消費するには、ストリームのリスト表示、ストリームの説明、シャードイテレーターの取得、レコードのフェッチという 4 つの API を順番に呼び出します。
-
listStream(ListStreamRequest)を呼び出して、インスタンス内のすべての Stream が有効化されたテーブルのstreamId値を一覧表示します。 -
describeStream(DescribeStreamRequest)を呼び出して、ストリームのメタデータ (作成時刻、有効期限、現在のステータス) とShardオブジェクトのリストを取得します。 -
getShardIterator(GetShardIteratorRequest)を呼び出して、特定のShardの読み取りイテレーター (shardIterator) を取得します。 イテレーターは、増分レコードのフェッチを開始する場所を示します。 -
shardIteratorを使用してgetStreamRecord(GetStreamRecordRequest)を呼び出し、増分レコードのバッチ (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 ストリームをエンドツーエンドで消費し、各レコードのタイプとプライマリキーを出力します。
String demoTable = "stream_test_demo";
// 1. インスタンス内でストリームが有効なすべてのテーブルをリスト表示し、ターゲットテーブルの 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. ストリームのすべてのシャードをクエリします。
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. シャードの初期読み取りイテレーターを取得します。
GetShardIteratorRequest iterRequest =
new GetShardIteratorRequest(targetStreamId, shardId);
GetShardIteratorResponse iterResponse = client.getShardIterator(iterRequest);
String shardIterator = iterResponse.getShardIterator();
// 4. イテレーターを使用して、シャードから増分レコードをプルします。
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"));
}
パラメーター
listStream リクエスト
ListStreamRequest には、次のパラメーターが含まれます。
|
名前 |
タイプ |
説明 |
|
tableName (オプション) |
String |
テーブルの名前。省略した場合、リクエストはインスタンス内のストリームが有効なすべてのテーブルのストリーム情報を返します。指定した場合、指定されたテーブルの情報のみを返します。 |
describeStream リクエスト
DescribeStreamRequest には、以下のパラメータが含まれます。
|
名前 |
タイプ |
説明 |
|
streamId (必須) |
String |
ストリームの一意の識別子で、 |
|
inclusiveStartShardId (オプション) |
String |
返されるシャードリストの開始 |
|
shardLimit (オプション) |
int |
レスポンスで返されるシャードの最大数。 |
getShardIterator リクエスト
GetShardIteratorRequest には、以下のパラメーターが含まれています。
|
名前 |
タイプ |
説明 |
|
streamId (必須) |
String |
ストリームの一意の識別子。 |
|
shardId (必須) |
String |
シャードの一意の識別子です。 |
|
timestamp (オプション) |
long |
イテレーターの開始タイムスタンプ (マイクロ秒単位)。省略した場合、読み取りはシャードの先頭から開始されます。 |
getStreamRecord リクエスト
GetStreamRecordRequest には、以下のパラメーターが含まれます。
|
名前 |
タイプ |
説明 |
|
shardIterator (必須) |
String |
|
|
limit (オプション) |
int |
レスポンスで返される |
|
tableName (オプション) |
String |
ターゲットシャードを含むテーブルの名前。 |
レスポンス
listStream レスポンス
ListStreamResponse には、以下の操作固有のフィールドが含まれています。
|
名前 |
タイプ |
説明 |
|
streams |
List<Stream> |
The Stream information list. Each element includes information such as the table name, Stream ID, and expiration time. Call |
describeStream レスポンス
DescribeStreamResponse には、以下の操作固有のフィールドが含まれています。
|
名前 |
タイプ |
説明 |
|
streamId |
String |
ストリーム ID。 |
|
tableName |
String |
データテーブルの名前。 |
|
creationTime |
long |
ストリームが作成された時刻。 |
|
expirationTime |
int |
ストリームの有効期限。 |
|
status |
StreamStatus |
ストリームのステータス。 |
|
shards |
List<StreamShard> |
現在のページで返されたシャード。 |
|
nextShardId |
String |
次のページの開始シャード ID です。 |
|
timeseriesDataTable |
boolean |
テーブルが時系列データテーブルであるかどうかを示します。 |
getShardIterator レスポンス
GetShardIteratorResponse は、以下の操作固有のフィールドを含みます。
|
名前 |
タイプ |
説明 |
|
shardIterator |
String |
指定されたシャードのイテレーター。最初の |
getStreamRecord レスポンス
GetStreamRecordResponse には、以下の操作固有のフィールドが含まれています。
|
名前 |
タイプ |
説明 |
|
records |
List<StreamRecord> |
レスポンスで返された増分レコード。 |
|
nextShardIterator |
String |
次の読み取りに使用するイテレーターです。 |
|
mayMoreRecord |
Boolean |
現在のシャードにさらにレコードが含まれている可能性があるかどうかを示します。 |
例
シャードリストのページング
多数のシャードを持つストリームの場合は、inclusiveStartShardId と shardLimit を使用してページ分割します。null の nextShardId は、すべてのシャードが返されたことを示します。
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 の場合、すべてのシャードが走査されたことを示します。
if (response.getNextShardId() == null) {
break;
}
startShardId = response.getNextShardId();
}
System.out.println("Total shards: " + totalShards);
増分データの継続的なポーリング
nextShardIterator を使用して getStreamRecord を繰り返し呼び出し、単一のシャードから増分レコードを取得します。nextShardIterator が null の場合は、現在のシャードが完全に消費されたことを示します。
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);