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

Spark读取Docker中Kafka Topic报错:找不到kafka数据源

问题解决思路

1. 修复Kafka依赖的作用域

你的spark-sql-kafka-0-10_2.13依赖被设置为<scope>test</scope>,这意味着该依赖只会在测试代码中生效,主程序运行时无法加载Kafka数据源,这是报错的核心原因。

修改pom.xml中的Kafka依赖,移除scope标签:

<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-sql-kafka-0-10_2.13</artifactId>
    <version>3.3.2</version>
</dependency>

2. 统一Scala版本依赖

你的Spark依赖使用的是Scala 2.13(比如spark-core_2.13),但Iceberg的runtime依赖用的是Scala 2.12(iceberg-spark-runtime-3.5_2.12),版本不兼容会导致类加载问题,也可能间接影响数据源加载。

替换Iceberg runtime依赖为对应Scala 2.13、Spark 3.3的版本:

<dependency>
    <groupId>org.apache.iceberg</groupId>
    <artifactId>iceberg-spark-runtime-3.3_2.13</artifactId>
    <version>1.4.3</version>
</dependency>

3. 修正结构化流代码逻辑

结构化流(readStream)不能直接调用show(),需要启动流查询并指定输出目标。修改Main类代码,添加流输出逻辑:

package com.dell;

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.streaming.StreamingQuery;
import org.apache.spark.sql.streaming.StreamingQueryException;

public class Main {
    public static void main(String[] args) throws StreamingQueryException {
        SparkSession spark = SparkSession.builder()
                .appName("Kafka Spark Demo")
                .master("local[*]")
                .getOrCreate();
                
        String kafkaBrokers= "localhost:9092";
        String kafkaTopic = "Topic_Project";
        
        Dataset<Row> kafkaStreamDf = spark
                .readStream()
                .format("kafka")
                .option("kafka.bootstrap.servers", kafkaBrokers)
                .option("subscribe", kafkaTopic)
                .option("startingOffsets", "earliest")
                .load();

        // 启动流查询,输出到控制台
        StreamingQuery query = kafkaStreamDf.writeStream()
                .outputMode("append")
                .format("console")
                .start();
                
        query.awaitTermination();
    }
}

4. 检查Docker Kafka的网络配置

确保Docker中运行的Kafka服务:

  • 已正确映射9092端口到主机(Docker run命令包含-p 9092:9092)
  • Kafka的advertised.listeners配置为PLAINTEXT://localhost:9092,否则Spark会尝试连接容器内部地址,导致无法通信

内容的提问来源于stack exchange,提问作者AROY

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 14:42:03