K8s环境下Spark读取Kafka报NoClassDefFoundError求助
在Minikube搭建的K8s环境中,使用Spark Streaming读取Kafka Topic数据时触发错误:java.lang.NoClassDefFoundError: org/apache/spark/kafka010/KafkaConfigUpdater。
环境配置:
- Bitnami Spark 3.5.1(1主节点+2工作节点)
- Bitnami Kafka 3.7.0(3个Kafka节点+1个客户端)
可能的原因及解决方法
1. Spark Kafka连接器依赖缺失或版本不匹配
KafkaConfigUpdater类属于Spark官方Kafka连接器库,出现该错误核心原因是作业运行时缺少对应依赖,或依赖版本与Spark 3.5.1不兼容。
- 解决步骤:
- 确认依赖版本匹配:Spark 3.5.x对应Kafka连接器版本为
org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1(Scala版本需与你的环境一致,此处默认2.12)。 - 提交作业时通过
--packages参数引入依赖:spark-submit \ --class com.your.package.YourStreamingApp \ --master k8s://https://<k8s-api-server-url> \ --deploy-mode cluster \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 \ your-app.jar - 若K8s集群无法访问公网Maven仓库,可提前将依赖包打包到应用镜像,或在集群内搭建私有Maven仓库。
- 确认依赖版本匹配:Spark 3.5.x对应Kafka连接器版本为
2. 依赖包冲突
若应用Jar中包含旧版Spark Kafka连接器,或集群环境存在冲突依赖,会导致类加载失败。
- 解决步骤:
- 检查应用依赖树,排除冲突的Kafka连接器依赖(以Maven为例):
<dependency> <groupId>your.dependency.group</groupId> <artifactId>your-dependency</artifactId> <version>x.x.x</version> <exclusions> <exclusion> <groupId>org.apache.spark</groupId> <artifactId>spark-sql-kafka-0-10_2.12</artifactId> </exclusion> </exclusions> </dependency> - 清理Spark集群节点classpath中多余的Kafka连接器Jar包。
- 检查应用依赖树,排除冲突的Kafka连接器依赖(以Maven为例):
3. Bitnami Spark镜像默认无Kafka连接器
Bitnami官方Spark镜像未预装Kafka连接器,需自定义镜像添加依赖,或提交作业时强制引入。
- 解决步骤(自定义镜像):
构建镜像后替换集群默认的Spark镜像即可。FROM bitnami/spark:3.5.1 USER root RUN wget -P /opt/bitnami/spark/jars/ https://repo1.maven.org/maven2/org/apache/spark/spark-sql-kafka-0-10_2.12/3.5.1/spark-sql-kafka-0-10_2.12-3.5.1.jar RUN wget -P /opt/bitnami/spark/jars/ https://repo1.maven.org/maven2/org/apache/spark/spark-token-provider-kafka-0-10_2.12/3.5.1/spark-token-provider-kafka-0-10_2.12-3.5.1.jar USER 1001
4. 作业提交配置错误
若使用--jars参数引入本地依赖,需确保包含Kafka连接器及其关联依赖(如kafka-clients、spark-token-provider-kafka-0-10等),避免遗漏导致类缺失。
示例Spark Streaming代码(Scala)
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object KafkaStreamingApp { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("KafkaStreamingDemo") .getOrCreate() import spark.implicits._ val kafkaDF = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "kafka-headless.default.svc.cluster.local:9092") .option("subscribe", "your-target-topic") .load() val query = kafkaDF.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") .writeStream .outputMode("append") .format("console") .start() query.awaitTermination() } }
完整报错日志示例
Exception in thread "main" java.lang.NoClassDefFoundError: org/apache/spark/kafka010/KafkaConfigUpdater
at org.apache.spark.sql.kafka010.KafkaSourceProvider.createSource(KafkaSourceProvider.scala:142)
at org.apache.spark.sql.execution.datasources.DataSource.createSource(DataSource.scala:324)
at org.apache.spark.sql.execution.streaming.StreamExecution.createSource(StreamExecution.scala:278)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$15(MicroBatchExecution.scala:463)
at scala.collection.TraversableLike.$anonfun$flatMap$1(TraversableLike.scala:293)
at scala.collection.Iterator.foreach(Iterator.scala:943)
at scala.collection.Iterator.foreach$(Iterator.scala:943)
at scala.collection.AbstractIterator.foreach(Iterator.scala:1431)
at scala.collection.IterableLike.foreach(IterableLike.scala:74)
at scala.collection.IterableLike.foreach$(IterableLike.scala:73)
at scala.collection.AbstractIterable.foreach(Iterable.scala:56)
at scala.collection.TraversableLike.flatMap(TraversableLike.scala:293)
at scala.collection.TraversableLike.flatMap$(TraversableLike.scala:290)
at scala.collection.AbstractTraversable.flatMap(Traversable.scala:108)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runBatch(MicroBatchExecution.scala:459)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:239)
at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:474)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:239)
at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:386)
at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:283)
Caused by: java.lang.ClassNotFoundException: org.apache.spark.kafka010.KafkaConfigUpdater
at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:641)
at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:188)
at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:520)
... 20 more
内容的提问来源于stack exchange,提问作者vivekdesai

