JavaInputDStream异常求助:Spark集成Kafka报AbstractMethodError
问题分析与解决方案
你遇到的java.lang.AbstractMethodError是典型的依赖版本不兼容问题,从你的POM配置里能直接找到根源:
核心问题点
- Spark版本不匹配:你的
spark-streaming_2.11用的是2.3.0版本,但spark-streaming-kafka-0-10_2.10用的是2.0.0版本。Spark的各个组件(包括Kafka整合包)必须使用完全一致的主版本号,否则会出现类方法签名不匹配的问题。 - Scala版本不兼容:
spark-streaming_2.11对应的是Scala 2.11,而spark-streaming-kafka-0-10_2.10对应的是Scala 2.10。不同Scala版本编译的类文件无法兼容,这直接导致了抽象方法找不到的错误。
修正后的POM依赖配置
把Kafka整合包的版本和Scala后缀调整为与Spark Streaming一致:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming_2.11</artifactId> <version>2.3.0</version> <scope>provided</scope> </dependency> <!-- Spark Streaming Kafka 0.10 整合包,版本与Spark一致,Scala后缀统一为2.11 --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.11</artifactId> <version>2.3.0</version> </dependency>
额外注意事项
- 所有Spark相关依赖的主版本号必须严格一致(比如这里全用2.3.0),包括
spark-core、spark-streaming、spark-sql等组件。 - 确认你的项目使用的Scala版本与依赖后缀匹配:Spark 2.3.x官方推荐使用Scala 2.11,所以所有依赖都要带
_2.11后缀。 - 如果是在本地Eclipse中运行,注意
providedscope的依赖:spark-streaming标记为provided意味着运行时由环境提供,但本地运行时可能需要移除这个scope,或者手动添加Spark相关的jar包到项目依赖中。
原错误信息
Exception in thread "main" java.lang.AbstractMethodError at org.apache.spark.internal.Logging$class.initializeLogIfNecessary(Logging.scala:99) at org.apache.spark.streaming.kafka010.KafkaUtils$.initializeLogIfNecessary(KafkaUtils.scala:40) at org.apache.spark.internal.Logging$class.log(Logging.scala:46) at org.apache.spark.streaming.kafka010.KafkaUtils$.log(KafkaUtils.scala:40) at org.apache.spark.internal.Logging$class.logWarning(Logging.scala:66) at org.apache.spark.streaming.kafka010.KafkaUtils$.logWarning(KafkaUtils.scala:40) at org.apache.spark.streaming.kafka010.KafkaUtils$.fixKafkaParams(KafkaUtils.scala:157) at org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.<init>(DirectKafkaInputDStream.scala:65) at org.apache.spark.streaming.kafka010.KafkaUtils$.createDirectStream(KafkaUtils.scala:126) at org.apache.spark.streaming.kafka010.KafkaUtils$.createDirectStream(KafkaUtils.scala:149) at org.apache.spark.streaming.kafka010.KafkaUtils.createDirectStream(KafkaUtils.scala) at com.spark.kafka.JavaDirectKafkaWordCount.main(JavaDirectKafkaWordCount.java:50) 18/05/29 18:05:43 INFO SparkContext: Invoking stop() from shutdown hook 18/05/29 18:05:43 INFO SparkUI: Stopped Spark web UI at
原代码片段
package com.spark.kafka; import java.util.HashMap; import java.util.HashSet; import java.util.Arrays; import java.util.Map; import java.util.Set; import java.util.regex.Pattern; import scala.Tuple2; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.spark.SparkConf; import org.apache.spark.streaming.api.java.*; import org.apache.spark.streaming.kafka010.ConsumerStrategies; import org.apache.spark.streaming.kafka010.KafkaUtils; import org.apache.spark.streaming.kafka010.LocationStrategies; import org.apache.spark.streaming.Durations; public final class JavaDirectKafkaWordCount { private static final Pattern SPACE = Pattern.compile(" "); public static void main(String[] args) throws Exception { String brokers = "localhost:9092"; String topics = "sparktestone"; SparkConf sparkConf = new SparkConf().setAppName("JavaDirectKafkaWordCount").setMaster("local[*]"); JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, Durations.seconds(2)); Set<String> topicsSet = new HashSet<>(Arrays.asList(topics.split(","))); Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", brokers); JavaInputDStream<ConsumerRecord<String, String>> messages = KafkaUtils.createDirectStream( jssc, LocationStrategies.PreferConsistent(), ConsumerStrategies.Subscribe(topicsSet, kafkaParams)); JavaDStream<String> lines = messages.map(ConsumerRecord::value); JavaDStream<String> words = lines.flatMap(x -> Arrays.asList(SPACE.split(x)).iterator()); JavaPairDStream<String, Integer> wordCounts = words.mapToPair(s -> new Tuple2<>(s, 1)) .reduceByKey((i1, i2) -> i1 + i2); wordCounts.print(); jssc.start(); jssc.awaitTermination(); } }
原POM依赖
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming_2.11</artifactId> <version>2.3.0</version> <scope>provided</scope> </dependency> <!-- https://mvnrepository.com/artifact/org.apache.spark/spark-streaming-kafka-0-10 --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.10</artifactId> <version>2.0.0</version> </dependency>
内容的提问来源于stack exchange,提问作者andani
相关产品推荐
相关产品推荐

