Todos os produtos
Search
Central de documentação

Tablestore:Use TableStoreReader to concurrently read data

Última atualização: Jun 30, 2026

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.

Nota

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

TableStoreReaderConfig parameters

Parâmetro

Tipo

Descrição

checkTableMeta

boolean

Define se a verificação de schema deve ser ativada. O valor padrão é true. Quando ativado, o TableStoreReader executa as seguintes verificações antes de gravar os dados no buffer:

  • Existência da tabela de dados.

  • Correspondência entre o schema da chave primária dos dados consultados e a chave primária da tabela.

bufferSize

int

Tamanho da fila de buffer. Deve ser uma potência de 2. O valor padrão é 1024.

concurrency

int

Número máximo de requisições concorrentes ao enviar dados do buffer para o Tablestore. O valor padrão é 10.

maxBatchRowsCount

int

Quantidade máxima de linhas lidas em uma requisição em lote. O valor padrão e o limite máximo são 100.

defaultMaxVersions

int

Número de versões de dados a serem lidas. O valor padrão é 1, indicando que apenas a versão mais recente será recuperada.

flushInterval

int

Intervalo de tempo para envio automático dos dados do buffer ao Tablestore. O valor padrão é 10000 milissegundos.

logInterval

int

Intervalo de tempo para impressão do status da tarefa durante o envio de dados do buffer ao Tablestore. O valor padrão é 10000 milissegundos.

bucketCount

int

Quantidade de buckets. Cada bucket funciona como um buffer independente. O valor padrão é 4.

O código de exemplo a seguir demonstra como configurar esses parâmetros.

// Specify whether to enable schema checking.
config.setCheckTableMeta(false);
// Specify the buffer queue size.
config.setBufferSize(1024);
// Specify the concurrency.
config.setConcurrency(10);
// Specify the maximum number of rows for batch requests.
config.setMaxBatchRowsCount(100);
// Specify the number of data versions to read.
config.setDefaultMaxVersions(1);
// Specify the flush interval.
config.setFlushInterval(10000);
// Specify the log interval.
config.setLogInterval(10000);
// Specify the number of buckets.
config.setBucketCount(4);
  • 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

  1. 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);
  2. 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);
    }
}