【问题标题】:Writing to a BigQuery table with date in table name from a DataFlow streaming pipeline从 DataFlow 流式传输管道写入表名称中包含日期的 BigQuery 表
【发布时间】:2018-01-12 16:38:59
【问题描述】:

我的表名格式:tableName_YYYYMMDD。我正在尝试从流式数据流管道写入此表。我想每天写一个新表的原因是因为我想在 30 天后使表过期,并且只想一次保留 30 个表的窗口。

当前代码:

tableRow.apply(BigQueryIO.Write
                .named("WriteBQTable")
                .to(String.format("%1$s:%2$s.%3$s",projectId, bqDataSet, bqTable))
                .withSchema(schema)
                .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
                .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND));

我确实意识到上面的代码不会滚动到新的一天并开始在那里编写。

正如this 的回答所暗示的,我可以对表进行分区和使分区过期,但是流式管道似乎不支持写入分区表。

有什么想法可以解决这个问题吗?

【问题讨论】:

    标签: google-bigquery google-cloud-platform google-cloud-dataflow


    【解决方案1】:

    在 Dataflow 2.0 SDK 中有一种方法可以指定 DynamicDestinations

    BigQuery Dynamic Destionations 中查看to(DynamicDestinations<T,?> dynamicDestinations)

    另外,请参阅 TableDestination 版本,它应该更简单且代码更少。虽然不幸的是 javadoc 中没有示例。

    to(SerializableFunction<ValueInSingleWindow<T>,TableDestination> tableFunction)
    

    https://beam.apache.org/documentation/sdks/javadoc/2.0.0/

    【讨论】:

    • 我一直在使用从数据流写入的分区表对此进行测试,看起来大查询正在正确分配值,我没有从数据流分配任何_PARTITIONTIME。下面返回正确的数据。 SELECT * FROM Mytable WHERE _PARTITIONTIME = TIMESTAMP("2017-08-04")
    【解决方案2】:

    This 是一个开源管道,可用于将 pub/sub 连接到大查询。我认为谷歌还添加了对流式管道的支持以支持日期分区表。详情here.

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-10-25
      • 2019-05-01
      • 2020-11-12
      • 1970-01-01
      • 1970-01-01
      • 2018-10-18
      • 1970-01-01
      相关资源
      最近更新 更多