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

引入bahir的spark-streaming-mqtt包仍找不到MQTTStreamSinkProvider类如何解决

问题排查与解决步骤

1. 修正依赖配置

你当前引入的org.apache.bahir:spark-streaming-mqtt是针对Spark DStream API的依赖包,不包含Structured Streaming需要的MQTT Sink/Source实现。需要替换为Structured Streaming对应的依赖:
在build.sbt中修改依赖为:

"org.apache.bahir" %% "spark-sql-streaming-mqtt" % "2.4.0"

注意依赖版本需要和你当前使用的Spark版本完全匹配,避免兼容性问题。

2. 验证依赖是否正确加载

  • 本地IDE运行场景:执行sbt clean compile后,检查项目依赖库中是否存在spark-sql-streaming-mqtt对应的jar包,确认类路径下存在org.apache.bahir.sql.streaming.mqtt.MQTTStreamSinkProvider类
  • 提交Spark集群运行场景:需要将依赖打入提交的fat jar,或者在spark-submit时通过--packages参数指定依赖:
spark-submit --packages org.apache.bahir:spark-sql-streaming-mqtt_2.11:2.4.0 其他提交参数

注意替换命令中的Scala版本为你项目使用的Scala主版本。

3. 简化格式声明(可选)

依赖配置正确后,format参数可以简化为mqtt,无需写全类名:

df
  .writeStream
  .format("mqtt")
  .outputMode("complete")
  .option("topic", "mytopic")
  .option("brokerUrl", "tcp://localhost:1883")
  .start()
  .awaitTermination(20000)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 05:36:01