Spark 2.1版本运行时出现Logging类缺失错误,求解决方法
Spark 2.1版本中Spark Streaming读取Kafka时NoClassDefFoundError问题解决
问题场景
我在使用以下Spark依赖配置:
<dependencies> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.11</artifactId> <version>2.1.0</version> <exclusions> <exclusion> <groupId>log4j</groupId> <artifactId>log4j</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-8_2.11</artifactId> <version>2.1.1</version> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> </exclusion> <exclusion> <groupId>log4j</groupId> <artifactId>log4j</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.11</artifactId> <version>2.1.0</version> </dependency> <!-- https://mvnrepository.com/artifact/org.apache.spark/spark-streaming --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming_2.11</artifactId> <version>2.1.0</version> </dependency> </dependencies>
运行这段Spark结构化流读取Kafka的代码时:
val spark = SparkSession .builder .appName("Test Data") .master("local[*]") .getOrCreate() import spark.implicits._ val df = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", "192.168.0.40:9092") .option("zookeeper.connect", "192.168.0.40:2181") .option("subscribe", "topic") .option("startingOffsets", "earliest") .option("max.poll.records", 100) .option("failOnDataLoss", false) .load() import org.apache.spark.sql.Encoders val schema = Encoders.product[event].schema val ds = df.select(from_json($"value" cast "string", schema)).as[event] val query = ds.writeStream .outputMode("append") .queryName("table") .format("console") .start() query.awaitTermination()
抛出了如下错误:
Exception in thread "main" java.lang.NoClassDefFoundError: org/apache/spark/internal/Logging at java.lang.ClassLoader.defineClass1(Native Method) at java.lang.ClassLoader.defineClass(ClassLoader.java:763) at java.security.SecureClassLoader.defineClass(SecureClassLoader.java:142)
我希望保留当前Spark 2.1版本的jar包不降级,想知道错误原因和解决办法。
错误原因分析
咱们从两个核心点拆解这个问题:
- Spark组件版本不一致:你引入的
spark-streaming-kafka-0-8_2.11版本是2.1.1,但其他Spark核心、SQL、Streaming依赖都是2.1.0。Spark不同小版本之间内部API和类结构会有细微调整,org.apache.spark.internal.Logging这个类的加载路径在版本不匹配时会出现冲突,导致找不到类。 - 依赖与API不匹配:你用的是Spark SQL的结构化流API(
readStream.format("kafka")),但引入的却是旧版的spark-streaming-kafka-0-8依赖——这个依赖是针对传统DStream API的,和结构化流的依赖体系完全不兼容,这才是引发类找不到错误的关键原因。
解决办法
按照以下步骤调整,就能在保留Spark 2.1版本的前提下解决问题:
- 统一所有Spark组件版本:把
spark-streaming-kafka-0-8_2.11的版本从2.1.1改成2.1.0,确保所有Spark相关依赖的版本完全一致,消除版本冲突。 - 替换为结构化流对应的Kafka依赖:移除旧的
spark-streaming-kafka-0-8_2.11依赖,换成结构化流专用的spark-sql-kafka-0-10_2.11依赖,版本保持2.1.0:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql-kafka-0-10_2.11</artifactId> <version>2.1.0</version> <exclusions> <exclusion> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> </exclusion> <exclusion> <groupId>log4j</groupId> <artifactId>log4j</artifactId> </exclusion> </exclusions> </dependency>
- 清理冗余依赖:
spark-streaming_2.11依赖在使用结构化流时是多余的,可以直接移除,因为结构化流的功能已经包含在spark-sql_2.11中。 - 代码配置调整:你的代码里的
zookeeper.connect是旧版DStream的配置项,结构化流读取Kafka不需要这个配置,直接移除该option即可,结构化流只需要依赖Kafka原生的kafka.bootstrap.servers等配置。
内容的提问来源于stack exchange,提问作者Pinnacle
相关产品推荐
相关产品推荐

