【发布时间】:2021-03-16 17:13:47
【问题描述】:
我是 Flink 的新手。有时我想在 DataStream 上进行聚合而不需要先执行 keyBy。为什么 Flink 不支持 DataStream 上的聚合(sum、min、max 等)?
谢谢你, 艾哈迈德。
【问题讨论】:
-
如果您觉得这很有用,您可以发表评论、点赞或接受答案。否则这个问题可能对未来的观众没有用处。
标签: apache-flink flink-streaming
我是 Flink 的新手。有时我想在 DataStream 上进行聚合而不需要先执行 keyBy。为什么 Flink 不支持 DataStream 上的聚合(sum、min、max 等)?
谢谢你, 艾哈迈德。
【问题讨论】:
标签: apache-flink flink-streaming
Flink 支持非 keyed 流的聚合,但你必须先应用 windowAll 操作,然后才能应用聚合。 windowAll 函数会将并行度值降低到1,这意味着所有数据都将流经单个任务槽。这是设计使然,因为当您有多个任务槽时,您只能针对该槽中可用的数据流进行聚合,而不是跨槽。
如果您的用例不适合使用具有并行性的 windowsAll(即,当您有更多来自源的记录时),那么您可以尝试应用 keyBy 函数然后聚合,这将获得聚合结果一组键,然后是 windowAll,最后是聚合函数。这样,您可以在不同的任务槽中通过键进行聚合,然后最终在单个任务槽中聚合减少的数据。
以下是没有keyBy操作的windowAll示例,
environment.fromCollection(list)
.windowAll(TumblingEventTimeWindows.of(Time.seconds(5)))
.max(1)
下面是keyBy操作后windowAll的例子,
environment.fromCollection(list)
.keyBy(1)
.window(TumblingEventTimeWindows.of(Time.seconds(5)))
.maxBy(1)
.windowAll(TumblingEventTimeWindows.of(Time.seconds(5)))
.max(1)
文档参考 - here
【讨论】:
对于FLIP-134,Flink 社区决定弃用 DataStream API 中的所有这些关系方法:
DataStream#projectWindowed/KeyedStream#sum,min,max,minBy,maxByDataStream#keyBy 其中key用字段名或索引指定(包括ConnectedStreams#keyBy)此决定背后的基本原理是 Table/SQL 是一个更完整、性能更高的关系 API,并且它已经支持批处理和流式处理。使用这些 API,您可以轻松执行全局聚合,而无需先执行 keyBy 或 GROUP BY。
一个例子:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
SingleOutputStreamOperator<Integer> numbers = env.fromElements(0, 1, 1, 0, 3, 2);
Table data = tableEnv.fromDataStream(numbers, $("n"));
Table results = data.select($("n").max());
tableEnv
.toRetractStream(results, Row.class)
.print();
env.execute();
【讨论】: