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

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包不降级,想知道错误原因和解决办法。


错误原因分析

咱们从两个核心点拆解这个问题:

  1. Spark组件版本不一致:你引入的spark-streaming-kafka-0-8_2.11版本是2.1.1,但其他Spark核心、SQL、Streaming依赖都是2.1.0。Spark不同小版本之间内部API和类结构会有细微调整,org.apache.spark.internal.Logging这个类的加载路径在版本不匹配时会出现冲突,导致找不到类。
  2. 依赖与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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:57:45