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

Spark Streaming读取Kafka数据时遇KafkaConfigUpdater类找不到错误求助

Spark Kafka流读取报错java.lang.NoClassDefFoundError: org/apache/spark/kafka010/KafkaConfigUpdater原因分析

问题场景

尝试使用Spark从Kafka流式读取数据,代码如下:

from pyspark import SparkConf, SparkContext
from pyspark.sql import SparkSession 
from pyspark.sql.functions import *
from pyspark.sql.types import *

spark = SparkSession.builder.master("yarn")
    .getOrCreate()

df = spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "11.11.11.11:9092") \
    .option("subscribe", "topic-1") \
    .option("kafka.security.protocol", "SASL_SSL") \
    .option("kafka.sasl.mechanism", "SCRAM-SHA-512") \
    .option("kafka.sasl.jaas.config", 'org.apache.kafka.common.security.scram.ScramLoginModule required username="bdsm" password="C28776666w";') \
    .option("kafka.ssl.truststore.location", "/usr/local/hadoop/spark-3.3.2-bin-hadoop3/kafka_broker_topic.trustst.jks") \
    .option("kafka.ssl.truststore.password", "storepass") \
    .option("kafka.ssl.keystore.location", "/usr/local/hadoop/spark-3.3.2-bin-hadoop3/kafka_broker_topic.keyst.jks") \
    .option("kafka.ssl.keystore.password", "storepass") \
    .option("kafka.ssl.key.password", "storepass") \
    .option("kafka.session.timeout.ms", "6000") \
    .option("startingOffsets", "latest") \
    .load()

data = df.selectExpr("CAST(value AS STRING) as json") \
    .select(from_json("json", schema).alias("data")) \
    .select("data.*")

query = data \
    .writeStream \
    .format("parquet") \
    .option("path", "hdfs://10.10.10.10:8020/user/spark/logs_spark") \
    .option("checkpointLocation", "hdfs://10.10.10.10:8020/user/spark/checkpoint/dir") \
    .trigger(processingTime='2 minutes') \
    .start()

query.awaitTermination()

报错信息

运行时触发如下错误:

23/06/16 08:42:46 WARN ResolveWriteToStream: spark.sql.adaptive.enabled is not supported in streaming DataFrames/Datasets and will be disabled.
23/06/16 08:42:46 ERROR MicroBatchExecution: Query [id = c8b77688-eb2e-42b6-9b1d-0f091bc5ded3, runId = 93ab5673-96de-4aee-bb5e-fe26a20f9c83] terminated with error
java.lang.NoClassDefFoundError: org/apache/spark/kafka010/KafkaConfigUpdater
    at org.apache.spark.sql.kafka010.KafkaSourceProvider$.kafkaParamsForDriver(KafkaSourceProvider.scala:645)
    at org.apache.spark.sql.kafka010.KafkaSourceProvider$KafkaScan.toMicroBatchStream(KafkaSourceProvider.scala:482)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$1.$anonfun$applyOrElse$4(MicroBatchExecution.scala:104)

环境信息

  • Spark版本:3.3.2
  • Scala版本:2.12
  • 本地依赖jar包:
    • spark-streaming-kafka-0-10-assembly_2.12-3.3.2.jar
    • spark-streaming-kafka-0-10_2.12-3.3.2.jar
    • spark-token-provider-kafka-0-10_2.12-3.3.2.jar
    • spark-sql-kafka-0-10_2.12-3.3.2.jar
    • kafka-clients-3.3.2.jar

问题原因

KafkaConfigUpdater类属于spark-sql-kafka-0-10模块,出现该错误的核心原因是集群Executor节点无法找到该类,具体分为以下几种情况:

  1. 依赖未分发到集群:本地环境虽有对应jar包,但提交Yarn任务时,未通过参数将jar包传递给集群Executor,导致节点缺失依赖。
  2. 版本冲突:集群环境中存在旧版本的Spark-Kafka相关jar包,覆盖了当前使用的3.3.2版本,引发类加载异常。
  3. 提交参数错误:使用Yarn模式提交任务时,未通过--packages或--jars显式声明Kafka依赖,Yarn无法自动加载对应类。
  4. 隐性依赖缺失:spark-sql-kafka-0-10依赖的部分关联jar包未被包含(比如Spark核心或Streaming模块的对应版本包)。

内容的提问来源于stack exchange,提问作者VLBigo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 05:02:04