【发布时间】:2017-06-01 19:00:12
【问题描述】:
是否可以在批处理 Dataflow 作业处理完所有数据后执行操作?具体来说,我想将管道刚刚处理的文本文件移动到不同的 GCS 存储桶。我不确定将其放置在我的管道中的哪个位置,以确保它在数据处理完成后执行一次。
【问题讨论】:
是否可以在批处理 Dataflow 作业处理完所有数据后执行操作?具体来说,我想将管道刚刚处理的文本文件移动到不同的 GCS 存储桶。我不确定将其放置在我的管道中的哪个位置,以确保它在数据处理完成后执行一次。
【问题讨论】:
我不明白您为什么需要执行此后期管道执行。您可以使用侧输出将文件写入多个存储桶,并在管道完成后保存自己的副本。
如果这对您不起作用(无论出于何种原因),那么您可以简单地在blocking execution 模式下运行您的管道,即使用pipeline.run().waitUntilFinish(),然后只需编写其余代码(执行复制)之后。
[..]
// do some stuff before the pipeline runs
Pipeline pipeline = ...
pipeline.run().waitUntilFinish();
// do something after the pipeline finishes here
[..]
【讨论】:
BlockingDataflowPipelineRunner 运行作业就可以了。 waitUntilFinish() 在 1.x Java API 中似乎不可用。
我从阅读apache beam的PassThroughThenCleanup.java的源代码中得到的一个小技巧。
在您的阅读器之后,创建一个“组合”整个集合的侧输入(在源代码中,它是View.asIterable() PTransform)并将其输出连接到DoFn。这个DoFn 只有在读者读完所有元素后才会被调用。
附:代码按字面意思命名操作,cleanupSignalView 我发现它非常聪明
请注意,您可以使用Combine.globally() (java) 或beam.CombineGlobally() (python) 实现相同的效果。欲了解更多信息,请查看第 4.2.4.3 节 here
【讨论】:
我认为这里有两个选项可以帮助您:
1) 使用TextIO 写入您想要的存储桶或文件夹,指定确切的 GCS 路径(例如 gs://sandbox/other-bucket)
2) 将Object Change Notifications 与Cloud Functions 结合使用。您可以在 here 和 JS 中的 GCS 开发工具包 here 中找到很好的入门指南。您在此选项中所做的基本上是在某个存储桶中掉落某物时设置触发器,然后使用您自己编写的 Cloud Function 将其移动到另一个存储桶。
【讨论】: