Spark消费Kafka程序IDE运行正常 java -jar执行报找不到kafka数据源错误
问题描述
编写Java Structured Streaming脚本读取Kafka数据,核心代码如下:
import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.functions; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; public class App { public static void main( String[] args ) { System.out.println( "Hello World!" ); SparkSession spark = SparkSession.builder().appName("Kafka_Load").config("spark.driver.allowMultipleContexts", "true").config("spark.master", "local").getOrCreate(); Dataset<Row> df = spark.readStream().format("kafka").option("kafka.bootstrap.servers", "localhost:9092").option("subscribe", "my_topic").load(); System.out.println( "Hello World1!" ); } }
- 运行表现:Eclipse中直接以Java Application模式运行无异常,执行
java -jar my_jar.jar运行打包后的jar包时抛出异常:
Exception in thread "main" org.apache.spark.sql.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".;
项目pom.xml依赖配置如下:
<dependencies> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>4.12</version> <scope>test</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.11</artifactId> <version>2.4.0</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.11</artifactId> <version>2.4.0</version> <scope>compile</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming_2.11</artifactId> <version>2.4.0</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.11</artifactId> <version>2.4.0</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql-kafka-0-10_2.11</artifactId> <version>2.4.0</version> </dependency> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-hdfs</artifactId> <version>2.2.0</version> </dependency> </dependencies>
根因分析
报错核心原因是打包生成的jar包没有正确包含spark-sql-kafka-0-10相关组件,或者打包过程中丢失了Spark SQL数据源注册需要的SPI声明文件:
- Eclipse运行时会自动把所有Maven依赖加入运行时classpath,因此可以正常加载Kafka数据源
- 用
java -jar运行时,程序classpath仅包含jar包内部打包的内容,如果打包配置错误,要么缺失Kafka连接器依赖,要么缺失META-INF/services目录下的数据源注册文件,Spark就无法识别kafka格式的数据源。
排查解决步骤
按以下顺序排查修复即可:
- 修正打包插件配置
如果用maven-shade-plugin打fat jar,必须添加ServicesResourceTransformer,否则打包时会覆盖依赖中的SPI服务注册文件,直接导致Kafka数据源无法被Spark识别。正确的shade插件配置示例:<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.2.4</version> <executions> <execution> <phase>package</phase> <goals> <goal>shade</goal> </goals> <configuration> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>App</mainClass> <!-- 替换为实际主类全限定名 --> </transformer> <!-- 必须添加该转换器,合并所有依赖的SPI服务文件 --> <transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/> </transformers> </configuration> </execution> </executions> </plugin> - 验证依赖是否被正确打入包内
打包完成后,执行以下命令查看jar包内是否包含spark-sql-kafka相关类:
如果输出中找不到jar tf my_jar.jar | grep kafkaorg/apache/spark/sql/kafka010相关路径,说明依赖没有被打进jar包,需要检查pom中spark-sql-kafka-0-10依赖的scope是否被覆盖为provided(当前配置为compile,正常会被打入,除非父pom修改了scope配置)。 - 排查依赖冲突
检查打包后的jar包内是否存在多个版本的spark-sql、spark-core依赖,版本不一致也会导致数据源加载失败,要确保所有Spark相关依赖版本统一为2.4.0,对应Scala版本统一为2.11。 - 临时验证方案
如果暂时不想重新打包,可以在执行java -jar命令时通过--jars参数手动指定Kafka连接器包路径,示例:java -jar my_jar.jar --jars /path/to/spark-sql-kafka-0-10_2.11-2.4.0.jar
注意:当前代码仅调用了
load()方法,没有添加writeStream的输出触发逻辑,即使修复了打包问题,程序运行后也不会真正消费Kafka数据,需要补充输出逻辑(例如df.writeStream().format("console").start().awaitTermination())才能看到实际消费效果。
内容的提问来源于stack exchange,提问作者Pedro Alves
相关产品推荐
相关产品推荐

