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

