Tous les produits
Search
Centre de documentation

Tablestore:Analyse parallèle

Dernière mise à jour :Aug 08, 2026

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 :

  1. Appelez computeSplits pour obtenir le niveau de parallélisme maximal splitsSize et l'ID de session de tâche sessionId de l'index de recherche.

  2. Configurez ParallelScanRequest. Pour une analyse mono-worker, vous pouvez omettre maxParallel et currentParallelId. Pour une analyse multi-workers, utilisez la même requête, le même sessionId et la même valeur maxParallel pour tous les workers, et spécifiez une valeur currentParallelId différente pour chaque worker.

  3. Appelez createParallelScanIterator pour lire automatiquement toutes les pages, ou appelez parallelScan et utilisez nextToken pour récupérer manuellement les pages suivantes.

  4. 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.

Important

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.


Nom


Type


Description


tableName (obligatoire)


String


Le nom de la table de données.


splitsOptions (obligatoire)


SplitsOptions


La configuration du partitionnement. Pour un index de recherche, définissez ce paramètre sur SearchIndexSplitsOptions.

Configuration du partitionnement de l'index de recherche

Le paramètre splitsRequest.splitsOptions est de type SearchIndexSplitsOptions et contient le paramètre suivant.


Nom


Type


Description


indexName (obligatoire)


String


Le nom de l'index de recherche.

Requête d'analyse

Le paramètre request est de type ParallelScanRequest et contient les paramètres suivants.


Nom


Type


Description


tableName (obligatoire)


String


Le nom de la table de données.


indexName (obligatoire)


String


Le nom de l'index de recherche.


scanQuery (obligatoire)


ScanQuery


La condition d'analyse, le nombre de lignes renvoyées par requête et la configuration du parallélisme.


columnsToGet (facultatif)


SearchRequest.ColumnsToGet


Les colonnes à renvoyer. Si ce paramètre n'est pas spécifié, seules les colonnes de clé primaire sont renvoyées.


sessionId (facultatif)


byte[]


L'ID de session de tâche renvoyé par computeSplits. Nous vous recommandons de spécifier ce paramètre pour utiliser le même instantané de données tout au long de l'analyse.


timeoutInMillisecond (facultatif)


int


Le délai d'expiration au niveau de la requête, en millisecondes. La valeur par défaut est -1, ce qui signifie qu'aucun délai d'expiration n'est défini au niveau de la requête.

Configuration de l'analyse

Le paramètre request.scanQuery est de type ScanQuery et contient les paramètres suivants.


Nom


Type


Description


query (obligatoire)


Query


La condition de requête qui définit la portée de l'analyse. L'analyse parallèle prend en charge les requêtes term, match, range, geo, nested, etc. Configurez la requête de la même manière que pour l'opération Search. Pour analyser toutes les lignes de l'index de recherche, définissez le type de requête sur MatchAllQuery.


limit (facultatif)


Integer


Le nombre maximal de lignes renvoyées par requête. Valeur par défaut : 2000. Nous vous recommandons d'utiliser la valeur par défaut.


maxParallel (facultatif)


Integer


Le parallélisme de la tâche d'analyse. La valeur ne peut pas dépasser ComputeSplitsResponse.splitsSize. Valeur par défaut : 1.


currentParallelId (facultatif)


Integer


L'ID du worker actuel. Ce paramètre est requis si maxParallel est supérieur à 1. Attribuez à chaque worker une valeur unique dans la plage [0, maxParallel).


aliveTime (facultatif)


Integer


L'intervalle maximal entre deux demandes de pages pour la tâche d'analyse, en secondes. Valeurs valides : 1 à 600. Valeur par défaut : 60. La période de validité est actualisée à chaque fois que des lignes sont récupérées avec succès.


token (facultatif)


byte[]


Le jeton de pagination. Omettez ce paramètre lors de la première requête. Pour une pagination manuelle, définissez-le sur la valeur nextToken de la réponse précédente. Le SDK gère ce paramètre lorsque vous utilisez createParallelScanIterator.

Remarque

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.


Nom


Type


Description


columns (facultatif)


List<String>


Les noms des champs de l'index de recherche à renvoyer. Un champ qui existe uniquement dans la table de données mais qui n'est pas inclus dans l'index de recherche ne peut pas être renvoyé. Les champs Date, Geo-point, IP, Vector, JSON/Nested et array peuvent être renvoyés.


returnAllFromIndex (facultatif)


boolean


Indique s'il faut renvoyer tous les champs de l'index de recherche. Valeur par défaut : false. Si vous définissez ce paramètre sur true, vous n'avez pas besoin de spécifier columns.


returnAll (facultatif)


boolean


L'analyse parallèle ne prend pas en charge ce paramètre. Ne le définissez pas sur true.

Valeurs de retour

Informations sur le partitionnement

La méthode computeSplits renvoie un objet ComputeSplitsResponse, qui contient les champs suivants.


Nom


Type


Description


sessionId


byte[]


L'ID de session de tâche utilisé pour analyser les lignes dans le même instantané de données.


splitsSize


Integer


Le parallélisme maximal pris en charge par l'index de recherche.

Résultats de l'analyse

La méthode parallelScan renvoie un objet ParallelScanResponse, qui contient les champs suivants.


Nom


Type


Description


rows


List<Row>


Les lignes renvoyées dans la réponse actuelle.


nextToken


byte[]


Le jeton pour la page suivante. Si cette valeur est null, le worker actuel a récupéré toutes les lignes.


bodyBytes


long


La taille du corps de la réponse en octets.

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);