【发布时间】:2018-12-10 21:44:21
【问题描述】:
具体来说,Beam 中的Flatten PTransform 是否执行任何类型的操作:
- 重复数据删除
- 过滤
- 清除现有元素
还是只是“合并”两个不同的 PCollection?
【问题讨论】:
标签: google-cloud-dataflow apache-flink apache-beam
具体来说,Beam 中的Flatten PTransform 是否执行任何类型的操作:
还是只是“合并”两个不同的 PCollection?
【问题讨论】:
标签: google-cloud-dataflow apache-flink apache-beam
Flatten 转换不执行任何类型的重复数据删除或任何类型的过滤。如前所述,它只是将多个 PCollection 合并为一个,其中包含每个输入的元素。
这意味着:
with beam.Pipeline() as p:
c1 = p | "Branch1" >> beam.Create([1, 2, 3, 4])
c2 = p | "Branch2" >> beam.Create([4, 4, 5, 6])
result = (c1, c2) | beam.Flatten()
在这种情况下,result PCollection 包含以下元素:[1, 2, 3, 4, 4, 4, 5, 6]。
注意元素4 在c1 中出现一次,在c2 中出现两次。这不会以任何方式进行重复数据删除、过滤或删除。
关于Flatten 的一个奇怪事实是,一些跑步者对其进行了优化,并在两个分支中简单地添加了下游转换。因此,简而言之,没有特殊的过滤或重复数据删除。只需合并 PCollections。
【讨论】: