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
相关产品推荐
相关产品推荐

