このトピックでは、Spark の使用に関するよくある質問にお答えします。
-
Spark Core
-
Spark SQL
-
PySpark
-
Spark Streaming
-
spark-submit
過去の Spark ジョブの確認方法
EMR コンソールで、対象クラスターの Access Links and Ports タブに移動し、Spark UI リンクをクリックすると、過去の Spark ジョブを確認できます。コンポーネントの Web UI へのアクセス方法の詳細については、「EMR コンソールでのオープンソースコンポーネントの Web UI へのアクセス」をご参照ください。
スタンドアロンモードでの Spark ジョブのサブミット
いいえ。E-MapReduce (EMR) は、Spark on YARN および Spark on Kubernetes を介したジョブのサブミットのみをサポートしています。スタンドアロンモードと Mesos モードはサポートされていません。
Spark 2 CLI のログ出力の削減方法
デフォルトでは、EMR DataLake クラスター上の spark-sql や spark-shell などのコマンドラインインターフェイス (CLI) ツールは、INFO レベルのログを出力します。
-
CLI ツールを実行するノード (例えば、マスターノード) で、log4j.properties 設定ファイルを作成します。デフォルトの設定ファイルをコピーするには、次のコマンドを使用します。
cp /etc/emr/spark-conf/log4j.properties /new/path/to/log4j.properties -
新しい設定ファイルでログレベルを変更します。
log4j.rootCategory=WARN, console -
Spark サービスの spark-defaults.conf ファイルで、spark.driver.extraJavaOptions プロパティを更新します。-Dlog4j.configuration=file:/etc/emr/spark-conf/log4j.properties を -Dlog4j.configuration=file:/new/path/to/log4j.properties に置き換えます。
重要パスにはプレフィックスとして file: を付ける必要があります。
Spark 3 の small ファイルマージの使用方法
spark.sql.adaptive.merge.output.small.files.enabled パラメーターを true に設定すると、small ファイルが自動的にマージされます。マージされたファイルは圧縮されます。マージされたファイルが小さすぎる場合は、spark.sql.adaptive.advisoryOutputFileSizeInBytes パラメーターの値を増やしてください。デフォルト値は 256 MB です。
Spark SQL でのデータスキューの処理方法
-
Spark 2 の場合は、次のいずれかのアプローチを使用します。
-
テーブルを読み取る際に、null 値などの無関係なデータを除外します。
-
小さい方のテーブルをブロードキャストします。
select /*+ BROADCAST (table1) */ * from table1 join table2 on table1.id = table2.id -
スキューしたキーに基づいて、スキューしたデータを分離します。
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 -
スキューしたキーが既知の場合は、データを分散させます。
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) -
スキューしたキーが不明な場合は、データを分散させます。
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
-
-
Spark 3 の場合は、EMR コンソールの Spark 3 サービスの Configure タブに移動し、spark.sql.adaptive.enabled パラメーターと spark.sql.adaptive.skewJoin.enabled パラメーターを true に設定します。
PySpark で Python 3 を指定する方法
このセクションでは、Spark 2 を搭載した EMR V5.7.0 の DataLake クラスターを例に、PySpark で Python 3 を指定する方法を説明します。
Python のバージョンを変更するには、2 つの方法があります。
一時的な方法
-
SSH 経由でクラスターにログインします。詳細については、「クラスターへのログイン」をご参照ください。
-
次のコマンドを実行して、Python のバージョンを変更します。
export PYSPARK_PYTHON=/usr/bin/python3 -
次のコマンドを実行して、Python のバージョンを確認します。
pyspark出力に次のメッセージが含まれている場合、Python のバージョンは Python 3 に変更されています。
Using Python version 3.6.8
恒久的な方法
-
SSH 経由でクラスターにログインします。詳細については、「クラスターへのログイン」をご参照ください。
-
設定ファイルを変更します。
-
次のコマンドを実行して profile ファイルを開きます。
vi /etc/profile -
iを押して挿入モードに入ります。 -
profile ファイルの末尾に次の行を追加します。
export PYSPARK_PYTHON=/usr/bin/python3 -
Escを押して挿入モードを終了します。次に、:wqと入力してファイルを保存し、閉じます。
-
-
次のコマンドを実行して設定ファイルを再読み込みし、変更をすぐに適用します。
source /etc/profile -
次のコマンドを実行して、Python のバージョンを確認します。
pyspark出力に次のメッセージが含まれている場合、Python のバージョンは Python 3 に変更されています。
Using Python version 3.6.8
Spark Streaming ジョブが予期せず停止する理由
-
ご利用の Spark のバージョンが 1.6 より前の場合は、アップグレードしてください。
1.6 より前のバージョンの Spark にはメモリリークがあり、コンテナが終了する原因となる可能性があります。
-
コードがメモリ使用量に対して最適化されていることを確認してください。
停止したジョブが実行中と表示される理由
この問題は、ジョブを YARN-client モードでサブミットした場合に発生する可能性があります。このモードでは、EMR がジョブのステータスを正確に監視できないためです。正しいステータスレポートを保証するには、代わりにジョブを YARN-cluster モードでサブミットしてください。
java.lang.ClassNotFoundException エラー:spark-submit を YARN-cluster モードで実行時
以下はエラーメッセージの例です。
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)
原因:Kerberos が有効な EMR クラスターでは、ジョブが YARN-cluster モードで実行される際に、ドライバーのクラスパスに必要な JAR が自動的に入力されません。これにより、ClassNotFoundException エラーが発生します。
解決策:Kerberos が有効な EMR クラスターでは、spark-submit を使用して YARN-cluster モードでジョブをサブミットする際に、--jars パラメーターを使用する必要があります。アプリケーションの JAR に加えて、/opt/apps/METASTORE/metastore-current/hive2 ディレクトリのすべての JAR パッケージも含める必要があります。
YARN-cluster モードでは、--jars パラメーター内のすべてのファイルパスはカンマで区切る必要があります。ディレクトリはサポートされていません。
たとえば、アプリケーションの JAR が /opt/apps/SPARK3/spark3-current/examples/jars/spark-examples_2.12-3.5.3-emr.jar の場合、次の spark-submit コマンドを実行します。
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