【问题标题】:Using mqtt with pyspark streaming将 mqtt 与 pyspark 流结合使用
【发布时间】: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


【解决方案1】:

对于 spark 版本 2.*,我们可以通过包含 Bahir Jar 在Structured Streaming 中使用 MQTT。

从 pyspark 连接到 MQTT 代理:

(spark
    .readStream
    .format("org.apache.bahir.sql.streaming.mqtt.MQTTStreamSourceProvider")
    .option("topic","mytopic")
    .load("tcp://{}".format(broker_uri)))

【讨论】:

    猜你喜欢
    • 2019-03-03
    • 2018-01-24
    • 2020-05-23
    • 1970-01-01
    • 2020-05-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多