【问题标题】:Pyspark custom receivers to read mongo change stream logs using spark streamingPyspark 自定义接收器使用火花流读取 mongo 更改流日志
【发布时间】: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 更改流数据。

[1]:https://docs.mongodb.com/manual/changeStreams/

【问题讨论】:

  • 如果遇到这种情况,谁能帮忙
  • 对我们有用的解决方案是分两步,而不是直接尝试通过 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


【解决方案1】:

对我们有用的解决方案出现在评论部分

【讨论】:

    猜你喜欢
    • 2020-04-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多