Utilisez le SDK Tablestore pour Go afin de diviser les données d'index de recherche et d'analyser les partitions en parallèle pour une exportation efficace à grande échelle.
Prérequis
Avant de commencer, préparez l'environnement comme suit :
Installez le SDK Tablestore pour Go et initialisez un client.
L'analyse parallèle nécessite le SDK Tablestore pour Go version 1.6.0 ou ultérieure. Nous vous recommandons d'utiliser la dernière version.
Description
L'exportation parallèle appelle d'abord ComputeSplits pour créer une session d'analyse et obtenir la concurrence recommandée. Chaque worker simultané appelle ensuite ParallelScan avec un CurrentParallelID unique et utilise NextToken pour lire sa partition en continu. L'analyse parallèle ne garantit pas l'ordre de l'ensemble des résultats et convient aux exportations volumineuses indépendantes de l'ordre des résultats. Vous pouvez également définir MaxParallel sur 1 et CurrentParallelID sur 0 pour une analyse à worker unique. Le code est alors plus simple et le débit est généralement supérieur à celui de Search, mais inférieur à celui d'une analyse multi-workers.
L'analyse parallèle ne prend pas en charge le tri ni l'agrégation. Utilisez Search si vous devez trier les résultats, effectuer des agrégations ou renvoyer les résultats de recherche aux utilisateurs finaux.
Les workers utilisant le même SessionId établissent un instantané des données lors du premier appel à ParallelScan. Les données insérées ou mises à jour pendant la tâche ne sont pas incluses dans cet instantané. Bien que SessionId soit facultatif, des modifications telles que l'équilibrage de charge côté serveur peuvent entraîner la duplication d'un petit nombre de lignes. Nous vous recommandons d'appeler d'abord ComputeSplits et d'inclure le SessionId renvoyé dans les requêtes suivantes.
Jusqu'à 10 tâches d'analyse parallèle peuvent s'exécuter simultanément sur le même index de recherche. Pour connaître les autres limites, consultez les Limites des index de recherche.
L'exemple suivant démarre plusieurs goroutines selon la concurrence recommandée par ComputeSplits, attribue un CurrentParallelID unique à chaque worker et lit toutes les partitions.
splits, err := client.ComputeSplits(
(&tablestore.ComputeSplitsRequest{}).
SetTableName("example_table").
SetSearchIndexSplitsOptions(tablestore.SearchIndexSplitsOptions{
IndexName: "example_index",
}),
)
if err != nil {
log.Fatal(err)
}
var waitGroup sync.WaitGroup
var mutex sync.Mutex
totalRows := 0
errors := make(chan error, splits.SplitsSize)
waitGroup.Add(int(splits.SplitsSize))
for workerID := int32(0); workerID < splits.SplitsSize; workerID++ {
currentWorkerID := workerID
go func() {
defer waitGroup.Done()
scanQuery := search.NewScanQuery().
SetQuery(&search.MatchAllQuery{}).
SetLimit(1000).
SetMaxParallel(splits.SplitsSize).
SetCurrentParallelID(currentWorkerID)
request := (&tablestore.ParallelScanRequest{}).
SetTableName("example_table").
SetIndexName("example_index").
SetScanQuery(scanQuery).
SetSessionId(splits.SessionId).
SetColumnsToGet(&tablestore.ColumnsToGet{
ReturnAllFromIndex: true,
})
for {
response, err := client.ParallelScan(request)
if err != nil {
errors <- err
return
}
// Process response.Rows here.
mutex.Lock()
totalRows += len(response.Rows)
mutex.Unlock()
if len(response.NextToken) == 0 {
return
}
request.SetScanQuery(scanQuery.SetToken(response.NextToken))
}
}()
}
waitGroup.Wait()
close(errors)
for err := range errors {
log.Fatal(err)
}
fmt.Println(totalRows)
Paramètres
Création d'une session d'analyse
|
Nom |
Type |
Description |
|
TableName (obligatoire) |
string |
Nom de la table de données. |
|
IndexName (obligatoire) |
string |
Nom de l'index de recherche. |
Requête d'analyse parallèle
|
Nom |
Type |
Description |
|
TableName (obligatoire) |
string |
Nom de la table de données. |
|
IndexName (obligatoire) |
string |
Nom de l'index de recherche. |
|
ScanQuery (obligatoire) |
search.ScanQuery |
Condition d'analyse et configuration de la concurrence. |
|
SessionId (facultatif) |
[]byte |
ID de session 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. |
|
ColumnsToGet (facultatif) |
*tablestore.ColumnsToGet |
Configuration des colonnes à renvoyer. Si ce paramètre est omis, seules les colonnes de clé primaire sont renvoyées. ParallelScan ne prend pas en charge ReturnAll. |
|
TimeoutMs (facultatif) |
*int32 |
Délai d'expiration de la requête en millisecondes. |
Configuration de l'analyse
|
Nom |
Type |
Description |
|
Query (obligatoire) |
search.Query |
Condition d'analyse. Tous les types de requête non vectoriels pris en charge par Search sont compatibles. |
|
Limit (facultatif) |
int32 |
Nombre maximal de lignes renvoyées par requête. Valeur par défaut : 2 000. Nous vous recommandons de conserver la valeur par défaut. Le serveur autorise une valeur allant jusqu'à 10 000, mais une valeur plus élevée augmente la latence des requêtes et l'utilisation des ressources. |
|
MaxParallel (facultatif) |
int32 |
Nombre total de workers. Valeur par défaut : 1. La valeur ne peut pas dépasser SplitsSize renvoyé par ComputeSplits. |
|
CurrentParallelID (facultatif) |
int32 |
ID du worker actuel. Ce paramètre est requis lorsque MaxParallel est supérieur à 1. Les ID de worker doivent être uniques et compris dans l'intervalle [0, MaxParallel). |
|
Token (facultatif) |
[]byte |
Valeur NextToken renvoyée par la réponse précédente. |
|
AliveTime (facultatif) |
int32 |
Intervalle maximal entre deux requêtes de pagination, en secondes. Valeurs valides : 1 à 600. Valeur par défaut : 60. Nous vous recommandons de conserver la valeur par défaut. Chaque requête de données réussie actualise la période de validité. |
Configuration des colonnes de retour
|
Nom |
Type |
Description |
|
Columns (facultatif) |
[]string |
Champs d'index de recherche à renvoyer. Un champ qui existe uniquement dans la table de données et qui n'est pas inclus dans l'index de recherche ne peut pas être renvoyé. |
|
ReturnAllFromIndex (facultatif) |
bool |
Indique s'il faut renvoyer tous les champs de l'index de recherche. Valeur par défaut : false. Si ce paramètre est défini sur true, vous n'avez pas besoin de spécifier Columns. |
|
ReturnAll (facultatif) |
bool |
L'analyse parallèle ne prend pas en charge ce paramètre. Ne le définissez pas sur true. |
Des modifications de schéma entraînant le basculement des index, un basculement de serveur ou un équilibrage de charge peuvent invalider prématurément une session et renvoyer l'erreur OTSSessionExpired. Des erreurs réseau côté client peuvent également interrompre une analyse. Si une telle erreur se produit, ignorez les résultats incomplets, appelez à nouveau ComputeSplits et redémarrez la tâche d'analyse complète depuis le début.
Réponse
Informations sur les partitions
|
Nom |
Type |
Description |
|
SessionId |
[]byte |
ID de session de la tâche utilisé pour analyser le même instantané de données. |
|
SplitsSize |
int32 |
Concurrence maximale prise en charge pour l'index de recherche. |
Résultat de l'analyse
|
Nom |
Type |
Description |
|
Rows |
[]*tablestore.Row |
Lignes renvoyées par l'analyse actuelle. |
|
NextToken |
[]byte |
Jeton pour la page suivante. Continuez l'analyse de la partition actuelle si la valeur n'est pas vide. |