【发布时间】:2020-11-06 11:38:05
【问题描述】:
我在 Apache Beam Go SDK 中有一个管道。
pcoll := beam.GroupByKey(s, src)
问题是,在 GroupByKey 转换之后,我想用 ParDo 转换进一步处理它。我有类型的问题,因为 Go 要我按如下方式定义 ParDo 函数输入:
value CoGBK<[]uint8,[]uint8>
但是 Go 中没有 CoGBK 类型。有没有办法在 Apache Beam Go SDK 中定义这种类型?
【问题讨论】:
-
要考虑的一件事:不幸的是,Beam 的 Go SDK 的文档在这方面写得很草率:它们似乎混合了 Java(任何 Apache 托管项目的通用语)和 Go 的术语,并且只有当读者相当精通这两种语言并且他们理解
CoGBK<[]uint8,[]uint8>是一个抛物线,而不是一个真实类型的定义时,这才有效。 -
感谢您的建议我遵循了他们,但是我仍然卡住了。
pcoll := beam.GroupByKey(s, src)现在你想在 PCollection 上应用另一个转换,比如说 ParDores, err := beam.TryParDo(s, &exampleFn{}, pcoll)和 exampleFnfunc (fn *exampleFn) ProcessElement (x beam.PCollection, y beam.PCollection, emit func([]bytes))Go SDK 期望 x,y 类型匹配 CoGBK,但我可以'在任何地方都找不到如何实现,learning/katas中有例子,但是在GroupByKey转换后他们不处理任何东西。 -
好的,当我们收到需要
CoGBK<[]uint8,[]uint8>类型的值的消息时,我们实际上应该做的是应用以下转换:beam.ParDo0(s, func(key []uint8, values func(*[]uint8) bool) {}, pcoll)Go 将其解释为 CoGBK。此问题应标记为已解决。
标签: go apache-beam dataflow