【问题标题】:Apache Flink weighted average based on two keys基于两个键的 Apache Flink 加权平均
【发布时间】:2022-09-30 23:11:40
【问题描述】:

我正在通过 Flink 作业从 websocket 传输数据,需要根据以下逻辑输出滚动加权平均值:

每条消息都有属性\"parent\"、\"name\"、\"amount\"、\"value\" 通过\"name\"获取最新消息,并与每个\"parent\"的其他最新消息相结合,得到基于\"amount\"和\"value\"的加权平均值

  1. 父 = \"a\";名称 = \"m\";金额=100;值=12.45
  2. 父 = \"a\";名称 = \"n\";金额=40;值=14.55
  3. 父 = \"a\";名称 = \"m\";金额=100;值=17.45
  4. 父 = \"a\";名称 = \"o\";金额=24;值=13.25
  5. 父 = \"a\";名称 = \"n\";金额=40;值=12.55

    消息 3、4 和 5 分别是 parent:name 的最新消息,因此这些消息应该用于获取“a”的当前加权平均值。 在任何时候,都不知道父母有多少孩子。 加权平均的逻辑很好。更多的是如何在 Flink 中进行 key、get latest、aggregation、average、keep state 等。

    我看过 RichFlatMapFunction、AggregateFunction 但证明很难将它们拼凑在一起。

    任何帮助或想法表示赞赏。

    标签: java apache-flink flink-streaming


    【解决方案1】:

    使用低级构建块,您可以使用KeyedProcessFunction 构建解决方案。您可以通过parent 键入事件流,然后使用MapState<String, Event> 跟踪每个名称的最新事件。随着事件的处理,您可以发出更新的结果。有关使用 MapState 的 KeyedProcessFunction 的示例,请参阅 the Flink docs

    如果要使用事件时间处理,则必须决定如何处理乱序事件。也许您可以忽略无序的事件,或者您可能需要先按时间戳对流进行排序。

    在更高级别工作,您可以使用 Flink SQL 代替。您可以使用按父级和名称组合划分的 OVER 窗口来跟踪每个父级/名称组合的最新事件,然后按父级分组并计算加权平均值(可能使用用户定义的聚合函数)。有关如何使用 OVER 窗口获取给定键的最新事件流的示例,请参阅the Immerok Cookbook

    免责声明:我为 Immerok 工作(我编写了 Flink 文档的那部分)。

    【讨论】:

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