【问题标题】:How to use delta trigger in flink?如何在 flink 中使用 delta 触发器?
【发布时间】:2017-12-05 20:03:07
【问题描述】:

我想在 apache flink (flink 1.3) 中使用 deltatrigger,但我在这段代码中遇到了一些问题:

.trigger(DeltaTrigger.of(100, new DeltaFunction[uniqStruct] {
    override def getDelta(oldFp: uniqStruct, newFp: uniqStruct): Double = newFp.time - oldFp.time
  }, TypeInformation[uniqStruct]))

我有这个错误:

error: object org.apache.flink.api.common.typeinfo.TypeInformation is not a value [ERROR] }, TypeInformation[uniqStruct]))

我不明白为什么 DeltaTrigger 需要 TypeSerializer[T] 而且我不知道该怎么做才能消除此错误。

非常感谢大家。

【问题讨论】:

    标签: scala apache-flink flink-streaming


    【解决方案1】:

    我会稍微读一下https://ci.apache.org/projects/flink/flink-docs-release-1.2/dev/types_serialization.html 听起来你可以在你的类型信息上使用typeInfo.createSerializer(config) 创建一个序列化程序。请注意,您当前传入的是类型本身,而不是类型信息,这就是您收到错误的原因。

    你需要做一些类似的事情

    val uniqStructTypeInfo: TypeInformation[uniqStruct] = createTypeInformation[uniqStruct]
    val uniqStrictTypeSerializer = typeInfo.createSerializer(config)
    

    引用上面关于创建序列化程序需要传递的配置参数的页面

    config 参数的类型为 ExecutionConfig,并保存 有关程序注册的自定义序列化程序的信息。在哪里 有可能,尝试通过适当的 ExecutionConfig 程序。你 通常可以通过调用从DataStream或DataSet中获取 获取执行配置()。内部函数(如 MapFunction),你可以得到 通过使函数成为丰富的函数并调用 getRuntimeContext().getExecutionConfig()。

    【讨论】:

    • getExecutionConfig 在 flink 1.3 中已弃用,使用 RichFunction 我有此错误:Cannot resolve symbol getRuntimeContext。或者没有其他方法可以每 n 毫秒触发一次?
    【解决方案2】:

    DeltaTrigger 需要一个TypeSerializer,因为它使用 Flink 的托管状态机制来存储每个元素以供以后与下一个比较(它只保留一个元素,即最后一个元素,它会随着新元素的到来而更新)。

    您将找到一个示例(Java 语言)here

    但如果您只需要一个每 100 毫秒触发一次的窗口,那么使用TimeWindow 会更容易,例如

    input
      .keyBy(<key selector>)
      .timeWindow(Time.milliseconds(100)))
      .apply(<window function>)
    

    更新:

    要拥有每 100 毫秒触发一次的长达一小时的窗口,您可以使用滑动窗口。但是,您将拥有 10 * 60 * 60 个窗口,并且每个事件都将放置在这 36000 个窗口中的每一个中。所以这不是一个好主意。

    如果您将GlobalWindowDeltaTrigger 一起使用,则仅当事件间隔超过100 毫秒时才会触发窗口,这不是您所说的想要的。

    我建议你看看ProcessFunction。用这种方式得到你想要的东西应该很简单。

    【讨论】:

    • 我想每 100 毫秒获取 1 小时以来的所有数据,我需要一个 deltatrigger 否?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-27
    • 2020-10-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-07-23
    相关资源
    最近更新 更多