Utilisez le SDK Tablestore pour Java pour analyser simultanément les lignes correspondantes dans un index de recherche et exporter l'ensemble complet des résultats lorsque l'ordre des résultats n'est pas requis.
Prérequis
Installez Tablestore SDK for Java et initialisez un client.
L'analyse parallèle nécessite le SDK Tablestore pour Java version 5.6.0 ou ultérieure. Pour plus d'informations sur les versions, consultez Historique des versions.
Fonctionnement
L'analyse parallèle parcourt toutes les lignes correspondant à une requête dans un index de recherche. Elle ne garantit pas l'ordre global des résultats et ne prend pas en charge le tri ni l'agrégation. Si vous avez besoin de résultats ordonnés, d'agrégations ou de résultats de recherche destinés aux utilisateurs finaux, utilisez l'opération Search.
Une analyse mono-worker est plus simple à configurer. Une analyse multi-workers lit plusieurs partitions simultanément et offre généralement un débit d'analyse plus élevé qu'une analyse mono-worker.
Une analyse parallèle se déroule selon les étapes suivantes :
Appelez
computeSplitspour obtenir le niveau de parallélisme maximalsplitsSizeet l'ID de session de tâchesessionIdde l'index de recherche.Configurez
ParallelScanRequest. Pour une analyse mono-worker, vous pouvez omettremaxParalleletcurrentParallelId. Pour une analyse multi-workers, utilisez la même requête, le mêmesessionIdet la même valeurmaxParallelpour tous les workers, et spécifiez une valeurcurrentParallelIddifférente pour chaque worker.Appelez
createParallelScanIteratorpour lire automatiquement toutes les pages, ou appelezparallelScanet utiliseznextTokenpour récupérer manuellement les pages suivantes.Attendez que tous les workers aient terminé et fusionnez leurs résultats non ordonnés.
Au sein d'un sessionId, l'instantané des données est figé lors du premier appel à parallelScan. Les lignes ajoutées ou mises à jour pendant l'exécution de la tâche ne sont pas incluses dans l'instantané. Vous pouvez omettre le paramètre sessionId. Toutefois, si un rééquilibrage de charge côté serveur ou une modification similaire survient pendant l'analyse, les résultats peuvent contenir un petit nombre de lignes en double. Nous vous recommandons d'appeler d'abord computeSplits et d'inclure le sessionId renvoyé dans les requêtes suivantes.
La session peut expirer prématurément en cas de changement dynamique de schéma modifiant l'index, ou lors d'un basculement ou d'un rééquilibrage de charge côté serveur. Dans ce cas, le serveur renvoie l'erreur OTSSessionExpired. Une erreur réseau côté client peut également interrompre l'analyse. Si l'une de ces erreurs se produit, ignorez les résultats incomplets, appelez à nouveau computeSplits et redémarrez entièrement la tâche d'analyse depuis le début. Un maximum de 10 tâches d'analyse parallèle peuvent s'exécuter simultanément sur un index de recherche. Pour connaître les autres limites, consultez Limites des index de recherche.
Appelez computeSplits pour calculer les partitions. Appelez parallelScan pour récupérer manuellement les pages, ou appelez createParallelScanIterator pour récupérer automatiquement toutes les pages.
ComputeSplitsResponse computeSplits(ComputeSplitsRequest request)
ParallelScanResponse parallelScan(ParallelScanRequest request)
RowIterator createParallelScanIterator(ParallelScanRequest request)
L'exemple suivant analyse toutes les lignes à l'aide d'un seul worker et renvoie les champs category et price. L'objet RowIterator récupère automatiquement les pages suivantes.
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);
}
Paramètres
Requête de partitionnement
Le paramètre splitsRequest est de type ComputeSplitsRequest et contient les paramètres suivants.
|
|
|
|
|
|
|
|
|
|
|
|
Configuration du partitionnement de l'index de recherche
Le paramètre splitsRequest.splitsOptions est de type SearchIndexSplitsOptions et contient le paramètre suivant.
|
|
|
|
|
|
|
|
Requête d'analyse
Le paramètre request est de type ParallelScanRequest et contient les paramètres suivants.
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
Configuration de l'analyse
Le paramètre request.scanQuery est de type ScanQuery et contient les paramètres suivants.
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
La valeur maximale de limit prise en charge par le serveur est 10000. Nous vous déconseillons de définir limit sur cette valeur.
Colonnes à renvoyer
Le paramètre request.columnsToGet est de type SearchRequest.ColumnsToGet et contient les paramètres suivants.
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
Valeurs de retour
Informations sur le partitionnement
La méthode computeSplits renvoie un objet ComputeSplitsResponse, qui contient les champs suivants.
|
|
|
|
|
|
|
|
|
|
|
|
Résultats de l'analyse
La méthode parallelScan renvoie un objet ParallelScanResponse, qui contient les champs suivants.
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
La méthode createParallelScanIterator renvoie un objet RowIterator. L'itérateur utilise automatiquement nextToken pour récupérer les pages suivantes et renvoie une Row par itération. Il ne prend pas en charge la récupération du nombre total de lignes correspondantes.
Exemples
Analyse avec plusieurs workers
L'exemple suivant crée plusieurs tâches d'analyse basées sur la valeur splitsSize. Chaque tâche utilise un currentParallelId unique, et toutes les tâches partagent le même sessionId et la même valeur maxParallel. Le pool de threads ne dépasse pas le nombre de cœurs CPU du client afin d'éviter une charge excessive côté client.
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();
}
Récupération manuelle des pages
L'exemple suivant appelle directement parallelScan et transmet la valeur nextToken de chaque réponse à la requête suivante jusqu'à ce que le worker actuel ait récupéré toutes les lignes.
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);