【发布时间】:2021-06-04 19:34:22
【问题描述】:
我想使用 spark 流从 mongodb 更改流中读取数据[链接在最后]。
这里想收集 30 秒的转储,然后推送到某个文件中。 我知道我可能需要编写一些自定义接收器(使用 pyspark)来从相关数据源接收数据,但我找不到任何讨论使用 PYTHON 进行 Spark Streaming 自定义接收器的文档。
下面的文档链接以及使用 java 或 scala 的提及。
http://spark.apache.org/docs/latest/streaming-custom-receivers.html
我正在使用简单的 python 代码从 ChangeStreams 中读取数据,但它不能满足我的要求。
注意:在下面的代码中,使用 for 循环逐一迭代 change_stream(而不是希望批量读取 30 秒时间范围内的文档,然后将其写入某个目标文件)
import os
import pymongo
from bson.json_util import dumps
STREAM_DB="mongodb://<username>:<pwd>@<host>:<port>/<database to be used> authSource=admin&retryWrites=true"
client = pymongo.MongoClient(STREAM_DB)
change_stream = client.<database name>.watch()
print(change_stream)
f = open("<filename>", "a")
for change in change_stream:
f.write(dumps(change) + '\n')
f.close()
我没有看到任何讨论使用 PYTHON 进行 Spark Streaming 的自定义接收器的文档。 下面的文档链接以及使用 java 或 scala 的提及。
http://spark.apache.org/docs/latest/streaming-custom-receivers.html
有没有办法我可以使用 spark 读取流式 mongodb 更改流数据。
【问题讨论】:
-
如果遇到这种情况,谁能帮忙
-
对我们有用的解决方案是分两步,而不是直接尝试通过 Spark。 1) 使用作为服务运行的 python 独立脚本将 Mongo 更改流数据发送到 Google Cloud PubSub。 2) 在 java 中编写自定义接收器并使用 pyspark (Spark Streaming 代码) 访问它 这个 github 链接具有专门用于 Spark-PubSub 链接的自定义接收器代码:github.com/SignifAi/Spark-PubSub
标签: python python-3.x mongodb apache-spark pyspark