【发布时间】:2022-09-30 23:11:40
【问题描述】:
我正在通过 Flink 作业从 websocket 传输数据,需要根据以下逻辑输出滚动加权平均值:
每条消息都有属性\"parent\"、\"name\"、\"amount\"、\"value\" 通过\"name\"获取最新消息,并与每个\"parent\"的其他最新消息相结合,得到基于\"amount\"和\"value\"的加权平均值
- 父 = \"a\";名称 = \"m\";金额=100;值=12.45
- 父 = \"a\";名称 = \"n\";金额=40;值=14.55
- 父 = \"a\";名称 = \"m\";金额=100;值=17.45
- 父 = \"a\";名称 = \"o\";金额=24;值=13.25
- 父 = \"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