You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark2 Streaming on Yarn模式下ClassPath问题求助

针对Yarn Client模式下Spark Streaming Kafka 0.10类缺失问题的解决方案

我之前在Cloudera平台的Yarn Client模式下也碰到过几乎一模一样的java.lang.ClassNotFoundException: org.apache.spark.streaming.kafka010.KafkaRDDPartition问题,结合你的排查情况,给你几个针对性的解决思路:

1. 先确认Uber Jar确实包含缺失的类

别默认Shade插件已经把目标类打包进去了,用命令直接验证:

jar tf your-uber-jar-name.jar | grep org.apache.spark.streaming.kafka010.KafkaRDDPartition

如果没有任何输出,说明Maven Shade插件的配置有问题,需要调整打包规则,确保没有排除Spark Kafka相关类。比如在Shade插件中添加明确的包含规则:

<plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-shade-plugin</artifactId>
    <version>3.2.4</version>
    <executions>
        <execution>
            <phase>package</phase>
            <goals>
                <goal>shade</goal>
            </goals>
            <configuration>
                <filters>
                    <filter>
                        <artifact>org.apache.spark:spark-streaming-kafka-0.10_2.11</artifact>
                        <includes>
                            <include>org/apache/spark/streaming/kafka010/**</include>
                            <include>org/apache/kafka/**</include>
                        </includes>
                    </filter>
                </filters>
                <!-- Spring Boot项目需注意避免和自带打包插件冲突,这里指定主类 -->
                <transformers>
                    <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                        <mainClass>your.main.ClassName</mainClass>
                    </transformer>
                </transformers>
            </configuration>
        </execution>
    </executions>
</plugin>

2. 修正Yarn模式下的类路径配置错误

你之前用spark.executor.extraClassPath指定的sparkHome + "/kafka-0.10"是本地机器的路径,但Yarn Executor运行在集群节点上,这些节点上不一定有这个目录,更不会有你本地的Jar包,这是核心问题之一!

推荐两种更可靠的方式:

  • 方式一:提交时用--packages自动拉取依赖
    直接在spark-submit命令中指定Kafka 0.10的依赖,Spark会自动下载并分发到所有Executor节点:
    spark-submit \
      --master yarn-client \
      --packages org.apache.spark:spark-streaming-kafka-0.10_2.11:2.4.0 \ # 版本要和你的Spark版本匹配
      --class com.your.package.KafkaSubscriber \
      your-uber-jar.jar
    
  • 方式二:用spark.jars指定HDFS上的Jar包
    把依赖Jar上传到HDFS(集群所有节点都能访问),然后配置:
    val conf: SparkConf = new SparkConf()
      .set("spark.streaming.concurrentJobs", "2")
      .set("spark.jars", "hdfs://path/to/spark-streaming-kafka-0.10_2.11-2.4.0.jar,hdfs://path/to/kafka-clients-0.10.2.2.jar")
      .setAppName(classOf[KafkaSubscriber].getSimpleName)
      .setMaster(sparkMaster)
    

3. 解决Cloudera平台的依赖冲突

Cloudera CDH默认会自带旧版本的Spark Kafka依赖(比如0.8或0.9版本),这些旧版本中没有KafkaRDDPartition类,而且Yarn的类加载优先级可能会优先加载平台自带的依赖,导致你的新版本类被覆盖。

解决方法是强制让用户Jar的类优先加载,在SparkConf中添加:

conf.set("spark.driver.userClassPathFirst", "true")
conf.set("spark.executor.userClassPathFirst", "true")

或者在spark-submit命令中添加参数:

--conf spark.driver.userClassPathFirst=true --conf spark.executor.userClassPathFirst=true

4. 放弃手动添加Jar的冗余代码

你之前尝试遍历本地目录添加Jar的方式在Yarn Client模式下不生效,因为addJar是把本地Jar上传到Spark分布式缓存,但如果你的Uber Jar已经包含了依赖,反而可能导致类重复或冲突,建议去掉这部分代码。


内容的提问来源于stack exchange,提问作者Cheloute

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 03:50:32