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节点无法找到该类,具体分为以下几种情况:
- 依赖未分发到集群:本地环境虽有对应jar包,但提交Yarn任务时,未通过参数将jar包传递给集群Executor,导致节点缺失依赖。
- 版本冲突:集群环境中存在旧版本的Spark-Kafka相关jar包,覆盖了当前使用的3.3.2版本,引发类加载异常。
- 提交参数错误:使用Yarn模式提交任务时,未通过
--packages或--jars显式声明Kafka依赖,Yarn无法自动加载对应类。 - 隐性依赖缺失:
spark-sql-kafka-0-10依赖的部分关联jar包未被包含(比如Spark核心或Streaming模块的对应版本包)。
内容的提问来源于stack exchange,提问作者VLBigo
相关产品推荐
相关产品推荐

