【发布时间】:2018-06-04 20:54:52
【问题描述】:
参考项目:Watching for new files matching a filepattern in Apache Beam
您可以将它用于简单的用例吗?我的用例是我让用户将数据上传到 Cloud Storage -> Pipeline (Process csv to json) -> Big Query。我知道云存储是有界集合,因此它代表批处理数据流。
我想做的是让管道以流模式运行,一旦文件上传到 Cloud Storage,它将通过管道进行处理。 watchfornewfiles 可以做到这一点吗?
我的代码如下:
p.apply(TextIO.read().from("<bucketname>")
.watchForNewFiles(
// Check for new files every 30 seconds
Duration.standardSeconds(30),
// Never stop checking for new files
Watch.Growth.<String>never()));
没有任何内容被转发到 Big Query,但管道显示它正在流式传输。
【问题讨论】:
标签: google-cloud-platform google-cloud-dataflow apache-beam