Spark集成Kafka报AnalysisException,如何无需spark-submit直接运行?
解决Spark Structured Streaming本地源码运行找不到Kafka数据源的问题
报错信息
org.apache.spark.sql.AnalysisException: Failed to find data source: kafka
环境信息
- Spark版本:3.3.0
- Java版本:1.8
- 项目类型:Spring Boot + Maven
解决方案
1. 补全Spark-Kafka依赖
在pom.xml中添加spark-sql-kafka-0-10依赖,本地源码运行时不要设置scope=provided,确保依赖被加入类路径:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql-kafka-0-10_2.12</artifactId> <version>3.3.0</version> <!-- 本地运行请移除provided,打包部署到Spark集群时可改为provided --> </dependency>
注意:artifactId中的_2.12对应Spark 3.3.0默认的Scala版本,需与项目使用的Scala版本保持一致。
2. 优化Maven Shade插件配置
当前配置已包含DataSourceRegister的合并逻辑,需确保以下两点:
ServicesResourceTransformer正确合并所有META-INF/services下的服务文件,保证Kafka数据源的注册信息被加载- 打包时不要排除Spark-Kafka相关的依赖文件,确保shaded包中包含完整的Kafka数据源实现
3. 配置本地运行的SparkSession
在代码中初始化SparkSession时,指定本地运行模式,确保本地环境能加载Kafka数据源:
SparkSession spark = SparkSession.builder() .master("local[*]") // 使用本地所有可用核心运行 .appName("KafkaStreamValidator") .getOrCreate();
如果是Spring Boot项目,需注意Spark上下文与Spring上下文的初始化顺序,避免线程冲突。
4. IDE运行的额外配置
- 在IDE中将
spark-sql-kafka-0-10依赖标记为Compile,确保IDE将其加入运行类路径 - 统一Kafka客户端版本:Spark 3.3.0默认依赖
kafka-clients:2.8.1,若项目中有其他Kafka相关依赖,需在dependencyManagement中统一版本,避免冲突:
<dependencyManagement> <dependencies> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>2.8.1</version> </dependency> </dependencies> </dependencyManagement>
内容的提问来源于stack exchange,提问作者유연수
相关产品推荐
相关产品推荐

