您可以在 apache Bahir 中看到以下拉取请求:Bahir Pull Request
您可以在哪里看到 MQTTUtils.createPairedStream 正在添加。
您使用以下工件导入您的 pom/gradle/sbt...:
spark-sql-streaming-mqtt_2.11 版本 2.3.2 来自组 org.apache.bahir。
您可以在 maven 中使用到 Spark 1.6:
<!-- https://mvnrepository.com/artifact/org.apache.spark/spark-streaming-mqtt -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming-mqtt_2.11</artifactId>
<version>1.6.3</version>
</dependency>
对于 Spark 2.3.2,您需要使用:
<dependency>
<groupId>org.apache.bahir</groupId>
<artifactId>spark-streaming-mqtt_2.11</artifactId>
<version>2.3.2</version>
</dependency>
或在 SBT 中:
libraryDependencies += "org.apache.bahir" %% "spark-streaming-mqtt" % "2.3.2"
您可以在以下位置找到更多信息:org.apache.bahir:spark-streaming-mqtt
bin/spark-shell --packages org.apache.bahir:spark-streaming-mqtt_2.11:2.3.0
您将使用 scala 导入包:
import org.apache.spark.streaming.mqtt._
并实例化:
val lines = MQTTUtils.createPairedStream(ssc, brokerUrl, topic)
我希望这会有所帮助。