【问题标题】:Group by key collection type in Apache Beam pipeline按 Apache Beam 管道中的键集合类型分组
【发布时间】: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 中定义这种类型?

【问题讨论】:

  • 基于this 和this,我会说你的声明«Go希望我定义ParDo函数输入如下:value CoGBK&lt;[]uint8,[]uint8&gt;»作为beam.ParDo的第三个参数是错误的是PCollection,而这正是beam.[Co]GroupByKey 返回的内容。
  • 要考虑的一件事:不幸的是,Beam 的 Go SDK 的文档在这方面写得很草率:它们似乎混合了 Java(任何 Apache 托管项目的通用语)和 Go 的术语,并且只有当读者相当精通这两种语言并且他们理解CoGBK&lt;[]uint8,[]uint8&gt; 是一个抛物线,而不是一个真实类型的定义时,这才有效。
  • 感谢您的建议我遵循了他们,但是我仍然卡住了。 pcoll := beam.GroupByKey(s, src) 现在你想在 PCollection 上应用另一个转换,比如说 ParDo res, err := beam.TryParDo(s, &amp;exampleFn{}, pcoll) 和 exampleFn func (fn *exampleFn) ProcessElement (x beam.PCollection, y beam.PCollection, emit func([]bytes)) Go SDK 期望 x,y 类型匹配 CoGBK,但我可以'在任何地方都找不到如何实现,learning/katas中有例子,但是在GroupByKey转换后他们不处理任何东西。
  • 好的,当我们收到需要CoGBK&lt;[]uint8,[]uint8&gt; 类型的值的消息时,我们实际上应该做的是应用以下转换:beam.ParDo0(s, func(key []uint8, values func(*[]uint8) bool) {}, pcoll) Go 将其解释为 CoGBK。此问题应标记为已解决。

标签: go apache-beam dataflow


【解决方案1】:

好的,当我们收到一条需要CoGBK&lt;[]uint8,[]uint8&gt; 类型的消息时,我们实际上应该做的是应用以下转换: beam.ParDo0(s, func(key []uint8, values func(*[]uint8) bool) {}, pcoll) 这被 Go 解释为 CoGBK&lt;[]uint8,[]uint8&gt;。

【讨论】:

    猜你喜欢
    • 2018-01-05
    • 1970-01-01
    • 1970-01-01
    • 2018-11-05
    • 1970-01-01
    • 2021-11-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多