【问题标题】:Flink event correlation and lookbackFlink 事件关联和回溯
【发布时间】:2021-03-17 14:21:58
【问题描述】:

我是 flink 新手,正在寻找有关构建实时事件关联系统的建议。我有两个主要用例:

  1. 事件关联逻辑由基于输入流中事件类型的静态规则组成。在最后 X 分钟内,根据这些规则关联不同事件类型的事件和具有业务价值的事件的输出数据。例如,在最后 1 分钟内,检查市场 A1 中事件类型 A 的价格是否
  2. 对于感兴趣/商业价值的事件,计算与最后 X 分钟的价格差异。例如,如果在应用所有规则后,我们确定事件 A 在最后 1 分钟窗口内是感兴趣的,在将事件数据添加到输出流中之前,我们还想计算事件 A 与过去 10 分钟的价格差异。

为了实现这些用例,我正在评估通过输入数据中的产品类型 ID 对输入流应用密钥。这将为我提供该产品针对不同市场的多种事件类型的数据,然后使用回溯期的滑动事件时间窗口,例如最后 10 分钟,滑动窗口为 1 分钟,并应用 ProcessWindowFunction 为最后 1 分钟的数据编写相关逻辑和使用其他 9 分钟的数据进行回顾并计算感兴趣事件的价格差异。

我不完全确定这是否是实现这些的最佳方式。任何提示/建议将不胜感激!

【问题讨论】:

    标签: apache-flink flink-streaming flink-statefun


    【解决方案1】:

    总的来说,我会说您的选择是:

    • 按照您的建议使用滑动窗口。
    • 使用KeyedProcessFunction。这个较低级别的 API 提供了更多的控制,并可能导致更好的优化解决方案。有时这也更简单,所以如果您发现窗口 API 妨碍了您,请考虑这一点。
    • 使用 Flink SQL 和/或 Table API。如果规则是用 SQL 编写的,您可能会发现更容易表达和维护规则。也许MATCH_RECOGNIZE 是相关的。

    【讨论】:

    • 感谢大卫的帮助!有一个后续问题。虽然这种滑动窗口方法效果很好,但我正在考虑将回溯逻辑完全分离到一个单独的 KDA 应用程序中以增加模块化。因此输入流(流 I)进入两个 Flink 应用程序(Flink A:计算 1 分钟窗口上的相关性并在流 A 中输出感兴趣的事件,Flink B:将输入作为流 I,流 A 使用 connect 和 KeyedCoProcessFunction 计算价格差异和在流 B 中发出输出。想了解这种方法是否比在 1 个 Flink/KDA 应用程序中执行所有操作有缺点。
    • 我能想到的最大缺点是增加了延迟(可能不是那么重要)。
    猜你喜欢
    • 1970-01-01
    • 2012-07-10
    • 2020-05-27
    • 1970-01-01
    • 2021-07-25
    • 2019-03-19
    • 2019-01-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多