Use o Tablestore SDK for Java para varrer linhas correspondentes em um search index de forma concorrente e exportar o conjunto completo de resultados quando a ordenação não for necessária.
Pré-requisitos
Instale o Tablestore SDK for Java e inicialize um cliente.
O parallel scan exige o Tablestore SDK for Java 5.6.0 ou posterior. Para informações sobre versões, consulte Version history.
Funcionamento
O parallel scan percorre todas as linhas correspondentes a uma consulta em um search index. Essa operação não garante a ordem global dos resultados nem suporta classificação ou agregação. Caso precise de resultados ordenados, agregações ou dados de busca para usuários finais, utilize a operação Search.
A varredura com único worker é mais simples de configurar. A varredura com múltiplos workers lê várias divisões simultaneamente e geralmente oferece maior throughput do que a abordagem com apenas um worker.
Uma varredura paralela consiste nas seguintes etapas:
Chame
computeSplitspara obter o paralelismo máximosplitsSizee o ID da sessão da tarefasessionIddo search index.Configure
ParallelScanRequest. Para varredura com único worker, omitamaxParallelecurrentParallelId. Em cenários com múltiplos workers, use a mesma consulta,sessionIdemaxParallelpara todos os workers e especifique umcurrentParallelIddiferente para cada worker.Chame
createParallelScanIteratorpara ler automaticamente todas as páginas ou chameparallelScane usenextTokenpara recuperar manualmente as páginas subsequentes.Aguarde a conclusão de todos os workers e mescle os resultados não ordenados.
Dentro de uma sessionId, o snapshot de dados é fixado na primeira chamada ao parallelScan. Linhas adicionadas ou atualizadas durante a execução da tarefa não são incluídas no snapshot. É possível omitir o parâmetro sessionId. No entanto, se ocorrer balanceamento de carga no servidor ou mudança similar durante a varredura, os resultados poderão conter algumas linhas duplicadas. Recomendamos chamar primeiro o computeSplits e incluir o sessionId retornado nas solicitações seguintes.
A sessão pode expirar antecipadamente quando uma alteração dinâmica de schema troca o índice ou quando ocorre failover ou balanceamento de carga no servidor. Nesses casos, o servidor retorna OTSSessionExpired. Um erro de rede no lado do cliente também pode interromper a varredura. Se algum desses erros ocorrer, descarte os resultados incompletos, chame computeSplits novamente e reinicie toda a tarefa de varredura desde o início. No máximo 10 tarefas de parallel scan podem ser executadas simultaneamente em um search index. Para outros limites, consulte Search index limits.
Use computeSplits para calcular as divisões. Chame parallelScan para recuperar páginas manualmente ou chame createParallelScanIterator para obter automaticamente todas as páginas.
ComputeSplitsResponse computeSplits(ComputeSplitsRequest request)
ParallelScanResponse parallelScan(ParallelScanRequest request)
RowIterator createParallelScanIterator(ParallelScanRequest request)
O exemplo abaixo varre todas as linhas usando um único worker e retorna os campos category e price. O RowIterator recupera automaticamente as páginas subsequentes.
String tableName = "example_table";
String indexName = "example_index";
ComputeSplitsRequest splitsRequest = ComputeSplitsRequest.newBuilder()
.tableName(tableName)
.splitsOptions(new SearchIndexSplitsOptions(indexName))
.build();
ComputeSplitsResponse splitsResponse =
client.computeSplits(splitsRequest);
ScanQuery scanQuery = ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.build();
ParallelScanRequest request = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(scanQuery)
.addColumnsToGet("category", "price")
.sessionId(splitsResponse.getSessionId())
.build();
RowIterator iterator = client.createParallelScanIterator(request);
while (iterator.hasNext()) {
Row row = iterator.next();
System.out.println(row);
}
Parâmetros
Solicitação de divisão
O objeto splitsRequest é do tipo ComputeSplitsRequest e contém os seguintes parâmetros.
|
Nome |
Tipo |
Descrição |
|
tableName (obrigatório) |
String |
Nome da tabela de dados. |
|
splitsOptions (obrigatório) |
SplitsOptions |
Configuração da divisão. Para um search index, defina este parâmetro como |
Configuração de divisão do search index
O atributo splitsRequest.splitsOptions é do tipo SearchIndexSplitsOptions e contém o seguinte parâmetro.
|
Nome |
Tipo |
Descrição |
|
indexName (obrigatório) |
String |
Nome do search index. |
Solicitação de varredura
O objeto request é do tipo ParallelScanRequest e contém os seguintes parâmetros.
|
Nome |
Tipo |
Descrição |
|
tableName (obrigatório) |
String |
Nome da tabela de dados. |
|
indexName (obrigatório) |
String |
Nome do search index. |
|
scanQuery (obrigatório) |
ScanQuery |
Condição da varredura, quantidade de linhas retornadas por solicitação e configuração de paralelismo. |
|
columnsToGet (opcional) |
SearchRequest.ColumnsToGet |
Colunas a serem retornadas. Se este parâmetro não for especificado, apenas as colunas de chave primária serão retornadas. |
|
sessionId (opcional) |
|
ID da sessão da tarefa retornado por |
|
timeoutInMillisecond (opcional) |
int |
Tempo limite da solicitação em milissegundos. O valor padrão é |
Configuração da varredura
O atributo request.scanQuery é do tipo ScanQuery e contém os seguintes parâmetros.
|
Nome |
Tipo |
Descrição |
|
query (obrigatório) |
Query |
Condição de consulta que define o escopo da varredura. O parallel scan suporta consultas term, match, range, geo, nested, entre outras. Configure a consulta da mesma forma que na operação Search. Para varrer todas as linhas no search index, defina o tipo de consulta como |
|
limit (opcional) |
Integer |
Número máximo de linhas retornadas por solicitação. Valor padrão: 2000. Recomendamos manter o valor padrão. |
|
maxParallel (opcional) |
Integer |
Paralelismo da tarefa de varredura. O valor não pode exceder |
|
currentParallelId (opcional) |
Integer |
ID do worker atual. Este parâmetro é obrigatório se |
|
aliveTime (opcional) |
Integer |
Intervalo máximo entre duas solicitações de página para a tarefa de varredura, em segundos. Valores válidos: 1 a 600. Valor padrão: 60. O período de validade é renovado sempre que linhas são recuperadas com sucesso. |
|
token (opcional) |
|
Token de paginação. Omita este parâmetro na primeira solicitação. Para paginação manual, defina-o como |
O valor máximo de limit suportado pelo servidor é 10000. Recomendamos não definir limit com esse valor.
Colunas a retornar
O atributo request.columnsToGet é do tipo SearchRequest.ColumnsToGet e contém os seguintes parâmetros.
|
Nome |
Tipo |
Descrição |
|
columns (opcional) |
|
Nomes dos campos do search index a serem retornados. Um campo existente apenas na tabela de dados, mas não incluído no search index, não pode ser retornado. Campos Date, Geo-point, IP, Vector, JSON/Nested e array podem ser retornados. |
|
returnAllFromIndex (opcional) |
boolean |
Indica se todos os campos do search index devem ser retornados. Valor padrão: |
|
returnAll (opcional) |
boolean |
O parallel scan não suporta este parâmetro. Não o defina como |
Valores de retorno
Informações da divisão
O método computeSplits retorna ComputeSplitsResponse, que contém os seguintes campos.
|
Nome |
Tipo |
Descrição |
|
sessionId |
|
ID da sessão da tarefa usado para varrer linhas no mesmo snapshot de dados. |
|
splitsSize |
Integer |
Paralelismo máximo suportado pelo search index. |
Resultados da varredura
O método parallelScan retorna ParallelScanResponse, que contém os seguintes campos.
|
Nome |
Tipo |
Descrição |
|
rows |
|
Linhas retornadas na resposta atual. |
|
nextToken |
|
Token para a próxima página. Se este valor for |
|
bodyBytes |
long |
Tamanho do corpo da resposta em bytes. |
O método createParallelScanIterator retorna um RowIterator. O iterador usa automaticamente o nextToken para recuperar páginas subsequentes e retorna uma Row por iteração. Ele não suporta a recuperação do número total de linhas correspondentes.
Exemplos
Varredura com múltiplos workers
O exemplo abaixo cria múltiplas tarefas de varredura com base no splitsSize. Cada tarefa utiliza um currentParallelId exclusivo, enquanto todas compartilham o mesmo sessionId e maxParallel. O pool de threads não excede o número de núcleos de CPU do cliente para evitar sobrecarga excessiva.
String tableName = "example_table";
String indexName = "example_index";
ComputeSplitsResponse splitsResponse = client.computeSplits(
ComputeSplitsRequest.newBuilder()
.tableName(tableName)
.splitsOptions(new SearchIndexSplitsOptions(indexName))
.build());
int maxParallel = splitsResponse.getSplitsSize();
int workerCount = Math.min(
maxParallel, Runtime.getRuntime().availableProcessors());
ExecutorService executor = Executors.newFixedThreadPool(workerCount);
List<Future<Integer>> futures = new ArrayList<Future<Integer>>();
try {
for (int parallelId = 0; parallelId < maxParallel; parallelId++) {
final int currentParallelId = parallelId;
futures.add(executor.submit(new Callable<Integer>() {
@Override
public Integer call() {
ScanQuery scanQuery = ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.maxParallel(maxParallel)
.currentParallelId(currentParallelId)
.build();
ParallelScanRequest request =
ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(scanQuery)
.returnAllColumnsFromIndex(true)
.sessionId(splitsResponse.getSessionId())
.build();
int rowCount = 0;
RowIterator iterator =
client.createParallelScanIterator(request);
while (iterator.hasNext()) {
Row row = iterator.next();
System.out.println(row);
rowCount++;
}
return rowCount;
}
}));
}
long totalRows = 0;
for (Future<Integer> future : futures) {
totalRows += future.get();
}
System.out.println("Total rows: " + totalRows);
} finally {
executor.shutdown();
}
Recuperação manual de páginas
Este exemplo chama diretamente o parallelScan e passa o nextToken de cada resposta para a próxima solicitação até que o worker atual tenha recuperado todas as linhas.
String tableName = "example_table";
String indexName = "example_index";
ComputeSplitsResponse splitsResponse = client.computeSplits(
ComputeSplitsRequest.newBuilder()
.tableName(tableName)
.splitsOptions(new SearchIndexSplitsOptions(indexName))
.build());
ScanQuery scanQuery = ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.maxParallel(1)
.currentParallelId(0)
.build();
ParallelScanRequest request = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(scanQuery)
.addColumnsToGet("category", "price")
.sessionId(splitsResponse.getSessionId())
.build();
long totalRows = 0;
do {
ParallelScanResponse response = client.parallelScan(request);
for (Row row : response.getRows()) {
System.out.println(row);
totalRows++;
}
scanQuery.setToken(response.getNextToken());
} while (scanQuery.getToken() != null);
System.out.println("Total rows: " + totalRows);