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

Spark Structured Streaming读取Kafka及spark-submit提交报错求助

问题原因及解决方案

第一个报错(找不到kafka数据源)

有两个触发原因:

  • 代码存在语法错误:.option("subscribe, "mytopic")中subscribe参数的双引号未闭合,正确写法为.option("subscribe", "mytopic")
  • 运行环境缺少Spark Structured Streaming和Kafka的集成依赖包

第二个报错(找不到主类/文件不存在)

你提交命令中的...是官方文档的占位符,需要替换为你实际的应用程序路径,Spark提交任务时必须指定要运行的程序文件。


具体操作步骤

  1. 先修正业务代码的语法问题,修正后代码如下:
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "mytopic") \
.load()
  1. 根据你的开发语言调整提交命令:
  • 如果你用的是Python开发,假设你的脚本文件名为kafka_stream.py,存放在/home/myname/目录下,提交命令为:
./bin/spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2 /home/myname/kafka_stream.py
  • 如果你用的是Scala/Java开发,假设你打的jar包名为kafka-stream.jar,主类全限定名为com.example.KafkaStreamApp,jar包存放在/home/myname/目录下,提交命令为:
./bin/spark-submit --class com.example.KafkaStreamApp --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2 /home/myname/kafka-stream.jar

注意:spark-sql-kafka的版本号必须和你当前使用的Spark版本完全一致,后缀的_2.12为Scala版本,需和Spark安装包对应的Scala版本匹配,你使用的spark-3.1.2-bin-hadoop3.2默认对应Scala2.12,因此依赖配置无需调整。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 23:09:01