DISTRIBUTED MAPJOIN est une version optimisée de MAPJOIN. Utilisez cette fonctionnalité pour joindre une petite table à une grande table. DISTRIBUTED MAPJOIN et MAPJOIN réduisent les opérations de brassage (shuffling) et de tri sur la grande table.
Notes d'utilisation
Les tables à joindre doivent avoir des tailles différentes. La grande table doit dépasser 10 To, tandis que la petite table doit avoir une taille comprise entre 1 Go et 100 Go.
Les données de la petite table doivent être réparties uniformément. Si la petite table présente des déséquilibres importants (longues traînes), un volume excessif de données peut s'accumuler dans un seul shard. Cela risque de provoquer une erreur d'épuisement de la mémoire (OOM) ou un dépassement du délai d'attente des appels de procédure distante (RPC).
Si l'exécution d'une tâche SQL dure plus de 20 minutes, utilisez DISTRIBUTED MAPJOIN pour l'optimiser.
-
L'exécution d'une tâche consomme des ressources importantes. Évitez d'exécuter ce type de tâche dans un petit groupe de quotas.
RemarqueModifiez le groupe de quotas sur la page Quotas. Pour plus d'informations, consultez la rubrique Gérer les quotas dans la nouvelle console MaxCompute.
Utiliser DISTRIBUTED MAPJOIN
Pour utiliser DISTRIBUTED MAPJOIN dans une instruction SELECT, ajoutez l'indicateur (hint) /*+distmapjoin(<table_name>(shard_count=<n>,replica_count=<m>))*/ à l'instruction. Les paramètres shard_count et replica_count déterminent le parallélisme des tâches. Calculez le parallélisme avec la formule suivante : Parallelism = shard_count × replica_count.
-
Paramètres
table_name : nom de la petite table à joindre.
-
shard_count=<n> : nombre de shards de données de la petite table à joindre. Les shards de la petite table sont distribués sur chaque nœud de calcul pour le traitement des données. Le paramètre n spécifie le nombre de shards. En général, définissez ce paramètre sur un nombre impair.
RemarqueSpécifiez manuellement le paramètre shard_count. Estimez sa valeur en fonction de la taille de la petite table. Le volume de données traité par un seul nœud de shard doit se situer entre 200 Mo et 500 Mo.
Une valeur trop élevée pour le paramètre shard_count dégrade les performances et la stabilité du traitement des données. À l'inverse, une valeur trop faible risque de provoquer une erreur due à une consommation excessive de mémoire.
-
replica_count=<m> : nombre de réplicas de la petite table. Le paramètre m spécifie le nombre de réplicas. Valeur par défaut : 1.
RemarquePour limiter le nombre excessif de requêtes d'accès et éviter qu'une défaillance d'un nœud unique n'entraîne l'échec de toute la tâche, créez plusieurs réplicas des données au sein d'un même shard. Si un nœud redémarre fréquemment en raison d'un parallélisme élevé ou de performances environnementales instables, augmentez la valeur du paramètre replica_count. Définissez ce paramètre sur 2 ou 3.
-
Syntaxe
-- Recommended: Specify shard_count. replica_count defaults to 1. /*+distmapjoin(a(shard_count=5))*/ -- Recommended: Specify both shard_count and replica_count. /*+distmapjoin(a(shard_count=5,replica_count=2))*/ -- Use DISTRIBUTED MAPJOIN for multiple small tables. /*+distmapjoin(a(shard_count=5,replica_count=2),b(shard_count=5,replica_count=2)) */ -- Use DISTRIBUTED MAPJOIN and MAPJOIN together. /*+distmapjoin(a(shard_count=5,replica_count=2)),mapjoin(b)*/
Exemples
Cet exemple illustre l'utilisation de DISTRIBUTED MAPJOIN lors de l'insertion de données dans une table partitionnée nommée tmall_dump_lasttable.
-
Syntaxe JOIN
insert OVERWRITE table tmall_dump_lasttable partition(ds='20211130') select t1.* from ( select nid, doc,type from search_ods.dump_lasttable where ds='20211203' )t1 join ( select distinct item_id from tbcdm.dim_tb_itm where ds='20211130' and bc_type='B' and is_online='Y' )t2 on t1.nid=t2.item_id; -
Syntaxe DISTRIBUTED MAPJOIN
insert OVERWRITE table tmall_dump_lasttable partition (ds='20211130') select /*+ distmapjoin(t2(shard_count=35)) */ t1.* from ( select nid, doc, type from search_ods.dump_lasttable where ds='20211203' )t1 join ( select distinct item_id from tbcdm.dim_tb_itm where ds='20211130' and bc_type='B' and is_online='Y' )t2 on t1.nid=t2.item_id;