【发布时间】:2017-10-06 13:41:10
【问题描述】:
我是 spark 和 mqtt 的新手。我正在尝试使用我在网上获得的名为 wordcount.py 的 MQTTUtils 代码
import sys
from pyspark import SparkContext
from pyspark.streaming import StreamingContext
from pyspark.streaming.mqtt import MQTTUtils
if __name__ == "__main__":
if len(sys.argv) != 3:
print >> sys.stderr, "Usage: mqtt_wordcount.py <broker url> <topic>"
exit(-1)
sc = SparkContext(appName="PythonStreamingMQTTWordCount")
ssc = StreamingContext(sc, 1)
brokerUrl = sys.argv[1]
topic = sys.argv[2]
lines = MQTTUtils.createStream(ssc, brokerUrl, topic)
counts = lines.flatMap(lambda line: line.split(" ")) \
.map(lambda word: (word, 1)) \
.reduceByKey(lambda a, b: a+b)
counts.pprint()
ssc.start()
ssc.awaitTermination()
我按照说明安装了 mosquitto 代理(它正在工作),下载 spark-streaming-mqtt-assembly_2.11-1.6.2.jar 并使用以下命令运行 python 脚本: ~$ spark-submit --jars spark-streaming-mqtt-assembly_*.jar wordcount.py
但显示的错误:
从 pyspark.streaming.mqtt 导入 MQTTUtils
ImportError: 没有名为 mqtt 的模块
我错过了什么吗? 谢谢
【问题讨论】:
-
如何创建minimal reproducible example。 Spark 2.0+ 也不再提供 MQTT 后端。它已移至 Spark 包中。
-
我遇到了同样的问题,但我使用的是 2.0 版,现在我使用的是 1.6.2 版并且脚本正在运行。
标签: apache-spark spark-streaming mqtt