Spark流数据写入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)
错误原因
- 版本不兼容:
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)不匹配,导致类找不到。 - 依赖包错误:Structured Streaming操作Kafka需要的是
spark-sql-kafka-0-10系列包,而非spark-streaming-kafka-0-10(后者是旧的DStream API依赖),你可能选错了依赖包。 - 类路径加载问题:手动放置的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

