【发布时间】: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