【问题标题】:Writing to text files in Apache Beam / Dataflow Python streaming在 Apache Beam / Dataflow Python 流中写入文本文件
【发布时间】:2018-11-19 17:10:52
【问题描述】:

我有一个非常基本的 Python 数据流作业,它从 Pub/Sub 读取一些数据,应用 FixedWindow 并写入 Google Cloud Storage。

transformed = ...
transformed | beam.io.WriteToText(known_args.output)

输出写入--output中指定的位置,但只是临时阶段,即

gs://MY_BUCKET/MY_DIR/beam-temp-2a5c0e1eec1c11e8b98342010a800004/...some_UUID...

文件永远不会使用分片模板放置到正确命名的位置。

在本地和 DataFlow 运行器上测试。


当进一步测试时,我注意到 streaming_wordcount 示例有同样的问题,但是标准 wordcount 示例写得很好。也许问题在于窗口化或从 pubsub 读取?


WriteToText 似乎与 PubSub 的流媒体源不兼容。可能有解决方法,或者 Java 版本可能兼容,但我选择完全使用不同的解决方案。

【问题讨论】:

  • 可以发一下代码吗?
  • 这里有同样的问题。你找到解决办法了吗?我尝试使用触发器,但它与触发器无关。当我在 Dataflow 作业中调用“Drain”时,数据将写入正确的文件夹中。
  • 试试fileio.WriteToFiles

标签: google-cloud-storage google-cloud-dataflow apache-beam google-cloud-pubsub


【解决方案1】:

Python SDK 中的WriteToText 转换不支持流式传输。

相反,您可以考虑apache_beam.io.fileio 中的转换。在这种情况下,你可以这样写(假设 10 分钟的窗口):

my_pcollection = (p | ReadFromPubSub(....)
                    |  WindowInto(FixedWindows(10*60))
                    |  fileio.WriteToFiles(path=known_args.output))

这足以为每个窗口写出单独的文件,并随着流的前进继续这样做。

你会看到这样的文件(假设输出是gs://mybucket/)。这些文件将在窗口被触发时被打印出来:

gs://mybucket/output-1970-01-01T00_00_00-1970-01-01T00_10_00-0000-00002
gs://mybucket/output-1970-01-01T00_00_00-1970-01-01T00_10_00-0001-00002
gs://mybucket/output-1970-01-01T00_10_00-1970-01-01T00_20_00-0000-00002
gs://mybucket/output-1970-01-01T00_10_00-1970-01-01T00_20_00-0001-00002
...

默认情况下,文件具有$prefix-$start-$end-$pane-$shard-of-$numShards$suffix$compressionSuffix 名称 - 其中前缀默认为output,但您可以传递更复杂的文件命名函数。


如果您想自定义文件的写入方式(例如,文件的命名、数据的格式等),您可以查看WriteToFiles 中的额外参数。

您可以看到在 Beam 测试中使用的转换示例 here,其中包含更复杂的参数 - 但听起来默认行为对您来说应该足够了。

【讨论】:

  • 谢谢,我试图弄清楚为什么我不能在我的代码中使用 io.fileio。可能存在一些版本问题,我使用 Python 3.7 和 SDK 2.16.0,但 fileio 不存在。这篇文章告诉我它已重命名,但文件系统模块没有 WriteToText stackoverflow.com/questions/46787428/…
  • 让我知道这是否有效,或者如果您需要帮助以找到正确的参数集。
  • 关于窗口大小:窗口使用的时间戳是多少?
【解决方案2】:

Python 流式管道执行在实验上可用(有一些限制)。

不支持的功能适用于所有跑步者。 状态和计时器 API, 自定义源 API, 可拆分的 DoFn API, 处理迟到的数据, 用户自定义的自定义 WindowFn.

此外,DataflowRunner 目前不支持以下具有 Python 流式执行的 Cloud Dataflow 特定功能。

流式自动缩放 更新现有管道 云数据流模板 一些监控功能,例如毫秒计数器、显示数据、指标和转换的元素计数。但是,支持源的日志记录、水印和元素计数。

https://beam.apache.org/documentation/sdks/python-streaming/

由于您使用的是 FixedWindowFn 并且管道能够将数据写入 tmp 位置,请重新检查输出位置--output gs://<your-gcs-bucket>/<you-gcs-folder>/<your-gcs-output-filename>

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-03-31
    • 1970-01-01
    • 2018-01-03
    • 1970-01-01
    • 1970-01-01
    • 2017-09-03
    • 2018-10-17
    • 1970-01-01
    相关资源
    最近更新 更多