O TableStoreReader é uma classe utilitária simples e de alto desempenho para leitura de dados, implementada pelo Tablestore com base no SDK para Java. Ele encapsula interfaces para leitura de dados com alta concorrência e alto throughput, permitindo leitura simultânea com suporte a callbacks no nível de linha e configuração personalizada. Este tópico descreve como usar o TableStoreReader para leituras concorrentes.
Pré-requisitos
Crie um par de AccessKey para sua conta Alibaba Cloud ou para um usuário RAM com permissões de acesso ao Tablestore.
Procedimento
Etapa 1: Instalar o SDK do Tablestore para Java
Se você usa o Maven para gerenciar projetos Java, adicione a seguinte dependência ao arquivo pom.xml:
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>tablestore</artifactId>
<version>5.17.4</version>
</dependency>
Para mais informações, consulte Instalar o SDK do Tablestore para Java.
Etapa 2: Inicialização
Antes de inicializar o TableStoreReader, crie uma instância de conexão do cliente Tablestore. Personalize os parâmetros e as funções de callback do TableStoreReader conforme necessário. O código de exemplo de inicialização está abaixo.
Ao usar múltiplas threads, recomendamos que elas compartilhem um único objeto TableStoreReader.
public static TableStoreReader createReader() {
// Initialize a Tablestore client.
client = new AsyncClient(endpoint, accessKeyId, accessKeySecret, instanceName);
// TableStoreReader configuration
TableStoreReaderConfig config = new TableStoreReaderConfig();
// Thread pool
executorService = new ThreadPoolExecutor(4, 4, 0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>(1024), new ThreadPoolExecutor.CallerRunsPolicy());
// Callback function
TableStoreCallback<PrimaryKeyWithTable, RowReadResult> callback = new TableStoreCallback<PrimaryKeyWithTable, RowReadResult>() {
@Override
public void onCompleted(PrimaryKeyWithTable primaryKeyWithTable, RowReadResult rowReadResult) {
succeedRows.incrementAndGet();
System.out.println(rowReadResult.getRowResult());
}
@Override
public void onFailed(PrimaryKeyWithTable primaryKeyWithTable, Exception e) {
failedRows.incrementAndGet();
System.out.println("Failed Rows: " + primaryKeyWithTable.getTableName() + " | " + primaryKeyWithTable.getPrimaryKey() + " | " + e.getMessage());
}
};
return new DefaultTableStoreReader(client, config, executorService, callback);
}
-
Se não precisar de uma função de callback, defina o parâmetro correspondente como null no método de inicialização.
return new DefaultTableStoreReader(client, config, executorService, null);
Etapa 3: Consultar dados
-
Antes de usar o TableStoreReader para consultar dados, adicione as informações de chave primária das linhas ao buffer.
PrimaryKey primaryKey = PrimaryKeyBuilder.createPrimaryKeyBuilder() .addPrimaryKeyColumn("id", PrimaryKeyValue.fromString("row1")) .build(); tableStoreReader.addPrimaryKey("test_table", primaryKey);-
Para recuperar os dados da linha após a consulta, use o método
addPrimaryKeyWithFuture.Future<ReaderResult> readerResult = tableStoreReader.addPrimaryKeyWithFuture("test_table", primaryKey); -
Defina também parâmetros de consulta, como número máximo de versões, intervalo de versões de dados e filtros.
RowQueryCriteria rowQueryCriteria = new RowQueryCriteria("test_version"); // Specify the maximum number of versions to read. rowQueryCriteria.setMaxVersions(1); // Specify the data version range to read. rowQueryCriteria.setTimeRange(new TimeRange(System.currentTimeMillis() - 86400*1000, System.currentTimeMillis())); // Specify the attribute columns to return. rowQueryCriteria.addColumnsToGet("col1"); // Specify filter conditions. SingleColumnValueFilter singleColumnValueFilter = new SingleColumnValueFilter("col1", SingleColumnValueFilter.CompareOperator.EQUAL, ColumnValue.fromString("val1")); rowQueryCriteria.setFilter(singleColumnValueFilter); // Add query criteria. tableStoreReader.setRowQueryCriteria(rowQueryCriteria);
-
-
Após adicionar as informações de chave primária ao buffer, o TableStoreReader envia automaticamente os dados para o Tablestore conforme o intervalo especificado (padrão: 10 segundos). Envie os dados do buffer manualmente, se preferir.
-
Transmissão síncrona
tableStoreReader.flush(); -
Transmissão assíncrona
tableStoreReader.send();
-
Etapa 4: Liberar recursos
Após concluir a consulta de dados, libere os recursos se nenhuma outra operação for necessária. Isso evita impactos no sistema de negócios.
tableStoreReader.close();
client.shutdown();
executorService.shutdown();
Código de exemplo completo
O exemplo abaixo usa o TableStoreReader para consultar concorrentemente 200 linhas de dados na tabela test_table e imprime os resultados por meio da função de callback.
public class TableStoreReaderExample {
private static final String endpoint = "https://n01k********.cn-hangzhou.ots.aliyuncs.com";
private static final String instanceName = "n01k********";
private static final String accessKeyId = System.getenv("TABLESTORE_ACCESS_KEY_ID");
private static final String accessKeySecret = System.getenv("TABLESTORE_ACCESS_KEY_SECRET");
private static AsyncClientInterface client;
private static ExecutorService executorService;
private static AtomicLong succeedRows = new AtomicLong();
private static AtomicLong failedRows = new AtomicLong();
public static void main(String[] args) throws InterruptedException, ExecutionException {
// Create TableStoreReader.
TableStoreReader tableStoreReader = createReader();
// Add primary keys for querying data.
for(int i=0; i<200; i++) {
PrimaryKey primaryKey = PrimaryKeyBuilder.createPrimaryKeyBuilder()
.addPrimaryKeyColumn("id", PrimaryKeyValue.fromString("row" + i))
.build();
tableStoreReader.addPrimaryKey("test_table", primaryKey);
}
// Send data in the buffer.
tableStoreReader.flush();
// Wait for callback function to complete.
Thread.sleep(1000L);
System.out.println("Succeed Rows Count: " + succeedRows.get());
System.out.println("Failed Rows Count: " + failedRows.get());
// Disable resources.
tableStoreReader.close();
client.shutdown();
executorService.shutdown();
}
public static TableStoreReader createReader() {
// Initialize a Tablestore client.
client = new AsyncClient(endpoint, accessKeyId, accessKeySecret, instanceName);
// TableStoreReader parameter configuration
TableStoreReaderConfig config = new TableStoreReaderConfig();
// Thread pool
executorService = new ThreadPoolExecutor(4, 4, 0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>(1024), new ThreadPoolExecutor.CallerRunsPolicy());
// Callback function
TableStoreCallback<PrimaryKeyWithTable, RowReadResult> callback = new TableStoreCallback<PrimaryKeyWithTable, RowReadResult>() {
@Override
public void onCompleted(PrimaryKeyWithTable primaryKeyWithTable, RowReadResult rowReadResult) {
succeedRows.incrementAndGet();
System.out.println(rowReadResult.getRowResult());
}
@Override
public void onFailed(PrimaryKeyWithTable primaryKeyWithTable, Exception e) {
failedRows.incrementAndGet();
System.out.println("Failed Rows: " + primaryKeyWithTable.getTableName() + " | " + primaryKeyWithTable.getPrimaryKey() + " | " + e.getMessage());
}
};
return new DefaultTableStoreReader(client, config, executorService, callback);
}
}