Si l'ordre des résultats de requête n'a pas d'importance, utilisez la fonctionnalité d'analyse parallèle pour obtenir les résultats efficacement.
Tablestore SDK for Java version 5.6.0 ou ultérieure prend en charge la fonctionnalité d'analyse parallèle. Avant d'utiliser cette fonctionnalité, assurez-vous de disposer de la version appropriée du SDK Tablestore for Java. Pour plus d'informations sur l'historique des versions du SDK Tablestore for Java, consultez Historique des versions du SDK Tablestore for Java.
Informations générales
La fonctionnalité d'index de recherche vous permet d'appeler l'opération Search pour utiliser toutes les fonctionnalités de requête et d'analyse, telles que le tri et l'agrégation. L'opération Search renvoie les résultats de requête dans un ordre spécifique.
Dans certains cas, par exemple lorsque vous connectez Tablestore à un environnement de calcul tel que Spark ou Presto, ou si vous souhaitez interroger un groupe spécifique d'objets, la vitesse de requête prime sur l'ordre des résultats. Pour accélérer les requêtes, Tablestore met à disposition l'opération ParallelScan pour la fonctionnalité d'index de recherche.
Contrairement à l'opération Search, l'opération ParallelScan prend en charge toutes les fonctionnalités de requête mais ne propose pas de capacités d'analyse telles que le tri et l'agrégation. Cette approche multiplie la vitesse de requête par plus de cinq. Vous pouvez appeler l'opération ParallelScan pour exporter des centaines de millions de lignes de données en moins d'une minute. La capacité d'exportation des données est évolutive horizontalement sans limite supérieure.
Le nombre maximal de lignes renvoyées par chaque appel ParallelScan est supérieur à celui de l'opération Search. L'opération Search renvoie jusqu'à 100 lignes par appel, tandis que l'opération ParallelScan en renvoie jusqu'à 2 000. La fonctionnalité d'analyse parallèle vous permet d'utiliser plusieurs threads pour lancer des requêtes en parallèle au sein d'une session, ce qui accélère l'exportation des données.
Scénarios
Si vous devez trier ou agréger les résultats de requête, ou si la requête est envoyée par un utilisateur final, utilisez l'opération Search.
Si le tri des résultats n'est pas nécessaire et que vous souhaitez renvoyer tous les résultats correspondants efficacement, ou si les données sont extraites par un environnement de calcul tel que Spark ou Presto, utilisez l'opération ParallelScan.
Fonctionnalités
Les éléments suivants décrivent les différences entre l'opération ParallelScan et l'opération Search.
-
Stabilité des résultats
Les tâches d'analyse parallèle conservent un état. Au sein d'une session, l'ensemble des résultats des données analysées est déterminé par l'état des données au moment du lancement de la première requête. Si des données sont insérées ou modifiées après l'envoi de la première requête, l'ensemble des résultats n'est pas affecté.
-
Sessions
ImportantS'il est difficile d'obtenir l'ID de session, vous pouvez appeler l'opération ParallelScan pour lancer une requête sans ID de session. Toutefois, l'envoi d'une requête sans ID de session peut entraîner, avec une très faible probabilité, la présence de données en double dans l'ensemble des résultats obtenus.
Les opérations liées à l'analyse parallèle utilisent des sessions. L'ID de session permet de déterminer l'ensemble des résultats des données analysées. Le processus suivant décrit comment obtenir et utiliser un ID de session :
Appelez l'opération ComputeSplits pour interroger le nombre maximal de tâches d'analyse parallèle et l'ID de session actuel.
Lancez plusieurs requêtes d'analyse parallèle pour lire les données. Vous devez spécifier l'ID de session actuel et les ID des tâches d'analyse parallèle dans ces requêtes.
Tablestore renvoie le code d'erreur OTSSessionExpired lorsque des exceptions réseau, des exceptions de thread, des modifications dynamiques des schémas ou des basculements d'index se produisent pendant le processus d'analyse parallèle et que les analyses de données s'arrêtent. Dans ces cas, lancez une autre tâche d'analyse parallèle pour analyser à nouveau les données.
Les tâches d'analyse parallèle partageant le même ID de session et la même valeur de paramètre ScanQuery sont considérées comme une seule tâche. Une tâche d'analyse parallèle commence lors de l'envoi de la première requête ParallelScan et se termine lorsque toutes les données ont été analysées ou que le jeton expire.
-
Nombre maximal de tâches d'analyse parallèle dans une seule requête
Le nombre maximal de tâches d'analyse parallèle prises en charge dans une seule requête par l'opération ParallelScan est déterminé par la réponse de la requête ComputeSplits. Un volume de données plus important nécessite davantage de tâches d'analyse parallèle au sein d'une session.
Une seule requête est définie par une instruction de requête. Par exemple, si vous utilisez l'opération Search pour interroger les résultats où la valeur de city est Hangzhou, toutes les données correspondant à cette condition sont renvoyées dans le résultat. En revanche, si vous utilisez l'opération ParallelScan et que le nombre de tâches d'analyse parallèle dans une session est de 2, chaque requête ParallelScan renvoie la moitié des résultats. L'ensemble complet des résultats est constitué des deux ensembles de résultats parallèles.
-
Performances
La vitesse de requête d'une requête ParallelScan incluant une tâche d'analyse parallèle est cinq fois plus rapide que celle d'une requête Search. Lorsque vous utilisez la fonctionnalité d'analyse parallèle, la vitesse de requête augmente avec le nombre de tâches d'analyse parallèle dans une session. Par exemple, si huit tâches d'analyse parallèle sont incluses dans une session, la vitesse de requête peut être multipliée par quatre.
-
Coût
Les requêtes ParallelScan consomment moins de ressources et sont proposées à un tarif inférieur. Pour exporter de grandes quantités de données, nous vous recommandons d'utiliser l'opération ParallelScan.
Limites
Le nombre maximal de tâches d'analyse parallèle est de 10. Vous pouvez ajuster cette limite en fonction de vos besoins métier.
-
Seules les colonnes existantes peuvent être renvoyées à partir des index de recherche. Cependant, les colonnes de type DATE et NESTED ne peuvent pas être renvoyées.
L'opération ParallelScan peut renvoyer les valeurs des colonnes ARRAY et GEOPOINT. Toutefois, les valeurs renvoyées sont formatées et peuvent différer des valeurs écrites dans la table de données. Par exemple, si vous écrivez [1,2, 3, 4] dans une colonne ARRAY, l'opération ParallelScan renvoie [1,2,3,4] comme valeur. Si vous écrivez
10,50dans une colonne GEOPOINT, l'opération ParallelScan renvoie10.0,50.0comme valeur.Vous pouvez définir le paramètre ReturnType sur RETURN_ALL_INDEX ou RETURN_SPECIFIED, mais pas sur RETURN_ALL.
Le nombre maximal de lignes renvoyées par chaque appel ParallelScan est spécifié par le paramètre limit. La valeur par défaut du paramètre limit est 2 000. Si vous spécifiez une valeur supérieure à 2 000, les performances ne changent pratiquement pas avec l'augmentation de la limite.
Opérations API
Vous pouvez appeler les opérations API suivantes pour utiliser la fonctionnalité d'analyse parallèle :
ComputeSplits : appelez cette opération pour interroger le nombre maximal de tâches d'analyse parallèle pour une seule requête ParallelScan.
ParallelScan : appelez cette opération pour exporter des données.
Utilisation des SDK Tablestore
Vous pouvez utiliser les SDK Tablestore suivants pour analyser les données en parallèle :
SDK Tablestore for Java : Analyse parallèle
SDK Tablestore for Go : Analyse parallèle
SDK Tablestore for Python : Analyse parallèle
SDK Tablestore for Node.js : Analyse parallèle
SDK Tablestore for .NET : Analyse parallèle
SDK Tablestore for PHP : Analyse parallèle
Paramètres
|
Paramètre |
Description |
|
|
tableName |
Nom de la table de données. |
|
|
indexName |
Nom de l'index de recherche. |
|
|
scanQuery |
query |
Instruction de requête pour l'index de recherche. L'opération prend en charge la requête term, la requête floue, la requête de plage, la requête géographique et la requête imbriquée, similaires à celles de l'opération Search. |
|
limit |
Nombre maximal de lignes renvoyées par chaque appel ParallelScan. |
|
|
maxParallel |
Nombre maximal de tâches d'analyse parallèle par requête. Ce nombre varie en fonction du volume de données. Un volume de données plus important nécessite davantage de tâches d'analyse parallèle par requête. Vous pouvez utiliser l'opération ComputeSplits pour interroger le nombre maximal de tâches d'analyse parallèle par requête. |
|
|
currentParallelId |
ID de la tâche d'analyse parallèle dans la requête. Valeurs valides : [0, Valeur de maxParallel) |
|
|
token |
Jeton utilisé pour la pagination des résultats de requête. Les résultats de la requête ParallelScan contiennent le jeton pour la page suivante. Vous pouvez utiliser ce jeton pour récupérer la page suivante. |
|
|
aliveTime |
Période de validité de la tâche d'analyse parallèle actuelle. Cette période de validité correspond également à celle du jeton. Unité : secondes. Valeur par défaut : 60. Nous vous recommandons d'utiliser la valeur par défaut. Si la requête suivante n'est pas lancée dans le délai de validité, aucune donnée supplémentaire ne peut être interrogée. Le temps de validité du jeton est actualisé à chaque envoi de requête. Remarque
Les sessions expirent prématurément si les index de basculement sont modifiés dynamiquement dans les schémas, si un serveur unique tombe en panne ou si un équilibrage de charge côté serveur est effectué. Dans ce cas, vous devez recréer les sessions. |
|
|
columnsToGet |
Nom de la colonne à renvoyer dans le résultat de regroupement. Vous pouvez ajouter le nom de la colonne à Columns. Si vous souhaitez que toutes les colonnes soient renvoyées dans l'index de recherche, vous pouvez utiliser l'opération ReturnAllFromIndex, plus concise. Important
ReturnAll ne peut pas être utilisé ici. |
|
|
sessionId |
ID de session de la tâche d'analyse parallèle. Vous pouvez appeler l'opération ComputeSplits pour créer une session et interroger le nombre maximal de tâches d'analyse parallèle prises en charge par la requête d'analyse parallèle. |
|
Exemple
Vous pouvez analyser les données en utilisant un seul thread ou plusieurs threads simultanément, selon vos besoins métier.
Analyse des données à l'aide d'un seul thread
Lorsque vous utilisez l'analyse parallèle, le code d'une requête utilisant un seul thread est plus simple que celui d'une requête utilisant plusieurs threads. Les paramètres currentParallelId et maxParallel ne sont pas requis pour une requête à thread unique. La requête ParallelScan à thread unique offre un débit supérieur à celui de la requête Search. Toutefois, son débit reste inférieur à celui d'une requête ParallelScan utilisant plusieurs threads.
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import com.alicloud.openservices.tablestore.SyncClient;
import com.alicloud.openservices.tablestore.model.ComputeSplitsRequest;
import com.alicloud.openservices.tablestore.model.ComputeSplitsResponse;
import com.alicloud.openservices.tablestore.model.Row;
import com.alicloud.openservices.tablestore.model.SearchIndexSplitsOptions;
import com.alicloud.openservices.tablestore.model.iterator.RowIterator;
import com.alicloud.openservices.tablestore.model.search.ParallelScanRequest;
import com.alicloud.openservices.tablestore.model.search.ParallelScanResponse;
import com.alicloud.openservices.tablestore.model.search.ScanQuery;
import com.alicloud.openservices.tablestore.model.search.SearchRequest.ColumnsToGet;
import com.alicloud.openservices.tablestore.model.search.query.MatchAllQuery;
import com.alicloud.openservices.tablestore.model.search.query.Query;
import com.alicloud.openservices.tablestore.model.search.query.QueryBuilders;
public class Test {
public static List<Row> scanQuery(final SyncClient client) {
String tableName = "<TableName>";
String indexName = "<IndexName>";
// Query the session ID and the maximum number of parallel scan tasks supported by the request.
ComputeSplitsRequest computeSplitsRequest = new ComputeSplitsRequest();
computeSplitsRequest.setTableName(tableName);
computeSplitsRequest.setSplitsOptions(new SearchIndexSplitsOptions(indexName));
ComputeSplitsResponse computeSplitsResponse = client.computeSplits(computeSplitsRequest);
byte[] sessionId = computeSplitsResponse.getSessionId();
int splitsSize = computeSplitsResponse.getSplitsSize();
/*
* Create a parallel scan request.
*/
ParallelScanRequest parallelScanRequest = new ParallelScanRequest();
parallelScanRequest.setTableName(tableName);
parallelScanRequest.setIndexName(indexName);
ScanQuery scanQuery = new ScanQuery();
// This query determines the range of the data to scan. You can create a nested and complex query.
Query query = new MatchAllQuery();
scanQuery.setQuery(query);
// Specify the maximum number of rows that can be returned by each ParallelScan call.
scanQuery.setLimit(2000);
parallelScanRequest.setScanQuery(scanQuery);
ColumnsToGet columnsToGet = new ColumnsToGet();
columnsToGet.setColumns(Arrays.asList("col_1", "col_2"));
parallelScanRequest.setColumnsToGet(columnsToGet);
parallelScanRequest.setSessionId(sessionId);
/*
* Use builder to create a parallel scan request that has the same features as the preceding request.
*/
ParallelScanRequest parallelScanRequestByBuilder = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(ScanQuery.newBuilder()
.query(QueryBuilders.matchAll())
.limit(2000)
.build())
.addColumnsToGet("col_1", "col_2")
.sessionId(sessionId)
.build();
List<Row> result = new ArrayList<>();
/*
* Use the native API operation to scan data.
*/
{
ParallelScanResponse parallelScanResponse = client.parallelScan(parallelScanRequest);
// Query the token of ScanQuery for the next request.
byte[] nextToken = parallelScanResponse.getNextToken();
// Obtain the data.
List<Row> rows = parallelScanResponse.getRows();
result.addAll(rows);
while (nextToken != null) {
// Specify the token.
parallelScanRequest.getScanQuery().setToken(nextToken);
// Continue to scan the data.
parallelScanResponse = client.parallelScan(parallelScanRequest);
// Obtain the data.
rows = parallelScanResponse.getRows();
result.addAll(rows);
nextToken = parallelScanResponse.getNextToken();
}
}
/*
* Recommended method.
* Use an iterator to scan all matched data. This method has the same query speed but is easier to use compared with the previous method.
*/
{
RowIterator iterator = client.createParallelScanIterator(parallelScanRequestByBuilder);
while (iterator.hasNext()) {
Row row = iterator.next();
result.add(row);
// Obtain the specific values.
String col_1 = row.getLatestColumn("col_1").getValue().asString();
long col_2 = row.getLatestColumn("col_2").getValue().asLong();
}
}
/*
* If the operation fails, retry the operation. If the caller of this function has a retry mechanism or if you do not want to retry the failed operation, you can ignore this part.
* To ensure availability, we recommend that you start a new parallel scan task when exceptions occur.
* The following exceptions may occur when you send a ParallelScan request:
* 1. A session exception occurs on the server side. The error code is OTSSessionExpired.
* 2. An exception such as a network exception occurs on the client side.
*/
try {
// Execute the processing logic.
{
RowIterator iterator = client.createParallelScanIterator(parallelScanRequestByBuilder);
while (iterator.hasNext()) {
Row row = iterator.next();
// Process rows of data. If you have sufficient memory resources, you can add the rows to a list.
result.add(row);
}
}
} catch (Exception ex) {
// Retry the processing logic.
{
result.clear();
RowIterator iterator = client.createParallelScanIterator(parallelScanRequestByBuilder);
while (iterator.hasNext()) {
Row row = iterator.next();
// Process rows of data. If you have sufficient memory resources, you can add the rows to a list.
result.add(row);
}
}
}
return result;
}
}
Analyse des données à l'aide de plusieurs threads
import com.alicloud.openservices.tablestore.SyncClient;
import com.alicloud.openservices.tablestore.model.ComputeSplitsRequest;
import com.alicloud.openservices.tablestore.model.ComputeSplitsResponse;
import com.alicloud.openservices.tablestore.model.Row;
import com.alicloud.openservices.tablestore.model.SearchIndexSplitsOptions;
import com.alicloud.openservices.tablestore.model.iterator.RowIterator;
import com.alicloud.openservices.tablestore.model.search.ParallelScanRequest;
import com.alicloud.openservices.tablestore.model.search.ScanQuery;
import com.alicloud.openservices.tablestore.model.search.query.QueryBuilders;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Semaphore;
import java.util.concurrent.atomic.AtomicLong;
public class Test {
public static void scanQueryWithMultiThread(final SyncClient client, String tableName, String indexName) throws InterruptedException {
// Query the number of CPU cores on the client.
final int cpuProcessors = Runtime.getRuntime().availableProcessors();
// Specify the number of parallel threads for the client. We recommend that you specify the number of CPU cores on the client as the number of parallel threads for the client to prevent impact on the query performance.
final Semaphore semaphore = new Semaphore(cpuProcessors);
// Query the session ID and the maximum number of parallel scan tasks supported by the request.
ComputeSplitsRequest computeSplitsRequest = new ComputeSplitsRequest();
computeSplitsRequest.setTableName(tableName);
computeSplitsRequest.setSplitsOptions(new SearchIndexSplitsOptions(indexName));
ComputeSplitsResponse computeSplitsResponse = client.computeSplits(computeSplitsRequest);
final byte[] sessionId = computeSplitsResponse.getSessionId();
final int maxParallel = computeSplitsResponse.getSplitsSize();
// Create an AtomicLong object if you need to obtain the row count for your business.
AtomicLong rowCount = new AtomicLong(0);
/*
* If you want to perform multithreading by using a function, you can build an internal class to inherit the threads.
* You can also build an external class to organize the code.
*/
final class ThreadForScanQuery extends Thread {
private final int currentParallelId;
private ThreadForScanQuery(int currentParallelId) {
this.currentParallelId = currentParallelId;
this.setName("ThreadForScanQuery:" + maxParallel + "-" + currentParallelId); // Specify the thread name.
}
@Override
public void run() {
System.out.println("start thread:" + this.getName());
try {
// Execute the processing logic.
{
ParallelScanRequest parallelScanRequest = ParallelScanRequest.newBuilder()
.tableName(tableName)
.indexName(indexName)
.scanQuery(ScanQuery.newBuilder()
.query(QueryBuilders.range("col_long").lessThan(10_0000)) // Specify the data to query.
.limit(2000)
.currentParallelId(currentParallelId)
.maxParallel(maxParallel)
.build())
.addColumnsToGet("col_long", "col_keyword", "col_bool") // Specify the fields to return from the search index. To return all fields from the search index, set returnAllColumnsFromIndex to true.
//.returnAllColumnsFromIndex(true)
.sessionId(sessionId)
.build();
// Use an iterator to obtain all the data.
RowIterator ltr = client.createParallelScanIterator(parallelScanRequest);
long count = 0;
while (ltr.hasNext()) {
Row row = ltr.next();
// Add a custom processing logic. The following sample code shows how to add a custom processing logic to count the number of rows:
count++;
}
rowCount.addAndGet(count);
System.out.println("thread[" + this.getName() + "] finished. this thread get rows:" + count);
}
} catch (Exception ex) {
// If exceptions occur, you can retry the processing logic.
} finally {
semaphore.release();
}
}
}
// Simultaneously execute threads. Valid values of currentParallelId: [0, Value of maxParallel).
List<ThreadForScanQuery> threadList = new ArrayList<ThreadForScanQuery>();
for (int currentParallelId = 0; currentParallelId < maxParallel; currentParallelId++) {
ThreadForScanQuery thread = new ThreadForScanQuery(currentParallelId);
threadList.add(thread);
}
// Simultaneously initiate the threads.
for (ThreadForScanQuery thread : threadList) {
// Specify a value for semaphore to limit the number of threads that can be initiated at the same time to prevent bottlenecks on the client.
semaphore.acquire();
thread.start();
}
// The main thread is blocked until all threads are complete.
for (ThreadForScanQuery thread : threadList) {
thread.join();
}
System.out.println("all thread finished! total rows:" + rowCount.get());
}
}