【问题标题】:Flink Streaming: How to implement windows which are defined by a start and end element?Flink Streaming:如何实现由 start 和 end 元素定义的窗口?
【发布时间】:2016-11-09 06:35:30
【问题描述】:

我有以下格式的数据,

SIP|2405463430|4115474257|8.205142580136622E12|美国标准时间 11 月 8 日星期二 16:58:58 2016|邀请 RTP|2405463430|4115474257|8.205142580136622E12|11 月 8 日星期二 16:58:58 IST 2016|0 RTP|2405463430|4115474257|8.205142580136622E12|星期二 2016 年 11 月 8 日 16:58:58 IST|1 RTP|2405463430|4115474257|8.205142580136622E12|11 月 8 日星期二 16:58:58 IST 2016|2 RTP|2405463430|4115474257|8.205142580136622E12|11 月 8 日星期二 16:58:58 IST 2016|3 RTP|2405463430|4115474257|8.205142580136622E12|星期二 2016 年 11 月 8 日 16:58:58 IST|4 RTP|2405463430|4115474257|8.205142580136622E12|11 月 8 日星期二 16:58:58 IST 2016|5 RTP|2405463430|4115474257|8.205142580136622E12|11 月 8 日星期二 16:58:58 IST 2016|6 RTP|2405463430|4115474257|8.205142580136622E12|星期二 2016 年 11 月 8 日 16:58:58 IST|7 RTP|2405463430|4115474257|8.205142580136622E12|11 月 8 日星期二 16:58:58 IST 2016|8 RTP|2405463430|4115474257|8.205142580136622E12|11 月 8 日星期二 16:58:58 IST 2016|9 SIP|2405463430|4115474257|8.205142580136622E12|星期二 2016 年 11 月 8 日 16:58:58 IST|再见

我希望我的窗口在遇到SIP-INVITE 消息时启动,并在遇到SIP-BYE 消息时触发一个事件,执行一些聚合。

我该怎么做? SIP-INVITE 消息在给定用户的任何时间点都会出现,而且我可能还会同时收到多个用户的多个 SIP-INVITE 消息。

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    我认为您可以使用由用户键入的全局窗口来解决您的用例。全局窗口收集每个键的所有数据,并将触发和清除窗口的责任推给用户定义的Trigger 函数。

    全局窗口定义如下:

    val input: DataStream[(String, Int, String)] = ??? // (userId, value, marker)
    val agg = input
      // one global window per user (handles overlapping SIP-INVITE events).
      .keyBy(_._1)
      // collect all data for each user until the trigger fires and purges the window.
      .window(GlobalWindows.create())
      // you have to implement a custom trigger which reacts on the marker.
      .trigger(new YourCustomTrigger())
      // the window function computes your aggregation.
      .apply(new YourWindowFunction())
    

    我认为执行以下操作的触发器应该可以工作(假设SIP-INVITE 事件总是在启动会话)。 Trigger.onElement() 方法应检查 SIP-BYE 字段并触发窗口评估并清除窗口,即返回 TriggerResult.FIRE_AND_PURGE。这将调用评估函数并移除窗口状态。

    注意,如果要支持乱序事件,需要特别注意(在这种情况下,您应该为关闭元素的时间戳设置一个事件时间计时器,以确保接收到时间戳之前的所有数据) .如果有数据应该被丢弃,因为它不是“介于”SIP-INVITE 和SIP-BYE 之间,你也需要处理它。

    有关详细信息,请参阅global windows 和 triggers 的文档、[Trigger][3] 的 JavaDocs 和此 blog post。

    【讨论】:

    • 谢谢 :) 这很有帮助!
    猜你喜欢
    • 2018-06-02
    • 1970-01-01
    • 2020-08-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多