A API Stream no Tablestore SDK for Java consome alterações incrementais (inserções, atualizações e exclusões) de uma tabela. Os casos de uso comuns incluem exportação offline, sincronização de dados e notificações de alteração.
Pré-requisitos
Tablestore SDK for Java instalado e cliente inicializado.
Stream ativado na tabela mediante a definição de
StreamSpecificationdurante a criação da tabela. Para mais detalhes, consulte Crie uma tabela de dados.
Como funciona
O Stream organiza as alterações incrementais de uma tabela em shards. Para consumir um stream, chame quatro APIs em sequência: listar streams, descrever um stream, obter um iterador de shard e buscar registros.
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
Etapas principais:
Chame
listStream(ListStreamRequest)para listar os valores destreamIdde todas as tabelas com Stream habilitado na instância.Chame
describeStream(DescribeStreamRequest)para recuperar metadados do stream (hora de criação, hora de expiração, status atual) e a lista de objetosShard.Chame
getShardIterator(GetShardIteratorRequest)para obter o iterador de leitura (shardIterator) de umShardespecífico. O iterador marca onde começar a buscar registros incrementais.Chame
getStreamRecord(GetStreamRecordRequest)com oshardIteratorpara buscar um lote de registros incrementais (uma lista de objetosStreamRecord). Use onextShardIteratorretornado para buscar os registros subsequentes.
O exemplo a seguir consome o stream stream_test_demo de ponta a ponta e imprime o tipo e a chave primária de cada registro.
String demoTable = "stream_test_demo";
// 1. List all tables in the instance with Stream enabled and find the streamId of the target table.
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. Query all shards of the Stream.
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. Get the initial read iterator of the shard.
GetShardIteratorRequest iterRequest =
new GetShardIteratorRequest(targetStreamId, shardId);
GetShardIteratorResponse iterResponse = client.getShardIterator(iterRequest);
String shardIterator = iterResponse.getShardIterator();
// 4. Use the iterator to pull incremental records from the 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 is used to continue pulling subsequent incremental records.
System.out.println("Next iterator: "
+ (recordResponse.getNextShardIterator() != null ? "yes" : "no"));
}
Parâmetros
ListStreamRequest
|
Nome |
Tipo |
Descrição |
|
tableName (opcional) |
String |
Nome da tabela. Se omitido, a solicitação retorna informações de stream para todas as tabelas com Stream habilitado na instância; caso contrário, retorna informações apenas para a tabela especificada. |
DescribeStreamRequest
|
Nome |
Tipo |
Descrição |
|
streamId (obrigatório) |
String |
Identificador exclusivo do stream. Retornado por |
|
inclusiveStartShardId (opcional) |
String |
O |
|
shardLimit (opcional) |
int |
Número máximo de shards a serem retornados na resposta. |
GetShardIteratorRequest
|
Nome |
Tipo |
Descrição |
|
streamId (obrigatório) |
String |
Identificador exclusivo do stream. Obtido por meio de |
|
shardId (obrigatório) |
String |
Identificador exclusivo do shard. Fornecido no objeto |
|
timestamp (opcional) |
long |
Timestamp inicial do iterador, em microssegundos. Se omitido, a leitura começa do início do shard. |
GetStreamRecordRequest
|
Nome |
Tipo |
Descrição |
|
shardIterator (obrigatório) |
String |
Iterador de leitura. Retornado por |
|
limit (opcional) |
int |
Quantidade máxima de objetos |
|
tableName (opcional) |
String |
Nome da tabela que contém o shard de destino. |
Exemplos
Paginar pela lista de shards
Para streams com muitos shards, pagine com inclusiveStartShardId e shardLimit. Um valor null em nextShardId indica que todos os shards foram retornados.
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();
// A null nextShardId indicates that all shards have been traversed.
if (response.getNextShardId() == null) {
break;
}
startShardId = response.getNextShardId();
}
System.out.println("Total shards: " + totalShards);
Consultar continuamente dados incrementais
Chame repetidamente getStreamRecord com nextShardIterator para buscar registros incrementais de um único shard. Quando nextShardIterator for null, significa que o shard atual foi totalmente consumido.
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);