【问题标题】:Filtering bounded data in Dataflow based on timestamp根据时间戳过滤Dataflow中的有界数据
【发布时间】:2016-06-11 18:10:48
【问题描述】:

在我的数据流管道中,我将有两个从 BigQuery 表中读取的 PCollections<TableRow>。我计划将这两个 PCollection 合并为一个 PCollection 和一个 flatten

由于 BigQuery 仅追加,因此目标是使用新的 PCollection 截断 BigQuery 中的第二个表。

我已通读文档,这是我感到困惑的中间步骤。对于我的新PCollection,计划是使用Comparator DoFn 查看最大上次更新日期并返回给定行。 我不确定是否应该使用过滤器转换,或者是否应该先按键分组然后使用过滤器?

所有PCollection<TableRow>s 都将包含相同的值:IE:字符串、整数和时间戳。当谈到键值对时,大多数关于云数据流的文档都只包含简单的字符串。 是否可以有一个键值对是PCollection<TableRow> 的整行?

这些行看起来类似于:

customerID, customerName, lastUpdateDate
0001, customerOne, 2016-06-01 00:00:00
0001, customerOne, 2016-06-11 00:00:00

在上面的示例中,我希望过滤 PCollection 以仅将第二行返回到将写入 BigQuery 的 PCollection。 另外,是否可以在不创建第四个的情况下将这些 Pardo 应用于第三个 PCollection?

【问题讨论】:

    标签: java google-cloud-dataflow


    【解决方案1】:

    你问了几个问题。我试图孤立地回答它们,但我可能误解了整个场景。如果您提供了一些示例代码,可能有助于澄清。

    对于我的新 PCollection,计划是使用 Comparator DoFn 查看最大上次更新日期并返回给定行。我不确定是否应该使用过滤器转换,或者是否应该按键分组然后使用过滤器?

    根据您的描述,您似乎想要获取 PCollection 的元素,并为每个 customerID(键)找到该客户记录的最新更新。您可以通过Top.largestPerKey(1, timestampComparator) 使用提供的转换来完成此操作,您可以在其中设置timestampComparator 以仅查看时间戳。

    是否可以有一个键值对是 PCollection 的整行?

    KV<K, V> 可以具有任何类型的键 (K) 和值 (V)。如果要按键分组,则键的编码器需要是确定性的。 TableRowJsonCoder 不是确定性的,因为它可能包含任意对象。但听起来您希望将customerID 用作键,将整个TableRow 用作值。

    是否可以在不创建第四个的情况下将这些 Pardo 应用于第三个 PCollection?

    当您将PTransform 应用于PCollection 时,会生成一个新的PCollection。没有办法解决这个问题,您无需尝试尽量减少管道中PCollections 的数量。

    PCollection 是一个概念对象;它没有内在成本。您的管道将被大量优化,因此许多中间PCollections - 尤其是ParDo 转换序列中的那些 - 无论如何都不会实现。

    【讨论】:

      猜你喜欢
      • 2018-09-11
      • 1970-01-01
      • 2017-09-08
      • 1970-01-01
      • 2023-02-13
      • 1970-01-01
      • 1970-01-01
      • 2021-11-22
      • 1970-01-01
      相关资源
      最近更新 更多