EMR Remote Shuffle Service (ESS) est une extension d'E-MapReduce (EMR) qui optimise les opérations de shuffle dans les moteurs de calcul.
Contexte
Les mécanismes de shuffle traditionnels présentent plusieurs défis :
Dans les scénarios impliquant de grands volumes de données, les opérations d'écriture de shuffle peuvent entraîner un débordement des données sur le disque, provoquant une amplification des écritures.
Les opérations de lecture de shuffle génèrent un grand nombre de petits paquets réseau, ce qui peut provoquer des erreurs de réinitialisation de connexion.
Les opérations de lecture de shuffle impliquent de nombreuses petites requêtes d'E/S et des lectures aléatoires, ce qui impose une charge importante aux disques et aux CPU.
Lorsque le nombre de mappeurs (M) et de réducteurs (N) atteint plusieurs milliers, le nombre total de connexions réseau (M × N) peut empêcher la fin d'une tâche.
Le NodeManager et le Spark Shuffle Service s'exécutent dans le même processus. Lorsque les volumes de données de shuffle sont extrêmement importants, le NodeManager peut redémarrer, ce qui affecte la stabilité de la planification YARN.
ESS offre les avantages suivants :
Il utilise un mécanisme de shuffle de type push plutôt que pull, ce qui réduit la pression sur la mémoire des mappeurs.
Il prend en charge l'agrégation des E/S, réduisant ainsi le nombre de connexions de lecture de shuffle de M × N à N et remplaçant les lectures aléatoires par des lectures séquentielles.
Il prend en charge un mécanisme à deux réplicas pour réduire la probabilité d'échecs de récupération.
Il prend en charge une architecture de séparation du calcul et du stockage, vous permettant de déployer le Shuffle Service dans un environnement matériel distinct, découplé du cluster de calcul.
Il élimine la dépendance aux disques locaux lors de l'exécution de Spark sur Kubernetes.
La figure suivante illustre l'architecture d'ESS.
Limitations
Cette rubrique s'applique uniquement aux versions d'EMR antérieures à EMR-3.39.1, aux versions de la série EMR-4.x et aux versions antérieures à EMR-5.5.0. Pour EMR-3.39.1 ou ultérieur et EMR-5.5.0 ou ultérieur, consultez RSS.
Créer un cluster
Par exemple, sur EMR-4.5.0, vous pouvez créer un cluster avec ESS de deux manières :
Créez un cluster E-MapReduce Shuffle Service. Sur la page Software Configuration, pour Cluster Type, sélectionnez Shuffle Service. Le service requis est ESS (1.0.0).
Créez un cluster E-MapReduce Hadoop. Sur la page Software Configuration, dans la section Cluster Type, sélectionnez un type tel que Hadoop, Kafka ou Druid. Configurez ensuite les Cloud Native Options (par exemple, sur ECS) et la Product Version (par exemple, EMR-4.5.0). La page affiche alors les services requis correspondants (tels que HDFS, YARN et Spark) et les services facultatifs (tels que ESS, HBase et Flink) avec leurs versions.
Pour plus d'informations sur la création d'un cluster, consultez Créer un cluster.
Utilisation d'ESS
Pour utiliser ESS avec Spark, ajoutez les paramètres suivants à la soumission de votre tâche Spark. Pour plus d'informations sur la configuration des paramètres, consultez Modifier les tâches.
Pour plus d'informations sur les paramètres Spark, consultez Configuration Spark.
|
Paramètre |
Description |
|
spark.shuffle.manager |
La valeur doit être org.apache.spark.shuffle.ess.EssShuffleManager. |
|
spark.ess.master.address |
Spécifiez l'adresse au format <ess-master-ip>:<ess-master-port>. Les paramètres sont les suivants :
|
|
spark.shuffle.service.enabled |
Définissez la valeur sur Vous devez désactiver le service externe de shuffle par défaut pour utiliser EMR Remote Shuffle Service. |
|
spark.shuffle.useOldFetchProtocol |
Définissez la valeur sur Cela permet la compatibilité avec l'ancien protocole de shuffle. |
|
spark.sql.adaptive.enabled |
Définissez la valeur sur EMR Remote Shuffle Service ne prend pas en charge l'exécution adaptative. |
|
spark.sql.adaptive.skewJoin.enabled |
Paramètres
La page de configuration du service ESS répertorie tous les paramètres ESS.
|
Paramètre |
Description |
Valeur par défaut |
|
ess.push.data.replicate |
Active ou désactive la fonctionnalité à deux réplicas. Valeurs possibles :
Remarque
Nous vous recommandons d'activer cette fonctionnalité dans les environnements de production. |
true |
|
ess.worker.flush.queue.capacity |
Le nombre de tampons de vidage par répertoire. Remarque
Pour améliorer les performances, vous pouvez configurer plusieurs disques. Pour un débit de lecture et d'écriture optimal, nous vous recommandons de ne pas utiliser plus de deux répertoires par disque. La mémoire heap consommée par le tampon de vidage pour chaque répertoire est ess.worker.flush.buffer.size ess.worker.flush.queue.capacity, soit |
512 |
|
ess.flush.timeout |
Le délai d'expiration pour le vidage des données vers la couche de stockage. |
240s |
|
ess.application.timeout |
Le délai d'expiration du heartbeat de l'application. Si aucun heartbeat n'est reçu pendant cette période, ESS nettoie les ressources de l'application. |
240s |
|
ess.worker.flush.buffer.size |
La taille du tampon de vidage. Lorsque le tampon dépasse cette taille, ESS vide les données sur le disque. |
256k |
|
ess.metrics.system.enable |
Active ou désactive la surveillance. Valeurs possibles :
|
false |
|
ess_worker_offheap_memory |
La taille de la mémoire off-heap pour un nœud core. |
4g |
|
ess_worker_memory |
La taille de la mémoire heap pour un nœud core. |
4g |
|
ess_master_memory |
La taille de la mémoire heap pour le nœud maître. |
4g |