All Products
Search
Document Center

Tablestore:Mengonsumsi data inkremental

Last Updated:Aug 06, 2026

API Stream dalam Tablestore SDK for Java mengonsumsi perubahan inkremental pada suatu tabel, termasuk operasi insert, update, dan delete.

Prasyarat

Deskripsi

Stream mengorganisasi perubahan inkremental dari suatu tabel ke dalam shard. Untuk mengonsumsi stream, panggil empat API secara berurutan: list streams, describe a stream, get a shard iterator, dan fetch records.

  1. Panggil listStream(ListStreamRequest) untuk mendaftar nilai streamId dari semua tabel yang diaktifkan Stream di instans tersebut.

  2. Panggil describeStream(DescribeStreamRequest) untuk mengambil metadata stream (waktu pembuatan, waktu kedaluwarsa, status saat ini) dan daftar objek Shard.

  3. Panggil getShardIterator(GetShardIteratorRequest) untuk mendapatkan iterator baca (shardIterator) untuk Shard tertentu. Iterator ini menandai posisi awal pengambilan catatan inkremental.

  4. Panggil getStreamRecord(GetStreamRecordRequest) dengan shardIterator untuk mengambil batch catatan inkremental (daftar objek StreamRecord). Gunakan nextShardIterator yang dikembalikan untuk mengambil catatan berikutnya.

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

Contoh berikut mengonsumsi stream stream_test_demo secara end-to-end dan mencetak tipe serta kunci primer setiap catatan.

String demoTable = "stream_test_demo";

// 1. Daftar semua tabel di instans yang diaktifkan Stream dan temukan streamId tabel target.
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. Kueri semua shard dari Stream.
DescribeStreamRequest describeRequest = new DescribeStreamRequest(targetStreamId);
DescribeStreamResponse describeResponse = client.describeStream(describeRequest);
List<StreamShard> shards = describeResponse.getShards();
System.out.println("Jumlah shard: " + shards.size());

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

    // 3. Dapatkan iterator baca awal dari shard.
    GetShardIteratorRequest iterRequest =
            new GetShardIteratorRequest(targetStreamId, shardId);
    GetShardIteratorResponse iterResponse = client.getShardIterator(iterRequest);
    String shardIterator = iterResponse.getShardIterator();

    // 4. Gunakan iterator untuk menarik catatan inkremental dari shard.
    GetStreamRecordRequest recordRequest = new GetStreamRecordRequest(shardIterator);
    recordRequest.setLimit(100);
    GetStreamRecordResponse recordResponse = client.getStreamRecord(recordRequest);

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

    // nextShardIterator digunakan untuk melanjutkan penarikan catatan inkremental berikutnya.
    System.out.println("Iterator berikutnya: "
            + (recordResponse.getNextShardIterator() != null ? "ya" : "tidak"));
}

Parameter

Permintaan daftar stream

ListStreamRequest berisi parameter berikut.

Nama

Tipe

Deskripsi

tableName (opsional)

String

Nama tabel. Jika dihilangkan, permintaan akan mengembalikan informasi stream untuk setiap tabel yang diaktifkan Stream di instans tersebut; jika tidak, hanya mengembalikan informasi untuk tabel yang ditentukan.

Permintaan deskripsi stream

DescribeStreamRequest berisi parameter berikut.

Nama

Type

Deskripsi

streamId (wajib)

String

Pengidentifikasi unik stream. Dikembalikan oleh listStream.

inclusiveStartShardId (opsional)

String

shardId awal dari daftar shard yang dikembalikan. Tentukan untuk melakukan paginasi melalui kumpulan shard yang besar.

shardLimit (opsional)

int

Jumlah maksimum shard yang dikembalikan dalam respons.

Permintaan mendapatkan iterator shard

GetShardIteratorRequest berisi parameter berikut.

Nama

Tipe

Deskripsi

streamId (wajib)

String

Pengidentifikasi unik stream. Dikembalikan oleh describeStream.

shardId (wajib)

String

Pengidentifikasi unik shard. Dikembalikan dalam objek StreamShard oleh describeStream.

timestamp (opsional)

long

Timestamp awal iterator, dalam mikrodetik. Jika dihilangkan, pembacaan dimulai dari awal shard.

Permintaan membaca data inkremental

GetStreamRecordRequest berisi parameter berikut.

Nama

Tipe

Deskripsi

shardIterator (wajib)

String

Iterator baca. Dikembalikan oleh getShardIterator atau oleh field nextShardIterator dari respons getStreamRecord sebelumnya.

limit (opsional)

int

Jumlah maksimum objek StreamRecord yang dikembalikan dalam respons.

tableName (opsional)

String

Nama tabel yang berisi shard target.

Respons

Daftar stream

ListStreamResponse berisi field spesifik operasi berikut.

Nama

Type

Deskripsi

streams

List<Stream>

Daftar informasi Stream. Setiap elemen mencakup informasi seperti nama tabel, ID Stream, dan waktu kedaluwarsa. Panggil getStreams() untuk mendapatkan daftar tersebut.

Informasi stream

DescribeStreamResponse berisi field spesifik operasi berikut.

Nama

Tipe

Deskripsi

streamId

String

Stream ID.

tableName

String

Nama tabel data.

creationTime

long

Waktu pembuatan Stream.

expirationTime

int

Waktu kedaluwarsa Stream.

status

StreamStatus

Status Stream.

shards

List<StreamShard>

Shard yang dikembalikan pada halaman saat ini.

nextShardId

String

ID shard awal untuk halaman berikutnya. Nilai null menunjukkan bahwa semua shard telah dikembalikan.

timeseriesDataTable

boolean

Menunjukkan apakah tabel tersebut merupakan tabel data deret waktu. Panggil isTimeseriesDataTable() untuk mendapatkan nilainya.

Shard iterator

GetShardIteratorResponse berisi field spesifik operasi berikut.

Nama

Type

Deskripsi

shardIterator

String

Iterator untuk shard yang ditentukan. Gunakan nilai ini dalam permintaan pertama getStreamRecord().

Data inkremental

GetStreamRecordResponse berisi field spesifik operasi berikut.

Nama

Tipe

Deskripsi

records

List<StreamRecord>

Catatan inkremental yang dikembalikan dalam respons.

nextShardIterator

String

Iterator untuk pembacaan berikutnya. Nilai null menunjukkan bahwa shard saat ini telah sepenuhnya dibaca.

mayMoreRecord

Boolean

Menunjukkan apakah shard saat ini mungkin masih berisi catatan tambahan.

Contoh

Lakukan paginasi melalui daftar shard

Untuk stream dengan banyak shard, lakukan paginasi menggunakan inclusiveStartShardId dan shardLimit. Nilai nextShardId yang null menunjukkan bahwa semua shard telah dikembalikan.

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

    // Nilai nextShardId null menunjukkan bahwa semua shard telah dilalui.
    if (response.getNextShardId() == null) {
        break;
    }
    startShardId = response.getNextShardId();
}
System.out.println("Total shard: " + totalShards);

Lakukan polling data inkremental secara terus-menerus

Panggil berulang kali getStreamRecord dengan nextShardIterator untuk mengambil catatan inkremental dari satu shard. Nilai nextShardIterator yang null menunjukkan bahwa shard saat ini telah sepenuhnya dikonsumsi.

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("Total catatan dari polling: " + totalRecords);