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

Spark流数据写入Kafka报NoClassDefFoundError问题咨询

问题解决:Spark Structured Streaming写Kafka报NoClassDefFoundError

问题重现

用于将流数据从Kafka Topic转发到另一个Topic的代码:

query = info_df_fin.selectExpr("to_json(struct(*)) AS value")
        .writeStream
        .format("kafka")
        .option("checkpointLocation", "chkpint_directory")
        .option("kafka.bootstrap.servers", "localhost:9092")
        .option("topic", "Dest_topic")
        .start()

已放置spark-streaming-kafka-0-10_2.12-2.4.0.jar,但仍出现错误:

23/09/18 11:15:58 WARN ResolveWriteToStream:
spark.sql.adaptive.enabled is not supported in streaming
DataFrames/Datasets and will be disabled. 23/09/18 11:15:58 ERROR
MicroBatchExecution: Query [id = 23d94359-dabf-41ef-8617-a62859a35adb,
runId = 57e0c21c-5245-40c5-9878-be1b5b06e106] terminated with error
java.lang.NoClassDefFoundError:
org/apache/spark/sql/internal/connector/SupportsStreamingUpdate at
java.base/java.lang.ClassLoader.defineClass1(Native Method) at
java.base/java.lang.ClassLoader.defineClass(ClassLoader.java:1022) at
java.base/java.security.SecureClassLoader.defineClass(SecureClassLoader.java:174)
at
java.base/jdk.internal.loader.BuiltinClassLoader.defineClass(BuiltinClassLoader.java:800)
at
java.base/jdk.internal.loader.BuiltinClassLoader.findClassOnClassPathOrNull(BuiltinClassLoader.java:698)
at
java.base/jdk.internal.loader.BuiltinClassLoader.loadClassOrNull(BuiltinClassLoader.java:621)
at
java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:579)
at
java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:178)
at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:527)
at
org.apache.spark.sql.kafka010.KafkaSourceProvider$KafkaTable.newWriteBuilder(KafkaSourceProvider.scala:397)
at
org.apache.spark.sql.execution.datasources.v2.V2Writes$.org$apache$spark$sql$execution$datasources$v2$V2Writes$$newWriteBuilder(V2Writes.scala:144)
at
org.apache.spark.sql.execution.datasources.v2.V2Writes$$anonfun$apply$1.applyOrElse(V2Writes.scala:90)
at
org.apache.spark.sql.execution.datasources.v2.V2Writes$$anonfun$apply$1.applyOrElse(V2Writes.scala:43)
at
org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:512)
at
org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:104)
at
org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:512)
at
org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:31)

错误原因

  1. 版本不兼容:org.apache.spark.sql.internal.connector.SupportsStreamingUpdate是Spark 3.x版本引入的接口,你使用的spark-streaming-kafka-0-10_2.12-2.4.0.jar对应Spark 2.4版本,和当前运行的Spark版本(推测为3.x)不匹配,导致类找不到。
  2. 依赖包错误:Structured Streaming操作Kafka需要的是spark-sql-kafka-0-10系列包,而非spark-streaming-kafka-0-10(后者是旧的DStream API依赖),你可能选错了依赖包。
  3. 类路径加载问题:手动放置的Jar包未被Spark正确加载,或者存在多个版本的依赖冲突。

解决方案

1. 匹配Spark与Kafka连接器版本

根据你实际使用的Spark版本,选择对应版本的spark-sql-kafka-0-10包:

  • 如果用Spark 3.x(比如3.3.0):选择spark-sql-kafka-0-10_2.12-3.3.0.jar
  • 如果用Spark 2.4.x:选择spark-sql-kafka-0-10_2.12-2.4.0.jar

版本号需严格对应,Spark和Kafka连接器的主版本号必须一致(比如Spark 3.3.x对应连接器3.3.x)。

2. 使用正确的依赖包

直接替换成spark-sql-kafka-0-10包,删除错误的spark-streaming-kafka-0-10包——Structured Streaming基于Spark SQL API,不需要旧的DStream依赖。

3. 推荐通过spark-submit自动管理依赖

避免手动放包出错,提交作业时用--packages参数指定依赖,Spark会自动下载并加载对应版本的包:

spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 your_streaming_script.py

(将版本号替换为你实际使用的Spark版本)

4. IDE环境检查

如果是在IDE(比如PyCharm、IntelliJ)中运行,确保:

  • 项目依赖的Spark版本和Kafka连接器版本完全一致
  • 所有依赖都已正确引入,没有冲突的旧版本包

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 14:00:56