【问题标题】:ClassNotFoundException: org.apache.beam.runners.spark.io.SourceRDD$SourcePartition during spark submitClassNotFoundException:org.apache.beam.runners.spark.io.SourceRDD$SourcePartition 在火花提交期间
【发布时间】:2022-12-09 13:23:28
【问题描述】:

我使用 spark-submit 来激发独立集群来执行我的着色 jar,但是执行程序出现错误:

22/12/06 15:21:25 INFO TaskSetManager: Starting task 0.1 in stage 0.0 (TID 1) (10.37.2.77, executor 0, partition 0, PROCESS_LOCAL, 5133 bytes) taskResourceAssignments Map()
22/12/06 15:21:25 INFO TaskSetManager: Lost task 0.1 in stage 0.0 (TID 1) on 10.37.2.77, executor 0: java.lang.ClassNotFoundException (org.apache.beam.runners.spark.io.SourceRDD$SourcePartition) [duplicate 1]
22/12/06 15:21:25 INFO TaskSetManager: Starting task 0.2 in stage 0.0 (TID 2) (10.37.2.77, executor 0, partition 0, PROCESS_LOCAL, 5133 bytes) taskResourceAssignments Map()
22/12/06 15:21:25 INFO TaskSetManager: Lost task 0.2 in stage 0.0 (TID 2) on 10.37.2.77, executor 0: java.lang.ClassNotFoundException (org.apache.beam.runners.spark.io.SourceRDD$SourcePartition) [duplicate 2]
22/12/06 15:21:25 INFO TaskSetManager: Starting task 0.3 in stage 0.0 (TID 3) (10.37.2.77, executor 0, partition 0, PROCESS_LOCAL, 5133 bytes) taskResourceAssignments Map()
22/12/06 15:21:25 INFO TaskSetManager: Lost task 0.3 in stage 0.0 (TID 3) on 10.37.2.77, executor 0: java.lang.ClassNotFoundException (org.apache.beam.runners.spark.io.SourceRDD$SourcePartition) [duplicate 3]
22/12/06 15:21:25 ERROR TaskSetManager: Task 0 in stage 0.0 failed 4 times; aborting job
22/12/06 15:21:25 INFO TaskSchedulerImpl: Removed TaskSet 0.0, whose tasks have all completed, from pool 
22/12/06 15:21:25 INFO TaskSchedulerImpl: Cancelling stage 0
22/12/06 15:21:25 INFO TaskSchedulerImpl: Killing all running tasks in stage 0: Stage cancelled
22/12/06 15:21:25 INFO DAGScheduler: ResultStage 0 (collect at BoundedDataset.java:96) failed in 1.380 s due to Job aborted due to stage failure: Task 0 in stage 0.0 failed 4 times, most recent failure: Lost task 0.3 in stage 0.0 (TID 3) (10.37.2.77 executor 0): java.lang.ClassNotFoundException: org.apache.beam.runners.spark.io.SourceRDD$SourcePartition
    at java.lang.ClassLoader.findClass(ClassLoader.java:523)
    at org.apache.spark.util.ParentClassLoader.findClass(ParentClassLoader.java:35)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:418)
    at org.apache.spark.util.ParentClassLoader.loadClass(ParentClassLoader.java:40)
    at org.apache.spark.util.ChildFirstURLClassLoader.loadClass(ChildFirstURLClassLoader.java:48)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:351)
    at java.lang.Class.forName0(Native Method)
    at java.lang.Class.forName(Class.java:348)
    at org.apache.spark.serializer.JavaDeserializationStream$$anon$1.resolveClass(JavaSerializer.scala:68)
    at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1988)
    at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1852)
    at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2186)
    at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1669)
    at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2431)
    at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2355)
    at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2213)
    at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1669)
    at java.io.ObjectInputStream.readObject(ObjectInputStream.java:503)
    at java.io.ObjectInputStream.readObject(ObjectInputStream.java:461)
    at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:76)
    at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:115)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:458)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:750)

