Utilisez le SDK Tablestore pour Python pour analyser en parallèle les données correspondantes dans un index de recherche et exporter l'intégralité des résultats non triés.
Prérequis
Installez le SDK Tablestore pour Python et initialisez un client.
Description
L'analyse parallèle exporte toutes les lignes correspondant à une requête dans un index de recherche. Les résultats ne sont pas triés globalement ; le tri et l'agrégation ne sont pas pris en charge. Pour trier ou agréger les résultats, ou pour renvoyer des résultats de recherche aux utilisateurs finaux, utilisez l'API Search. La configuration d'un seul worker est plus simple, tandis que plusieurs workers lisent simultanément plusieurs partitions (splits) et offrent généralement un débit plus élevé.
Le flux de travail est le suivant : appelez compute_splits pour obtenir la concurrence maximale splits_size et l'identifiant de session session_id. Configurez chaque worker avec la même requête, le même session_id et la même valeur max_parallel, mais avec un identifiant unique current_parallel_id. Chaque worker appelle parallel_scan et utilise next_token pour lire les pages suivantes. Enfin, attendez la fin de tous les workers et fusionnez leurs résultats non triés.
Les workers partageant le même session_id établissent un instantané des données au début de la première analyse. Des mises à jour dynamiques du schéma, un basculement ou un équilibrage de charge peuvent expirer prématurément une session et renvoyer l'erreur OTSSessionExpired ; des pannes réseau peuvent également interrompre une analyse. Dans ces cas, ignorez les résultats incomplets, appelez à nouveau compute_splits et redémarrez tous les workers depuis le début. Jusqu'à 10 tâches d'analyse parallèle peuvent s'exécuter simultanément sur un même index de recherche.
L'exemple suivant calcule les partitions, puis analyse toutes les données à l'aide d'un seul worker.
splits = client.compute_splits("example_table", "example_index")
next_token = None
rows = []
while True:
scan_query = ScanQuery(
MatchAllQuery(),
limit=2000,
next_token=next_token,
current_parallel_id=0,
max_parallel=1,
alive_time=60,
)
response = client.parallel_scan(
"example_table",
"example_index",
scan_query,
splits.session_id,
ColumnsToGet(return_type=ColumnReturnType.ALL_FROM_INDEX),
)
rows.extend(response.rows)
next_token = response.next_token
if not next_token:
break
print(len(rows))
Paramètres
Calcul des partitions
La méthode compute_splits(table_name, index_name) accepte les paramètres suivants.
|
Nom |
Type |
Description |
|
table_name (obligatoire) |
|
Le nom de la table de données. |
|
index_name (obligatoire) |
|
Le nom de l'index de recherche. |
Requête d'analyse
La méthode parallel_scan accepte les paramètres suivants.
|
Nom |
Type |
Description |
|
table_name (obligatoire) |
|
Le nom de la table de données. |
|
index_name (obligatoire) |
|
Le nom de l'index de recherche. |
|
scan_query (obligatoire) |
|
La condition d'analyse, la pagination et les configurations de concurrence. |
|
session_id (obligatoire) |
|
L'identifiant de session renvoyé par |
|
columns_to_get (facultatif) |
|
La configuration des colonnes à renvoyer. Si omis, seules les colonnes de clé primaire sont renvoyées. |
|
timeout_s (facultatif) |
|
Le délai d'expiration de la requête en secondes. |
Configuration de l'analyse
L'objet scan_query est de type ScanQuery et accepte les paramètres suivants.
|
Nom |
Type |
Description |
|
query (obligatoire) |
|
La condition de requête qui définit la portée de l'analyse. Utilisez |
|
limit (obligatoire) |
|
Le nombre maximal de lignes par requête. Nous recommandons la valeur par défaut |
|
next_token (obligatoire) |
|
Le jeton de pagination. Définissez-le sur |
|
current_parallel_id (obligatoire) |
|
L'ID du worker actuel. Valeurs valides : |
|
max_parallel (obligatoire) |
|
La concurrence de la tâche, qui ne peut pas dépasser |
|
alive_time (facultatif) |
|
La période de validité maximale entre deux requêtes de page, en secondes. Valeurs valides : 1 à 600. Valeur par défaut : |
Colonnes de retour
L'objet columns_to_get est de type ColumnsToGet et accepte les paramètres suivants.
|
Nom |
Type |
Description |
|
column_names (facultatif) |
|
Les noms des champs de l'index de recherche à renvoyer. Spécifiez ce paramètre uniquement lorsque |
|
return_type (facultatif) |
|
Le mode de retour des colonnes. L'analyse parallèle prend en charge |
Réponse
Informations sur les partitions
La méthode compute_splits renvoie les informations sur les partitions.
|
Champ |
Type |
Description |
|
session_id |
|
L'identifiant de session de la tâche. |
|
splits_size |
|
La concurrence maximale prise en charge par l'index de recherche. |
Résultat de l'analyse
La méthode parallel_scan renvoie un résultat d'analyse.
|
Champ |
Type |
Description |
|
rows |
|
Les lignes renvoyées par la requête. |
|
next_token |
|
Le jeton pour la page suivante. Une valeur vide indique que le worker actuel a terminé. |
Réponses compatibles avec les tuples
L'analyse parallèle est prise en charge à partir de la version 5.2.0 du SDK Tablestore pour Python, qui renvoie des objets de réponse. À partir de la version 5.2.1, vous pouvez appeler ComputeSplitsResponse.v1_response() et ParallelScanResponse.v1_response() pour obtenir des tuples. Pour le nouveau code, accédez directement aux attributs de l'objet de réponse.
session_id, splits_size = splits.v1_response()
rows, next_token = response.v1_response()
Exemples
Analyse avec plusieurs workers
L'exemple suivant crée des workers en fonction de splits_size. Le pool de threads ne dépasse pas le nombre de cœurs CPU du client et tous les workers partagent la session et la concurrence maximale.
from concurrent.futures import ThreadPoolExecutor
import os
def scan_split(parallel_id, max_parallel, session_id):
rows = []
next_token = None
while True:
scan_query = ScanQuery(
MatchAllQuery(),
limit=2000,
next_token=next_token,
current_parallel_id=parallel_id,
max_parallel=max_parallel,
alive_time=60,
)
response = client.parallel_scan(
"example_table",
"example_index",
scan_query,
session_id,
ColumnsToGet(return_type=ColumnReturnType.ALL_FROM_INDEX),
)
rows.extend(response.rows)
next_token = response.next_token
if not next_token:
return rows
splits = client.compute_splits("example_table", "example_index")
worker_count = min(splits.splits_size, os.cpu_count() or 1)
with ThreadPoolExecutor(max_workers=worker_count) as executor:
futures = [
executor.submit(
scan_split,
parallel_id,
splits.splits_size,
splits.session_id,
)
for parallel_id in range(splits.splits_size)
]
all_rows = [row for future in futures for row in future.result()]
print(len(all_rows))