Se a ordem dos resultados da consulta não for relevante, use o recurso de varredura paralela para obter os dados com eficiência.
Pré-requisitos
Instância do OTSClient inicializada. Para mais informações, consulte Inicializar uma instância do OTSClient.
Tabela de dados criada e preenchida. Para mais detalhes, consulte Criar tabelas de dados e Gravar dados.
Índice de pesquisa criado para a tabela de dados. Consulte Criar índices de pesquisa.
Parâmetros
A varredura paralela exige duas operações coordenadas: chame ComputeSplits para obter o ID da sessão e o número máximo de tarefas paralelas (MaxParallel). Em seguida, inicie uma tarefa ParallelScan para cada slot paralelo, cada uma com um CurrentParallelId exclusivo variando de 0 a MaxParallel - 1. Por exemplo, se MaxParallel for 4, inicie quatro tarefas de varredura com valores de CurrentParallelId iguais a 0, 1, 2 e 3 para cobrir todo o conjunto de dados.
|
Parâmetro |
Descrição |
|
|
TableName |
Nome da tabela de dados. |
|
|
IndexName |
Nome do índice de pesquisa. |
|
|
ScanQuery |
Query |
Condição de consulta. Tipos compatíveis: consulta por termo, consulta difusa, consulta por intervalo, consulta geográfica e consulta aninhada — os mesmos tipos aceitos pela operação Search. |
|
Limit |
Número máximo de linhas retornadas por chamada ao ParallelScan. |
|
|
MaxParallel |
Quantidade máxima de tarefas de varredura paralela por requisição. O valor limite depende do volume de dados; conjuntos maiores permitem mais tarefas simultâneas. Chame ComputeSplits para obter esse valor antes de iniciar a varredura. Cada requisição ParallelScan utiliza um CurrentParallelId no intervalo [0, MaxParallel). É necessário cobrir todos os IDs desse intervalo para varrer o conjunto completo de dados. Se MaxParallel for 4, por exemplo, inicie quatro tarefas com CurrentParallelId definidos como 0, 1, 2 e 3. MaxParallel e CurrentParallelId devem ser usados em conjunto. |
|
|
CurrentParallelId |
Identificador desta tarefa de varredura paralela. Valores válidos: [0, MaxParallel). Cada tarefa concorrente deve usar um ID exclusivo. Uso obrigatório em conjunto com MaxParallel. |
|
|
Token |
Token de paginação para buscar a próxima página de resultados. Cada resposta do ParallelScan inclui um token para a página seguinte. Passe-o na requisição subsequente para continuar a leitura. Quando o token for nulo, todos os dados dessa tarefa já foram recuperados. |
|
|
AliveTime |
Período de validade da sessão e do token associado, em segundos. Valor padrão: 60. Use o valor padrão. Se nenhuma requisição for enviada dentro desse período, a sessão expira e não é possível recuperar mais dados. A validade é renovada a cada requisição. Nota
As sessões podem expirar antecipadamente se o esquema do índice de pesquisa for modificado dinamicamente, se houver falha em um único servidor ou se ocorrer balanceamento de carga entre servidores. Nesses casos, crie a sessão novamente. |
|
|
ColumnsToGet |
Colunas a serem retornadas. Defina o parâmetro Columns para especificar os nomes das colunas. Para retornar todas as colunas do índice de pesquisa, defina ReturnAllFromIndex como true. Importante
O parâmetro ReturnAll não é compatível. |
|
|
SessionId |
ID da sessão para a varredura paralela. Chame ComputeSplits para criar uma sessão e obter tanto o ID da sessão quanto o número máximo de tarefas paralelas compatíveis (MaxParallel). |
|
Exemplo
O exemplo abaixo varre todos os dados usando uma única thread. Ele chama ComputeSplits para obter o ID da sessão e o MaxParallel e, em seguida, percorre as páginas de resultados usando NextToken até que todas as linhas sejam recuperadas.
/// <summary>
/// Scans all data in a search index using a single thread.
/// Calls ComputeSplits to get the session ID and MaxParallel,
/// then pages through results until NextToken is null.
/// </summary>
public class ParallelScan
{
public static void ParallelScanwithSingleThread(OTSClient otsClient)
{
SearchIndexSplitsOptions options = new SearchIndexSplitsOptions
{
IndexName = IndexName
};
ComputeSplitsRequest computeSplitsRequest = new ComputeSplitsRequest
{
TableName = TableName,
SplitOptions = options
};
ComputeSplitsResponse computeSplitsResponse = otsClient.ComputeSplits(computeSplitsRequest);
MatchAllQuery matchAllQuery = new MatchAllQuery();
ScanQuery scanQuery = new ScanQuery();
scanQuery.AliveTime = 60;
scanQuery.Query = matchAllQuery;
scanQuery.MaxParallel = computeSplitsResponse.SplitsSize;
scanQuery.Limit = 10;
ParallelScanRequest parallelScanRequest = new ParallelScanRequest();
parallelScanRequest.TableName = TableName;
parallelScanRequest.IndexName = IndexName;
parallelScanRequest.ScanQuery = scanQuery;
parallelScanRequest.ColumnToGet = new ColumnsToGet { ReturnAllFromIndex = true };
parallelScanRequest.SessionId = computeSplitsResponse.SessionId;
int total = 0;
List<Row> result = new List<Row>();
ParallelScanResponse parallelScanResponse = otsClient.ParallelScan(parallelScanRequest);
while (parallelScanResponse.NextToken != null)
{
List<Row> rows = new List<Row>(parallelScanResponse.Rows);
total += rows.Count;
result.AddRange(rows);
parallelScanRequest.ScanQuery.Token = parallelScanResponse.NextToken;
parallelScanResponse = otsClient.ParallelScan(parallelScanRequest);
}
foreach (Row row in result)
{
Console.WriteLine(JsonConvert.SerializeObject(row));
}
Console.WriteLine("Total Row Count: {0}", total);
}
}