Spark Structured Streaming读取Kafka及spark-submit提交报错求助
问题原因及解决方案
第一个报错(找不到kafka数据源)
有两个触发原因:
- 代码存在语法错误:
.option("subscribe, "mytopic")中subscribe参数的双引号未闭合,正确写法为.option("subscribe", "mytopic") - 运行环境缺少Spark Structured Streaming和Kafka的集成依赖包
第二个报错(找不到主类/文件不存在)
你提交命令中的...是官方文档的占位符,需要替换为你实际的应用程序路径,Spark提交任务时必须指定要运行的程序文件。
具体操作步骤
- 先修正业务代码的语法问题,修正后代码如下:
df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "mytopic") \ .load()
- 根据你的开发语言调整提交命令:
- 如果你用的是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
相关产品推荐
相关产品推荐