我的请求看起来像:

 curl -X POST http://xxxxxx:6066/v1/submissions/create --header "Content-Type:application/json;charset=UTF-8" --data '{
  "appResource": "/home/xxxx/xxxx-bundled-0.1.jar",
  "sparkProperties": {
    "spark.master": "spark://xxxxxxx:7077",
    "spark.driver.userClassPathFirst": "true",
    "spark.executor.userClassPathFirst": "true",
    "spark.app.name": "DataPipeline",
    "spark.submit.deployMode": "cluster",
    "spark.driver.supervise": "true"
  },
  "environmentVariables": {
    "SPARK_ENV_LOADED": "1"
  },
  "clientSparkVersion": "3.1.3",
  "mainClass": "com.xxxx.DataPipeline",
  "action": "CreateSubmissionRequest",
  "appArgs": [
    "--config=xxxx",
    "--runner=SparkRunner"
  ]

我设置了 "spark.driver.userClassPathFirst": "true" 和 "spark.executor.userClassPathFirst": "true" 因为在我的 jar 中使用了 proto3。 不确定为什么在执行器上找不到此类。 我的beam版本2.41.0,spark版本3.1.3,hadoop版本3.2.0。

最后,我将 shaded 插件升级到 3.4.0,然后 protobuf 的重定位开始工作,我删除了“spark.driver.userClassPathFirst”:“true”和“spark.executor.userClassPathFirst”:“true”。之后一切正常。 “spark-submit”在本地或通过 rest api 都有效。

【问题讨论】:

  • 请将您用于构建阴影 jar 的配置添加到问题中。另外,你有没有重新安排班级?您究竟是如何提交代码的?注意,如果使用userClassPathFirst,你必须小心地从你的 fat jar 中删除 Spark、Hadoop、Scala 类(以及更多)。
  • 1. 我已经尝试重新定位 protobuf3 的类,但它似乎不起作用,因此我设置了 userClassPathFirst=true 并且它起作用了。 2. 我首先构建了 shaded jar 包,然后将其复制到独立的 spark 主机,然后尝试在集群模式下运行 spark-submit(并且还尝试远程调用 rest api 来提交作业,如上)。两者都遇到相同的问题。客户端模式工作正常。 3.“删除”是指我将范围更改为“提供”还是“运行时”?
  • 谢谢,将 shaded 插件升级到 3.4.0 后,重定位工作正常,之后一切正常。
  • 通过删除我的意思是从 uber jar 中排除这些类。如果使用 userClassPathFirst 那很关键,但始终建议这样做。这些类已经存在于 Spark 类路径中,请在此处查看详细信息github.com/apache/beam/issues/23568#issuecomment-1286746306

标签: apache-spark apache-beam spark-submit


【解决方案1】:

最后,我将 shaded 插件升级到 3.4.0,然后 protobuf 的重定位开始工作,我删除了“spark.driver.userClassPathFirst”:“true”和“spark.executor.userClassPathFirst”:“true”。之后一切正常。 “spark-submit”在本地或通过 rest api 都有效。

shaded 插件中 protobuf 的重定位:

<relocations>
    <relocation>
        <pattern>com.google.protobuf</pattern>
        <shadedPattern>shaded.com.google.protobuf</shadedPattern>
    </relocation>
</relocations>

然后在这两个调用之后工作:

spark-submit --class XXXXPipeline --master spark://xxxx:7077 --deploy-mode cluster --supervise /xxxxx/xxxx-bundled-0.1.jar

curl -X POST http://xxxx:6066/v1/submissions/create --header "Content-Type:application/json;charset=UTF-8" --data '{
  "appResource": "xxxxx-bundled-0.1.jar",
  "sparkProperties": {
    "spark.master": "spark://XXXXXX:7077",
    "spark.jars": "xxxx-bundled-0.1.jar",
    "spark.app.name": "XXXPipeline",
    "spark.submit.deployMode": "cluster",
    "spark.driver.supervise": "true"
  },
  "environmentVariables": {
    "SPARK_ENV_LOADED": "1"
  },
  "clientSparkVersion": "3.1.3",
  "mainClass": "XXXXPipeline",
  "action": "CreateSubmissionRequest",
  "appArgs": []
}'

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-11-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多