Spark Streaming提交任务报错NoSuchMethodError:scala.collection.immutable.Map$.apply
解决Spark Streaming提交Kafka任务时的NoSuchMethodError错误
问题背景
我编写了一个基于Java 1.8的简单Spark Streaming程序,用于流式处理Kafka主题的数据。将程序打包为包含所有依赖的Uber JAR后,提交至Spark 3.3.1集群时出现NoSuchMethodError错误。
程序代码
import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.spark.SparkConf; import org.apache.spark.streaming.Durations; import org.apache.spark.streaming.api.java.*; import org.apache.spark.streaming.kafka010.*; import org.apache.kafka.common.serialization.StringDeserializer; import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.Map; public class KafkaStreamingExample { public static void main(String[] args) throws InterruptedException { SparkConf conf = new SparkConf().setAppName("KafkaStreamingExample").setMaster("local[2]"); JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(5)); Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", "localhost:9092"); kafkaParams.put("key.deserializer", StringDeserializer.class); kafkaParams.put("value.deserializer", StringDeserializer.class); kafkaParams.put("group.id", "group1"); kafkaParams.put("auto.offset.reset", "latest"); kafkaParams.put("enable.auto.commit", false); Collection<String> topics = Arrays.asList("dbserver1.inventory.customers"); JavaInputDStream<ConsumerRecord<String, String>> stream = KafkaUtils.createDirectStream( jssc, LocationStrategies.PreferConsistent(), ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams) ); stream.foreachRDD(rdd -> { rdd.foreach(record -> { System.out.println(record.value()); }); }); jssc.start(); jssc.awaitTermination(); } }
提交命令
./bin/spark-submit --class KafkaStreamingExample --master yarn --deploy-mode client --executor-memory 2g --num-executors 1 kafka-streaming-example-jar-with-dependencies.jar
错误日志
Exception in thread "main" java.lang.NoSuchMethodError: scala.collection.immutable.Map$.apply(Lscala/collection/immutable/Seq;)Ljava/lang/Object; at org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.<init>(DirectKafkaInputDStream.scala:68) at org.apache.spark.streaming.kafka010.KafkaUtils$.createDirectStream(KafkaUtils.scala:136) at org.apache.spark.streaming.kafka010.KafkaUtils$.createDirectStream(KafkaUtils.scala:115) at org.apache.spark.streaming.kafka010.KafkaUtils$.createDirectStream(KafkaUtils.scala:156) at org.apache.spark.streaming.kafka010.KafkaUtils.createDirectStream(KafkaUtils.scala) at KafkaStreamingExample.main(KafkaStreamingExample.java:28) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52) at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:958) at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:180) at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:203) at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:90) at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:1046) at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:1055) at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
解决方案
这个错误是Scala版本不兼容或Spark与Kafka依赖版本不匹配导致的:Spark 3.3.1依赖Scala 2.12,若打包时引入了Scala 2.11的库,或者spark-streaming-kafka-0-10版本与Spark集群版本不一致,就会出现方法签名不匹配的问题。
具体修改步骤:
对齐依赖版本:确保
spark-streaming-kafka-0-10和spark-streaming的版本与Spark集群版本(3.3.1)完全一致,且使用对应Scala 2.12的包。以Maven为例,依赖配置如下:<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.12</artifactId> <version>3.3.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming_2.12</artifactId> <version>3.3.1</version> <scope>provided</scope> </dependency>注意包名中的
_2.12标识,不能使用Scala 2.11版本的依赖。排除集群已有的依赖:将Spark和Scala相关依赖的scope设为
provided(Maven)或compileOnly(Gradle),避免将集群已自带的依赖打包进Uber JAR,防止版本冲突。重新打包提交:删除旧的Uber JAR,重新构建正确的J包后,用原提交命令重新提交任务。
内容的提问来源于stack exchange,提问作者Eugene Goldberg
相关产品推荐
相关产品推荐

