Se a ordem dos resultados da consulta não for relevante, use o recurso de varredura paralela para obter os dados com maior eficiência.
Pré-requisitos
Instância OTSClient inicializada. Para mais informações, consulte Inicializar um cliente do Tablestore.
Tabela de dados criada e preenchida. Para mais informações, consulte Criar uma tabela de dados e Gravar dados.
Índice de pesquisa criado para a tabela de dados. Para mais informações, consulte Criar um índice de pesquisa.
Parâmetros
|
Parâmetro |
Descrição |
|
|
table_name |
Nome da tabela de dados. |
|
|
index_name |
Nome do índice de pesquisa. |
|
|
scan_query |
query |
Tipo da consulta. A operação aceita consultas por termo, fuzzy, intervalo, geográfica e aninhada, semelhantes às suportadas pela operação Search. |
|
limit |
Número máximo de linhas retornadas em cada chamada ao ParallelScan. |
|
|
max_parallel |
Número máximo de tarefas de varredura paralela por requisição. Esse valor varia conforme o volume de dados: quanto maior o volume, mais tarefas são necessárias. Use a operação ComputeSplits para consultar o número máximo de tarefas de varredura paralela por requisição. |
|
|
current_parallel_id |
ID da tarefa de varredura paralela na requisição. Valores válidos: [0, max_parallel). |
|
|
token |
Token usado para paginar os resultados da consulta. A resposta da requisição ParallelScan inclui o token da próxima página. Use esse token para recuperar a próxima página. |
|
|
alive_time |
Período de validade da tarefa de varredura paralela atual. Esse período também define a validade do token. Unidade: segundos. Valor padrão: 60. Recomendamos usar o valor padrão. Se nenhuma nova requisiçã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 troca de esquemas entre o índice source e o canary, falha em um único servidor ou balanceamento de carga no lado do servidor. Nesses casos, crie novamente as sessões. |
|
|
columns_to_get |
Nomes das colunas a serem retornadas para cada linha que atenda às condições da consulta. Para retornar todas as colunas do índice de pesquisa, defina o parâmetro return_type como RETURN_ALL_FROM_INDEX. |
|
|
session_id |
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. |
|
Exemplos
O código de exemplo a seguir mostra como executar uma varredura paralela:
def fetch_rows_per_thread(query, session_id, current_thread_id, max_thread_num):
token = None
while True:
try:
scan_query = ScanQuery(query, limit = 20, next_token = token, current_parallel_id = current_thread_id,
max_parallel = max_thread_num, alive_time = 30)
response = client.parallel_scan(
table_name, index_name, scan_query, session_id,
columns_to_get = ColumnsToGet(return_type=ColumnReturnType.ALL_FROM_INDEX))
for row in response.rows:
print("%s:%s" % (threading.currentThread().name, str(row)))
if len(response.next_token) == 0:
break
else:
token = response.next_token
except OTSServiceError as e:
print (e)
except OTSClientError as e:
print (e)
def parallel_scan(table_name, index_name):
response = client.compute_splits(table_name, index_name)
query = TermQuery('d', 0.1)
params = []
for i in range(response.splits_size):
params.append((([query, response.session_id, i, response.splits_size], None)))
pool = threadpool.ThreadPool(response.splits_size)
requests = threadpool.makeRequests(fetch_rows_per_thread, params)
[pool.putRequest(req) for req in requests]
pool.wait()