Se a ordem dos resultados da consulta não for um requisito, use o recurso de varredura paralela para obter esses dados com eficiência.
Tablestore SDK for Java V5.6.0 ou posterior oferece suporte ao recurso de varredura paralela. Antes de usar esse recurso, certifique-se de ter obtido a versão correta do Tablestore SDK for Java. Para mais informações sobre o histórico de versões do Tablestore SDK for Java, consulte Histórico de versões do Tablestore SDK for Java.
Informações básicas
O recurso de índice de busca permite chamar a operação Search para usar todos os recursos de consulta e capacidades analíticas, como classificação e agregação. A operação Search retorna os resultados da consulta em uma ordem específica.
Em alguns cenários, como ao conectar o Tablestore a um ambiente de computação (Spark ou Presto) ou ao consultar um grupo específico de objetos, a velocidade da consulta pode ser mais importante que a ordenação dos resultados. Para aumentar a velocidade das consultas, o Tablestore fornece a operação ParallelScan para o recurso de índice de busca.
Comparada à operação Search, a operação ParallelScan suporta todos os recursos de consulta, mas não oferece capacidades analíticas como classificação e agregação. Isso aumenta a velocidade das consultas em mais de cinco vezes. É possível chamar a operação ParallelScan para exportar centenas de milhões de linhas de dados em menos de um minuto. Além disso, a capacidade de exportação de dados pode ser dimensionada horizontalmente sem limites superiores.
O número máximo de linhas retornadas por cada chamada ParallelScan é superior ao limite da operação Search. Enquanto a operação Search retorna até 100 linhas por chamada, a operação ParallelScan retorna até 2.000 linhas. O recurso de varredura paralela permite usar múltiplas threads para iniciar requisições paralelas dentro de uma sessão, o que acelera significativamente a exportação de dados.
Cenários
Use a operação Search quando precisar classificar ou agregar resultados de consulta, ou quando a requisição for enviada diretamente por um usuário final.
Use a operação ParallelScan se não houver necessidade de ordenar os resultados e você desejar retornar todas as correspondências com eficiência, ou ainda quando os dados forem consumidos por ambientes de computação como Spark ou Presto.
Recursos
Os itens a seguir descrevem as diferenças entre a operação ParallelScan e a operação Search.
-
Resultados estáveis
As tarefas de varredura paralela mantêm estado. Dentro de uma sessão, o conjunto de resultados dos dados varridos é definido pelo estado dos dados no momento em que a primeira requisição é iniciada. Inserções ou modificações de dados ocorridas após o envio da primeira requisição não afetam o conjunto de resultados.
-
Sessões
ImportanteSe for difícil obter o ID da sessão, chame a operação ParallelScan para iniciar uma requisição sem esse identificador. No entanto, ao enviar uma requisição sem ID de sessão, existe uma probabilidade muito baixa de ocorrência de dados duplicados no conjunto de resultados obtido.
As operações relacionadas à varredura paralela usam sessões. O ID da sessão determina o conjunto de resultados dos dados varridos. O processo abaixo descreve como obter e usar um ID de sessão:
Chame a operação ComputeSplits para consultar o número máximo de tarefas de varredura paralela e o ID da sessão atual.
Inicie múltiplas requisições de varredura paralela para ler os dados. Especifique o ID da sessão atual e os IDs das tarefas de varredura paralela nessas requisições.
O Tablestore retorna o código de erro OTSSessionExpired quando ocorrem exceções de rede, exceções de thread, modificações dinâmicas em esquemas ou alternâncias de índice durante o processo de varredura paralela, interrompendo a leitura dos dados. Nesses casos, inicie uma nova tarefa de varredura paralela para ler os dados novamente.
Tarefas de varredura paralela que compartilham o mesmo ID de sessão e o mesmo valor do parâmetro ScanQuery são consideradas uma única tarefa. Uma tarefa de varredura paralela começa no momento em que você envia a primeira requisição ParallelScan e termina quando todos os dados são varridos ou quando o token expira.
-
Número máximo de tarefas de varredura paralela em uma única requisição
A resposta da requisição ComputeSplits determina o número máximo de tarefas de varredura paralela suportadas em uma única requisição pela operação ParallelScan. Um volume maior de dados exige mais tarefas de varredura paralela em uma sessão.
Uma única requisição é definida por uma instrução de consulta. Por exemplo, ao usar a operação Search para buscar resultados onde o valor de city é Hangzhou, todos os dados correspondentes retornam no resultado. Contudo, se você usar a operação ParallelScan e o número de tarefas de varredura paralela na sessão for 2, cada requisição ParallelScan retornará metade dos resultados. O conjunto completo de resultados consiste na união dos dois conjuntos paralelos.
-
Desempenho
A velocidade de consulta de uma requisição ParallelScan que inclui uma tarefa de varredura paralela é cinco vezes maior que a de uma requisição Search. Ao usar o recurso de varredura paralela, a velocidade de consulta aumenta proporcionalmente ao número de tarefas de varredura paralela na sessão. Por exemplo, se uma sessão incluir oito tarefas de varredura paralela, a velocidade de consulta poderá quadruplicar.
-
Custo
Requisições ParallelScan consomem menos recursos e têm preço mais baixo. Para exportar grandes volumes de dados, use a operação ParallelScan.
Limites
O limite máximo de tarefas de varredura paralela é 10. Ajuste esse limite conforme as necessidades do seu negócio.
-
Índices de busca podem retornar apenas colunas existentes. No entanto, colunas dos tipos DATE e NESTED não podem ser retornadas.
A operação ParallelScan pode retornar valores das colunas ARRAY e GEOPOINT. Porém, os valores retornados são formatados e podem diferir dos valores gravados na tabela de dados. Por exemplo, se você gravar [1,2, 3, 4] em uma coluna ARRAY, a operação ParallelScan retornará [1,2,3,4]. Se gravar
10,50em uma coluna GEOPOINT, a operação ParallelScan retornará10.0,50.0como valor.Defina o parâmetro ReturnType como RETURN_ALL_INDEX ou RETURN_SPECIFIED, mas nunca como RETURN_ALL.
O parâmetro limit define o número máximo de linhas retornadas por cada chamada ParallelScan. O valor padrão de limit é 2.000. Caso especifique um valor superior a 2.000, o desempenho praticamente não se altera com o aumento desse limite.
Operações de API
Chame as seguintes operações de API para usar o recurso de varredura paralela:
ComputeSplits: Use esta operação para consultar o número máximo de tarefas de varredura paralela para uma única requisição ParallelScan.
ParallelScan: Use esta operação para exportar dados.
Uso dos SDKs do Tablestore
Use os seguintes SDKs do Tablestore para realizar varredura paralela de dados:
Tablestore SDK for Java: Varredura paralela
Tablestore SDK for Go: Varredura paralela
Tablestore SDK for Python: Varredura paralela
Tablestore SDK for Node.js: Varredura paralela
Tablestore SDK for .NET: Varredura paralela
Tablestore SDK for PHP: Varredura paralela
Parâmetros
|
Parâmetro |
Descrição |
|
|
tableName |
Nome da tabela de dados. |
|
|
indexName |
Nome do índice de busca. |
|
|
scanQuery |
query |
Instrução de consulta para o índice de busca. A operação suporta consulta por termo, consulta difusa, consulta por intervalo, consulta geográfica e consulta aninhada, semelhantes às da operação Search. |
|
limit |
Número máximo de linhas que cada chamada ParallelScan pode retornar. |
|
|
maxParallel |
Quantidade máxima de tarefas de varredura paralela por requisição. Esse número varia conforme o volume de dados; volumes maiores exigem mais tarefas. Use a operação ComputeSplits para consultar o limite máximo de tarefas paralelas por requisição. |
|
|
currentParallelId |
ID da tarefa de varredura paralela na requisição. Valores válidos: [0, Valor de maxParallel). |
|
|
token |
Token usado para paginar os resultados da consulta. Os resultados da requisição ParallelScan contêm o token da próxima página. Use esse token para recuperar a página seguinte. |
|
|
aliveTime |
Período de validade da tarefa de varredura paralela atual, que também corresponde à validade do token. Unidade: segundos. Valor padrão: 60. Mantenha o valor padrão. Se a próxima requisição não for iniciada dentro desse período, não será possível consultar mais dados. O tempo de validade do token é renovado a cada requisição enviada. Nota
As sessões expiram antecipadamente se houver alteração dinâmica de índices de alternância nos esquemas, falha em um único servidor ou balanceamento de carga no lado do servidor. Nessas situações, recrie as sessões. |
|
|
columnsToGet |
Nome da coluna a ser retornada no resultado do agrupamento. Adicione o nome da coluna a Columns. Para retornar todas as colunas do índice de busca, use a operação ReturnAllFromIndex, que é mais concisa. Importante
Não é possível usar ReturnAll neste contexto. |
|
|
sessionId |
ID da sessão da tarefa de varredura paralela. Chame a operação ComputeSplits para criar uma sessão e consultar o número máximo de tarefas de varredura paralela suportadas pela requisição. |
|
Exemplo
Realize a varredura de dados usando uma única thread ou múltiplas threads simultaneamente, conforme as necessidades do seu negócio.
Varredura de dados com thread única
Ao usar a varredura paralela, o código para uma requisição com thread única é mais simples do que para múltiplas threads. Os parâmetros currentParallelId e maxParallel não são necessários nesse caso. A requisição ParallelScan com thread única oferece throughput superior ao da requisição Search, porém inferior ao da requisição ParallelScan com múltiplas threads.
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import com.alicloud.openservices.tablestore.SyncClient;
import com.alicloud.openservices.tablestore.model.ComputeSplitsRequest;
import com.alicloud.openservices.tablestore.model.ComputeSplitsResponse;
import com.alicloud.openservices.tablestore.model.Row;
import com.alicloud.openservices.tablestore.model.SearchIndexSplitsOptions;
import com.alicloud.openservices.tablestore.model.iterator.RowIterator;
import com.alicloud.openservices.tablestore.model.search.ParallelScanRequest;
import com.alicloud.openservices.tablestore.model.search.ParallelScanResponse;
import com.alicloud.openservices.tablestore.model.search.ScanQuery;
import com.alicloud.openservices.tablestore.model.search.SearchRequest.ColumnsToGet;
import com.alicloud.openservices.tablestore.model.search.query.MatchAllQuery;
import com.alicloud.openservices.tablestore.model.search.query.Query;
import com.alicloud.openservices.tablestore.model.search.query.QueryBuilders;
public class Test {
public static List<Row> scanQuery(final SyncClient client) {
String tableName = "<TableName>";
String indexName = "<IndexName>";
// Query the session ID and the maximum number of parallel scan tasks supported by the request.
ComputeSplitsRequest computeSplitsRequest = new ComputeSplitsRequest();
computeSplitsRequest.setTableName(tableName);
computeSplitsRequest.setSplitsOptions(new SearchIndexSplitsOptions(indexName));
ComputeSplitsResponse computeSplitsResponse = client.computeSplits(computeSplitsRequest);
byte[] sessionId = computeSplitsResponse.getSessionId();
int splitsSize = computeSplitsResponse.getSplitsSize();
/*
* Create a parallel scan request.
*/
ParallelScanRequest parallelScanRequest = new ParallelScanRequest();
parallelScanRequest.setTableName(tableName);
parallelScanRequest.setIndexName(indexName);
ScanQuery scanQuery = new ScanQuery();
// This query determines the range of the data to scan. You can create a nested and complex query.
Query query = new MatchAllQuery();
scanQuery.setQuery(query);
// Specify the maximum number of rows that can be returned by each ParallelScan call.
scanQuery.setLimit(2000);
parallelScanRequest.setScanQuery(scanQuery);
ColumnsToGet columnsToGet = new ColumnsToGet();
columnsToGet.setColumns(Arrays.asList("col_1", "col_2"));
parallelScanRequest.setColumnsToGet(columnsToGet);
parallelScanRequest.setSessionId(sessionId);
/*
* Use builder to create a parallel scan request that has the same features as the preceding request.
*/
ParallelScanRequest parallelScanRequestByBuilder = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.build())
.addColumnsToGet("col_1", "col_2")
.sessionId(sessionId)
.build();
List<Row> result = new ArrayList<>();
/*
* Use the native API operation to scan data.
*/
{
ParallelScanResponse parallelScanResponse = client.parallelScan(parallelScanRequest);
// Query the token of ScanQuery for the next request.
byte[] nextToken = parallelScanResponse.getNextToken();
// Obtain the data.
List<Row> rows = parallelScanResponse.getRows();
result.addAll(rows);
while (nextToken != null) {
// Specify the token.
parallelScanRequest.getScanQuery().setToken(nextToken);
// Continue to scan the data.
parallelScanResponse = client.parallelScan(parallelScanRequest);
// Obtain the data.
rows = parallelScanResponse.getRows();
result.addAll(rows);
nextToken = parallelScanResponse.getNextToken();
}
}
/*
* Recommended method.
* Use an iterator to scan all matched data. This method has the same query speed but is easier to use compared with the previous method.
*/
{
RowIterator iterator = client.createParallelScanIterator(parallelScanRequestByBuilder);
while (iterator.hasNext()) {
Row row = iterator.next();
result.add(row);
// Obtain the specific values.
String col_1 = row.getLatestColumn("col_1").getValue().asString();
long col_2 = row.getLatestColumn("col_2").getValue().asLong();
}
}
/*
* If the operation fails, retry the operation. If the caller of this function has a retry mechanism or if you do not want to retry the failed operation, you can ignore this part.
* To ensure availability, we recommend that you start a new parallel scan task when exceptions occur.
* The following exceptions may occur when you send a ParallelScan request:
* 1. A session exception occurs on the server side. The error code is OTSSessionExpired.
* 2. An exception such as a network exception occurs on the client side.
*/
try {
// Execute the processing logic.
{
RowIterator iterator = client.createParallelScanIterator(parallelScanRequestByBuilder);
while (iterator.hasNext()) {
Row row = iterator.next();
// Process rows of data. If you have sufficient memory resources, you can add the rows to a list.
result.add(row);
}
}
} catch (Exception ex) {
// Retry the processing logic.
{
result.clear();
RowIterator iterator = client.createParallelScanIterator(parallelScanRequestByBuilder);
while (iterator.hasNext()) {
Row row = iterator.next();
// Process rows of data. If you have sufficient memory resources, you can add the rows to a list.
result.add(row);
}
}
}
return result;
}
}
Varredura de dados com múltiplas threads
import com.alicloud.openservices.tablestore.SyncClient;
import com.alicloud.openservices.tablestore.model.ComputeSplitsRequest;
import com.alicloud.openservices.tablestore.model.ComputeSplitsResponse;
import com.alicloud.openservices.tablestore.model.Row;
import com.alicloud.openservices.tablestore.model.SearchIndexSplitsOptions;
import com.alicloud.openservices.tablestore.model.iterator.RowIterator;
import com.alicloud.openservices.tablestore.model.search.ParallelScanRequest;
import com.alicloud.openservices.tablestore.model.search.ScanQuery;
import com.alicloud.openservices.tablestore.model.search.query.QueryBuilders;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Semaphore;
import java.util.concurrent.atomic.AtomicLong;
public class Test {
public static void scanQueryWithMultiThread(final SyncClient client, String tableName, String indexName) throws InterruptedException {
// Query the number of CPU cores on the client.
final int cpuProcessors = Runtime.getRuntime().availableProcessors();
// Specify the number of parallel threads for the client. We recommend that you specify the number of CPU cores on the client as the number of parallel threads for the client to prevent impact on the query performance.
final Semaphore semaphore = new Semaphore(cpuProcessors);
// Query the session ID and the maximum number of parallel scan tasks supported by the request.
ComputeSplitsRequest computeSplitsRequest = new ComputeSplitsRequest();
computeSplitsRequest.setTableName(tableName);
computeSplitsRequest.setSplitsOptions(new SearchIndexSplitsOptions(indexName));
ComputeSplitsResponse computeSplitsResponse = client.computeSplits(computeSplitsRequest);
final byte[] sessionId = computeSplitsResponse.getSessionId();
final int maxParallel = computeSplitsResponse.getSplitsSize();
// Create an AtomicLong object if you need to obtain the row count for your business.
AtomicLong rowCount = new AtomicLong(0);
/*
* If you want to perform multithreading by using a function, you can build an internal class to inherit the threads.
* You can also build an external class to organize the code.
*/
final class ThreadForScanQuery extends Thread {
private final int currentParallelId;
private ThreadForScanQuery(int currentParallelId) {
this.currentParallelId = currentParallelId;
this.setName("ThreadForScanQuery:" + maxParallel + "-" + currentParallelId); // Specify the thread name.
}
@Override
public void run() {
System.out.println("start thread:" + this.getName());
try {
// Execute the processing logic.
{
ParallelScanRequest parallelScanRequest = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(ScanQuery.newBuilder()
.query(QueryBuilders.range("col_long").lessThan(10_0000)) // Specify the data to query.
.limit(2000)
.currentParallelId(currentParallelId)
.maxParallel(maxParallel)
.build())
.addColumnsToGet("col_long", "col_keyword", "col_bool") // Specify the fields to return from the search index. To return all fields from the search index, set returnAllColumnsFromIndex to true.
//.returnAllColumnsFromIndex(true)
.sessionId(sessionId)
.build();
// Use an iterator to obtain all the data.
RowIterator ltr = client.createParallelScanIterator(parallelScanRequest);
long count = 0;
while (ltr.hasNext()) {
Row row = ltr.next();
// Add a custom processing logic. The following sample code shows how to add a custom processing logic to count the number of rows:
count++;
}
rowCount.addAndGet(count);
System.out.println("thread[" + this.getName() + "] finished. this thread get rows:" + count);
}
} catch (Exception ex) {
// If exceptions occur, you can retry the processing logic.
} finally {
semaphore.release();
}
}
}
// Simultaneously execute threads. Valid values of currentParallelId: [0, Value of maxParallel).
List<ThreadForScanQuery> threadList = new ArrayList<ThreadForScanQuery>();
for (int currentParallelId = 0; currentParallelId < maxParallel; currentParallelId++) {
ThreadForScanQuery thread = new ThreadForScanQuery(currentParallelId);
threadList.add(thread);
}
// Simultaneously initiate the threads.
for (ThreadForScanQuery thread : threadList) {
// Specify a value for semaphore to limit the number of threads that can be initiated at the same time to prevent bottlenecks on the client.
semaphore.acquire();
thread.start();
}
// The main thread is blocked until all threads are complete.
for (ThreadForScanQuery thread : threadList) {
thread.join();
}
System.out.println("all thread finished! total rows:" + rowCount.get());
}
}