【问题标题】:Simple approach to combine PCollection<T> into PCollection<Iterable<T>>将 PCollection<T> 组合成 PCollection<Iterable<T>> 的简单方法
【发布时间】:2017-03-08 03:35:03
【问题描述】:

我正在使用 Google Cloud Dataflow 并且有一个 ParDo 函数,该函数需要访问 PCollection 中的所有元素。为此,我想将 PCollection 转换为包含所有元素的单个 Iterable 的 PCollection>。我想知道是否有更清洁/更简单/更快的解决方案来解决我的想法。

第一种方法是创建一个虚拟键,执行 GroupByKey,然后获取值。

PCollection<MyType> myData;
// AddDummyKey() outputs KV.of(1, context.element()) for everything
PCollection<KV<Integer, MyType>> myDataKeyed = myData.apply(ParDo.of(new AddDummyKey())); 
// Group by dummy key
PCollection<KV<Integer, Iterable<MyType>>> myDataGrouped = myDataKeyed.apply(GroupByKey.create());
// Extract values
PCollection<Iterable<MyType>> myDataIterable = myDataGrouped.apply(Values.<Iterable<MyType>>create()

第二种方法遵循此处的建议:How do I make View's asList() sortable in Google Dataflow SDK?,但没有排序。我创建了一个 View.asList(),创建了一个虚拟 PCollection,然后在虚拟 PCollection 上应用 ParDo 函数,并将视图作为侧面输入,然后简单地返回视图。

PCollection<MyType> myData;
// Create view of the PCollection as a list
PCollectionView<List<MyType>> myDataView = myData.apply(View.asList()); 
// Create dummy PCollection
PCollection<Integer> dummy = pipeline.apply(Create.<Integer>of(1));
// Apply dummy ParDo that returns the view
PCollection<List<MyType>> myDataList = dummy.apply(
        ParDo.withSideInputs(myDataView).of(new DoFn<Integer, List<MyType>>() {
            @Override
            public void processElement(ProcessContext c) {
                c.output(c.sideInput(myDataView)); 
            }
        }));

这个任务似乎有一个预定义的组合函数,但我找不到。感谢您的帮助!

【问题讨论】:

    标签: google-cloud-dataflow


    【解决方案1】:

    如果您知道自己需要全部内容,那么您的两种方法都是合理的。两者都已在 Dataflow SDK 中使用,后来成为 Apache Beam SDK。

    1. 侧面输入然后输出整个事情:这就是DataflowAssert 的工作原理,事实上。在 Beam 中,不同的后端运行器可能以不同的方式实现侧输入,您应该更喜欢 View.asIterable(),因为它的假设更少,并且可能允许对非常大的侧输入进行更多流式传输。
    2. 按单个键分组,然后放下该键:这就是 Beam 的继任者 PAssert 的工作原理。它完成了同样的事情,需要对空集合多加注意,但更多的 Beam runner 拥有良好的 GroupByKey 支持而不是侧输入支持(尤其是当它们是新的且仍在开发中时)。

    所以View.asIterable() 基本上就是您所要求的。还有一些要求进行第二个版本的GroupGlobally 转换;这可能在某个时候发生。

    【讨论】:

      【解决方案2】:

      到目前为止,更重要的方法是将 Combine 与 AccumulatorFn 一起使用,例如:

      https://beam.apache.org/releases/javadoc/2.8.0/org/apache/beam/sdk/transforms/Combine.AccumulatingCombineFn.html

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-08-13
        • 1970-01-01
        • 1970-01-01
        • 2021-03-23
        相关资源
        最近更新 更多