Cette rubrique répond aux questions fréquentes sur l'utilisation de Spark.
-
Spark Core
-
Spark SQL
-
PySpark
-
Spark Streaming
-
spark-submit
Où consulter l'historique des tâches Spark ?
Dans la console EMR, accédez à l'onglet Access Links and Ports du cluster cible, puis cliquez sur le lien Spark UI pour consulter l'historique des tâches Spark. Pour plus d'informations sur l'accès aux interfaces utilisateur des composants, consultez la rubrique Accéder aux interfaces web des composants open source dans la console EMR.
Soumission de tâches Spark en mode standalone possible ?
Non. E-MapReduce prend uniquement en charge la soumission de tâches via Spark on YARN et Spark on Kubernetes. Les modes standalone et Mesos ne sont pas pris en charge.
Réduction de la verbosité des journaux des outils CLI Spark 2
Par défaut, les outils d'interface de ligne de commande (CLI) tels que spark-sql et spark-shell sur un cluster DataLake EMR génèrent des journaux de niveau INFO.
-
Sur le nœud où vous exécutez les outils CLI (par exemple, le nœud maître), créez un fichier de configuration log4j.properties. Pour copier le fichier de configuration par défaut, utilisez la commande suivante :
cp /etc/emr/spark-conf/log4j.properties /new/path/to/log4j.properties -
Modifiez le niveau de journalisation dans le nouveau fichier de configuration.
log4j.rootCategory=WARN, console -
Dans le fichier spark-defaults.conf du service Spark, mettez à jour la propriété spark.driver.extraJavaOptions. Remplacez -Dlog4j.configuration=file:/etc/emr/spark-conf/log4j.properties par -Dlog4j.configuration=file:/new/path/to/log4j.properties.
ImportantLe chemin doit être préfixé par file:.
Utilisation de la fusion des petits fichiers dans Spark 3
Définissez le paramètre spark.sql.adaptive.merge.output.small.files.enabled sur true pour fusionner automatiquement les petits fichiers. Les fichiers fusionnés sont compressés. Si les fichiers résultants sont encore trop petits, augmentez la valeur du paramètre spark.sql.adaptive.advisoryOutputFileSizeInBytes. La valeur par défaut est de 256 Mo.
Gestion du déséquilibre des données (data skew) dans Spark SQL
-
Pour Spark 2, utilisez l'une des approches suivantes :
Filtrez les données non pertinentes, telles que les valeurs nulles, lors de la lecture d'une table.
-
Diffusez (broadcast) la table la plus petite.
select /*+ BROADCAST (table1) */ * from table1 join table2 on table1.id = table2.id -
Séparez les données déséquilibrées en fonction de la clé concernée.
select * from table1_1 join table2 on table1_1.id = table2.id union all select /*+ BROADCAST (table1_2) */ * from table1_2 join table2 on table1_2.id = table2.id -
Si la clé déséquilibrée est connue, dispersez les données.
select id, value, concat(id, (rand() * 10000) % 3) as new_id from A select id, value, concat(id, suffix) as new_id from ( select id, value, suffix from B Lateral View explode(array(0, 1, 2)) tmp as suffix) -
Si la clé déséquilibrée est inconnue, dispersez les données.
select t1.id, t1.id_rand, t2.name from ( select id , case when id = null then concat('SkewData_', cast(rand() as string)) else id end as id_rand from test1 where statis_date = '20221130') t1 left join test2 t2 on t1.id_rand = t2.id
Pour Spark 3, accédez à l'onglet Configure du service Spark 3 dans la console EMR, puis définissez les paramètres spark.sql.adaptive.enabled et spark.sql.adaptive.skewJoin.enabled sur true.
Spécification de Python 3 pour PySpark
Cette section explique comment spécifier Python 3 pour PySpark, en prenant pour exemple un cluster DataLake sur EMR V5.7.0 avec Spark 2.
Il existe deux méthodes pour modifier la version de Python :
Méthode temporaire
Connectez-vous au cluster via SSH. Pour plus d'informations, consultez la rubrique Se connecter à un cluster.
-
Exécutez la commande suivante pour modifier la version de Python :
export PYSPARK_PYTHON=/usr/bin/python3 -
Exécutez la commande suivante pour vérifier la version de Python :
pysparkSi la sortie contient le message suivant, la version de Python a bien été modifiée en Python 3.
Using Python version 3.6.8
Méthode permanente
Connectez-vous au cluster via SSH. Pour plus d'informations, consultez la rubrique Se connecter à un cluster.
-
Modifiez le fichier de configuration.
-
Exécutez la commande suivante pour ouvrir le fichier profile :
vi /etc/profile Appuyez sur
ipour passer en mode insertion.-
À la fin du fichier profile, ajoutez la ligne suivante :
export PYSPARK_PYTHON=/usr/bin/python3 Appuyez sur
Escpour quitter le mode insertion. Saisissez ensuite:wqpour enregistrer et fermer le fichier.
-
-
Exécutez la commande suivante pour recharger le fichier de configuration et appliquer immédiatement les modifications :
source /etc/profile -
Exécutez la commande suivante pour vérifier la version de Python :
pysparkSi la sortie contient le message suivant, la version de Python a bien été modifiée en Python 3.
Using Python version 3.6.8
Arrêt inattendu des tâches Spark Streaming
-
Si votre version de Spark est antérieure à la 1,6, effectuez une mise à niveau.
Les versions de Spark antérieures à la 1,6 présentent une fuite de mémoire pouvant entraîner la terminaison des conteneurs.
Assurez-vous que votre code est optimisé en termes d'utilisation de la mémoire.
Affichage d'une tâche arrêtée comme Running
Ce phénomène peut se produire si vous avez soumis la tâche en mode YARN-client, car E-MapReduce ne peut pas surveiller avec précision l'état des tâches dans ce mode. Pour garantir un rapport d'état correct, soumettez plutôt la tâche en mode YARN-cluster.
Erreur java.lang.ClassNotFoundException avec spark-submit en mode YARN-cluster
Voici un exemple de message d'erreur :
Process Output>>> 24/09/30 15:41:24 WARN HiveConf: HiveConf of name hive.metastore.type does not exist
Process Output>>> 24/09/30 15:41:24 ERROR Hive: Unable to instantiate a metastore client factory com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory: java.lang.ClassNotFoundException: Class com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory not found
Process Output>>> java.lang.ClassNotFoundException: Class com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory not found
Process Output>>> at org.apache.hadoop.conf.Configuration.getClassByName(Configuration.java:2542)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.createMetaStoreClient(Hive.java:3711)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.getMSC(Hive.java:3794)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.getMSC(Hive.java:3774)
Process Output>>> at org.apache.hadoop.hive.ql.metadata.Hive.getDelegationToken(Hive.java:3924)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider.$anonfun$obtainDelegationTokens$4(HiveDelegationTokenProvider.scala:104)
Process Output>>> at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider$$anon$1.run(HiveDelegationTokenProvider.scala:139)
Process Output>>> at java.security.AccessController.doPrivileged(Native Method)
Process Output>>> at javax.security.auth.Subject.doAs(Subject.java:422)
Process Output>>> at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1730)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider.doAsRealUser(HiveDelegationTokenProvider.scala:138)
Process Output>>> at org.apache.spark.sql.hive.security.HiveDelegationTokenProvider.obtainDelegationTokens(HiveDelegationTokenProvider.scala:102)
Process Output>>> at org.apache.spark.deploy.security.HadoopDelegationTokenManager.$anonfun$obtainDelegationTokens$2(HadoopDelegationTokenManager.scala:164)
Process Output>>> at scala.collection.TraversableLike.$anonfun$flatMap$1(TraversableLike.scala:293)
Process Output>>> at scala.collection.Iterator.foreach(Iterator.scala:943)
Process Output>>> at scala.collection.Iterator.foreach$(Iterator.scala:943)
Process Output>>> at scala.collection.AbstractIterator.foreach(Iterator.scala:1431)
Process Output>>> at scala.collection.MapLike$DefaultValuesIterable.foreach(MapLike.scala:214)
Process Output>>> at scala.collection.TraversableLike.flatMap(TraversableLike.scala:293)
Process Output>>> at scala.collection.TraversableLike.flatMap$(TraversableLike.scala:290)
Cause : Sur un cluster EMR activé avec Kerberos, le classpath du driver n'est pas automatiquement renseigné avec les JAR requis lors de l'exécution d'une tâche en mode YARN-cluster. Cela provoque l'erreur ClassNotFoundException.
Solution : Sur un cluster EMR activé avec Kerberos, vous devez utiliser le paramètre --jars lors de la soumission d'une tâche avec spark-submit en mode YARN-cluster. En plus des JAR de votre application, vous devez inclure tous les packages JAR situés dans le répertoire /opt/apps/METASTORE/metastore-current/hive2.
En mode YARN-cluster, tous les chemins de fichiers dans le paramètre --jars doivent être séparés par des virgules. Les répertoires ne sont pas pris en charge.
Par exemple, si le JAR de votre application est /opt/apps/SPARK3/spark3-current/examples/jars/spark-examples_2.12-3.5.3-emr.jar, exécutez la commande spark-submit suivante :
spark-submit --deploy-mode cluster --class org.apache.spark.examples.SparkPi --master yarn \
--jars $(ls /opt/apps/METASTORE/metastore-current/hive2/*.jar | tr '\n' ',') \
/opt/apps/SPARK3/spark3-current/examples/jars/spark-examples_2.12-3.5.3-emr.jar