すべてのプロダクト
Search
ドキュメントセンター

Tablestore:増分データの消費

最終更新日:Aug 06, 2026

Tablestore SDK for Java のストリーム API は、挿入、更新、削除など、テーブルに対する増分データを消費します。

前提条件

  • Tablestore SDK for Java をインストールし、クライアントを初期化しておきます。

  • テーブルのストリームは、テーブル作成時に StreamSpecification を設定することで有効になります。詳細については、「データテーブルを作成する」をご参照ください。

説明

ストリームは、テーブルからの増分データをシャードに整理します。ストリームを消費するには、ストリームのリスト表示、ストリームの説明、シャードイテレーターの取得、レコードのフェッチという 4 つの API を順番に呼び出します。

  1. listStream(ListStreamRequest) を呼び出して、インスタンス内のすべての Stream が有効化されたテーブルの streamId 値を一覧表示します。

  2. describeStream(DescribeStreamRequest) を呼び出して、ストリームのメタデータ (作成時刻、有効期限、現在のステータス) とShard オブジェクトのリストを取得します。

  3. getShardIterator(GetShardIteratorRequest) を呼び出して、特定の Shard の読み取りイテレーター (shardIterator) を取得します。 イテレーターは、増分レコードのフェッチを開始する場所を示します。

  4. 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

ストリームの一意の識別子で、 listStream によって返されます。

inclusiveStartShardId (オプション)

String

返されるシャードリストの開始 shardId を示します。大量のシャードセットをページネーションで処理する場合に指定します。

shardLimit (オプション)

int

レスポンスで返されるシャードの最大数。

getShardIterator リクエスト

GetShardIteratorRequest には、以下のパラメーターが含まれています。

名前

タイプ

説明

streamId (必須)

String

ストリームの一意の識別子。describeStream によって返されます。

shardId (必須)

String

シャードの一意の識別子です。 describeStream によって StreamShard オブジェクトで返されます。

timestamp (オプション)

long

イテレーターの開始タイムスタンプ (マイクロ秒単位)。省略した場合、読み取りはシャードの先頭から開始されます。

getStreamRecord リクエスト

GetStreamRecordRequest には、以下のパラメーターが含まれます。

名前

タイプ

説明

shardIterator (必須)

String

getShardIterator、または直前の getStreamRecord レスポンスの nextShardIterator フィールドによって返される読み取りイテレータです。

limit (オプション)

int

レスポンスで返される StreamRecord オブジェクトの最大数。

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 getStreams() to obtain the list.

describeStream レスポンス

DescribeStreamResponse には、以下の操作固有のフィールドが含まれています。

名前

タイプ

説明

streamId

String

ストリーム ID。

tableName

String

データテーブルの名前。

creationTime

long

ストリームが作成された時刻。

expirationTime

int

ストリームの有効期限。

status

StreamStatus

ストリームのステータス。

shards

List<StreamShard>

現在のページで返されたシャード。

nextShardId

String

次のページの開始シャード ID です。null 値は、すべてのシャードが返されたことを示します。

timeseriesDataTable

boolean

テーブルが時系列データテーブルであるかどうかを示します。isTimeseriesDataTable() を呼び出して値を取得します。

getShardIterator レスポンス

GetShardIteratorResponse は、以下の操作固有のフィールドを含みます。

名前

タイプ

説明

shardIterator

String

指定されたシャードのイテレーター。最初の getStreamRecord() リクエストでこの値を使用します。

getStreamRecord レスポンス

GetStreamRecordResponse には、以下の操作固有のフィールドが含まれています。

名前

タイプ

説明

records

List<StreamRecord>

レスポンスで返された増分レコード。

nextShardIterator

String

次の読み取りに使用するイテレーターです。null 値は、現在のシャードが完全に読み取られたことを示します。

mayMoreRecord

Boolean

現在のシャードにさらにレコードが含まれている可能性があるかどうかを示します。

シャードリストのページング

多数のシャードを持つストリームの場合は、inclusiveStartShardIdshardLimit を使用してページ分割します。nullnextShardId は、すべてのシャードが返されたことを示します。

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 を繰り返し呼び出し、単一のシャードから増分レコードを取得します。nextShardIteratornull の場合は、現在のシャードが完全に消費されたことを示します。

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