DataProc集群模式下Kryo序列化遇ClassNotFoundException求助
问题描述
在Google Cloud DataProc上使用KryoSerializer对DataFrame进行序列化时,抛出ClassNotFoundException,缺失的类属于spark-submit提交的主JAR包。该问题仅在设置spark.submit.deployMode=cluster且master=yarn时出现,使用local[*]作为master时运行正常,推测是Executor的类路径中未包含主JAR。
已尝试的Spark配置
已尝试将主JAR添加到以下Spark配置的类路径中,但问题依旧(根据相关资料,此操作本不应必要):
spark.executor.extraClassPath spark.driver.extraClassPath spark.driver.userClassPathFirst spark.executor.userClassPathFirst spark.yarn.dist.jars spark.yarn.jars spark.jars
环境信息
使用DataProc 2.1最新镜像(具体URI:https://www.googleapis.com/compute/v1/projects/cloud-dataproc/global/images/dataproc-2-1-ubu20-20221201-035100-rc01),也尝试过其他2.1版本的Debian镜像,问题一致。获取稳定2.1镜像URI的gcloud命令:
gcloud compute images list --uri --project cloud-dataproc --filter "labels.goog-dataproc-version ~ ^2.1.0" --sort-by=~creationTimestamp
作业提交代码
主JAR存储在GCS,使用GCS云存储连接器,非集群模式下可正常运行:
SparkJob sparkJob = SparkJob.newBuilder().setMainJarFileUri(MAIN_JAR)
异常堆栈信息
22/12/21 03:57:18 WARN TaskSetManager: Lost task 0.0 in stage 0.0 (TID 0) (
executor 1): java.io.IOException: org.apache.spark.SparkException: Failed to register classes with Kryo
at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1477)
at org.apache.spark.broadcast.TorrentBroadcast.readBroadcastBlock(TorrentBroadcast.scala:228)
at org.apache.spark.broadcast.TorrentBroadcast.getValue(TorrentBroadcast.scala:105)
at org.apache.spark.broadcast.Broadcast.value(Broadcast.scala:70)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:84)
at org.apache.spark.scheduler.Task.run(Task.scala:136)
at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: org.apache.spark.SparkException: Failed to register classes with Kryo
at org.apache.spark.serializer.KryoSerializer.$anonfun$newKryo$5(KryoSerializer.scala:183)
at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
at org.apache.spark.util.Utils$.withContextClassLoader(Utils.scala:233)
at org.apache.spark.serializer.KryoSerializer.newKryo(KryoSerializer.scala:171)
at org.apache.spark.serializer.KryoSerializer$$anon$1.create(KryoSerializer.scala:102)
at com.esotericsoftware.kryo.pool.KryoPoolQueueImpl.borrow(KryoPoolQueueImpl.java:48)
at org.apache.spark.serializer.KryoSerializer$PoolWrapper.borrow(KryoSerializer.scala:109)
at org.apache.spark.serializer.KryoSerializerInstance.borrowKryo(KryoSerializer.scala:346)
at org.apache.spark.serializer.KryoDeserializationStream.(KryoSerializer.scala:302)
at org.apache.spark.serializer.KryoSerializerInstance.deserializeStream(KryoSerializer.scala:436)
at org.apache.spark.broadcast.TorrentBroadcast$.unBlockifyObject(TorrentBroadcast.scala:336)
at org.apache.spark.broadcast.TorrentBroadcast.$anonfun$readBroadcastBlock$4(TorrentBroadcast.scala:259)
at scala.Option.getOrElse(Option.scala:189)
at org.apache.spark.broadcast.TorrentBroadcast.$anonfun$readBroadcastBlock$2(TorrentBroadcast.scala:233)
at org.apache.spark.util.KeyLock.withLock(KeyLock.scala:64)
at org.apache.spark.broadcast.TorrentBroadcast.$anonfun$readBroadcastBlock$1(TorrentBroadcast.scala:228)
at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1470)
... 11 more
Caused by: java.lang.ClassNotFoundException: <.MyClass In the Main JAR sent to spark-submit>
at java.base/java.lang.ClassLoader.findClass(ClassLoader.java:719)
at org.apache.spark.util.ParentClassLoader.findClass(ParentClassLoader.java:35)
at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:589)
at org.apache.spark.util.ParentClassLoader.loadClass(ParentClassLoader.java:40)
at org.apache.spark.util.ChildFirstURLClassLoader.loadClass(ChildFirstURLClassLoader.java:48)
at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:522)
at java.base/java.lang.Class.forName0(Native Method)
at java.base/java.lang.Class.forName(Class.java:398)
at org.apache.spark.util.Utils$.classForName(Utils.scala:220)
at org.apache.spark.serializer.KryoSerializer.$anonfun$newKryo$6(KryoSerializer.scala:174)
at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
at org.apache.spark.serializer.KryoSerializer.$anonfun$newKryo$5(KryoSerializer.scala:173)
... 27 more
POM依赖配置
集群已安装的依赖均设置为provided scope:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.12</artifactId> <version>3.3.0</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.12</artifactId> <version>3.3.0</version> <scope>provided</scope> </dependency> <dependency> <groupId>com.google.cloud.spark</groupId> <artifactId>spark-bigquery-with-dependencies_2.12</artifactId> <version>0.27.1</version> <scope>provided</scope> </dependency>
额外尝试
还尝试通过初始化动作,使用gsutil手动将JAR复制到Executor的类路径中,但问题仍未解决——按正常逻辑,此操作本不应必要。
内容的提问来源于stack exchange,提问作者Mark Ballenger

