【问题标题】:Does anyone have an example of a generic ProcessFunction in Flink?有人有 Flink 中通用 ProcessFunction 的示例吗?
【发布时间】:2020-04-30 19:21:38
【问题描述】:

“通用”是指能够接受任何类型的对象作为输入并返回相同的对象作为输出。

假设该函数的工作是将每个元素序列化为 json 并将其作为辅助输出写入。

class MyProcessFunction() extends ProcessFunction[? , ?] {

    def processElement(element: ?, ctx: ProcessFunction[?, ?]#Context, out: Collector[?]): Unit = ??? 

    ... 
}

我能否以这样一种方式定义它,使其可供不同类型的输入使用?

【问题讨论】:

    标签: scala templates generics apache-flink


    【解决方案1】:

    您可以通过将 Your 类设为通用来做到这一点。所以,你会有类似的东西:

    class MyProcessFunction[T] extends ProcessFunction[T, T] {
      override def processElement(value: T, ctx: ProcessFunction[T, T]#Context, out: Collector[T]): Unit = ???
    }
    

    这样,您将能够在创建函数实例时确定类型。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-09-27
      • 2012-04-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多