API Stream dalam Tablestore SDK for Java mengonsumsi perubahan inkremental pada suatu tabel, termasuk operasi insert, update, dan delete.
Prasyarat
-
Instal Tablestore SDK for Java dan inisialisasi klien.
-
Aktifkan Stream pada tabel dengan menetapkan
StreamSpecificationsaat pembuatan tabel. Untuk detailnya, lihat Buat tabel data.
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.
-
Panggil
listStream(ListStreamRequest)untuk mendaftar nilaistreamIddari semua tabel yang diaktifkan Stream di instans tersebut. -
Panggil
describeStream(DescribeStreamRequest)untuk mengambil metadata stream (waktu pembuatan, waktu kedaluwarsa, status saat ini) dan daftar objekShard. -
Panggil
getShardIterator(GetShardIteratorRequest)untuk mendapatkan iterator baca (shardIterator) untukShardtertentu. Iterator ini menandai posisi awal pengambilan catatan inkremental. -
Panggil
getStreamRecord(GetStreamRecordRequest)denganshardIteratoruntuk mengambil batch catatan inkremental (daftar objekStreamRecord). GunakannextShardIteratoryang 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 |
|
inclusiveStartShardId (opsional) |
String |
|
|
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 |
|
shardId (wajib) |
String |
Pengidentifikasi unik shard. Dikembalikan dalam objek |
|
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 |
|
limit (opsional) |
int |
Jumlah maksimum objek |
|
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 |
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 |
|
timeseriesDataTable |
boolean |
Menunjukkan apakah tabel tersebut merupakan tabel data deret waktu. Panggil |
Shard iterator
GetShardIteratorResponse berisi field spesifik operasi berikut.
|
Nama |
Type |
Deskripsi |
|
shardIterator |
String |
Iterator untuk shard yang ditentukan. Gunakan nilai ini dalam permintaan pertama |
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 |
|
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);