【发布时间】:2017-03-08 03:35:03
【问题描述】:
我正在使用 Google Cloud Dataflow 并且有一个 ParDo 函数,该函数需要访问 PCollection 中的所有元素。为此,我想将 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));
}
}));
这个任务似乎有一个预定义的组合函数,但我找不到。感谢您的帮助!
【问题讨论】: