Gra*_*ley 1 google-cloud-dataflow
我们在BigQuery中有一个大表,其中数据正在流入.每天晚上,我们都想运行处理过去24小时数据的Cloud Dataflow管道.
在BigQuery中,可以使用" 表装饰器 " 执行此操作,并指定我们想要的范围,即24小时.
从BQ表读取时,Dataflow中是否可以以某种方式实现相同的功能?
我们已经看过Dataflow 的' Windows '文档,但我们无法确定这是否是我们需要的.到目前为止我们想出了这个(我们希望最后24小时的数据使用FixedWindows),但它仍然试图读取整个表格:
pipeline.apply(BigQueryIO.Read
.named("events-read-from-BQ")
.from("projectid:datasetid.events"))
.apply(Window.<TableRow>into(FixedWindows.of(Duration.standardHours(24))))
.apply(ParDo.of(denormalizationParDo)
.named("events-denormalize")
.withSideInputs(getSideInputs()))
.apply(BigQueryIO.Write
.named("events-write-to-BQ")
.to("projectid:datasetid.events")
.withSchema(getBigQueryTableSchema())
.withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE) .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED));
Run Code Online (Sandbox Code Playgroud)
我们走在正确的轨道上吗?
小智 5
谢谢你的问题.
此时,BigQueryIO.Read需要"project:dataset:table"格式的表信息,因此指定装饰器将不起作用.
在支持此功能之前,您可以尝试以下方法:
希望这可以帮助
| 归档时间: |
|
| 查看次数: |
190 次 |
| 最近记录: |