You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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格式的数据源。
排查解决步骤

按以下顺序排查修复即可:

  1. 修正打包插件配置
    如果用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>
    
  2. 验证依赖是否被正确打入包内
    打包完成后,执行以下命令查看jar包内是否包含spark-sql-kafka相关类:
    jar tf my_jar.jar | grep kafka
    
    如果输出中找不到org/apache/spark/sql/kafka010相关路径,说明依赖没有被打进jar包,需要检查pom中spark-sql-kafka-0-10依赖的scope是否被覆盖为provided(当前配置为compile,正常会被打入,除非父pom修改了scope配置)。
  3. 排查依赖冲突
    检查打包后的jar包内是否存在多个版本的spark-sql、spark-core依赖,版本不一致也会导致数据源加载失败,要确保所有Spark相关依赖版本统一为2.4.0,对应Scala版本统一为2.11。
  4. 临时验证方案
    如果暂时不想重新打包,可以在执行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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.01 18:57:49