上一篇讲了 Table 查询与过滤——select 投影和 filter 过滤,这两个是最基础的一对一转换操作。这篇讲分组与聚合——groupBy 和聚合函数,这是 Flink SQL 中最常用的多对一转换操作,也是统计分析、数据汇总、指标计算的核心。
几乎每一个数据分析作业都会用到分组聚合——按用户统计订单数、按地区统计销售额、按时间统计 PV/UV。但很多同学只是简单地写groupBy(...).select(sum(...)),对流处理聚合的状态管理、回撤机制、两阶段聚合、数据倾斜处理知之甚少。实际上,分组聚合是流处理中最容易出性能问题和状态膨胀的算子之一。
这篇从本质定义、四步执行原理、内置聚合函数分类、状态管理与回撤机制、两阶段聚合、微批预聚合、分组集、自定义聚合函数、完整代码实现、常见配置、六个常见坑、八条最佳实践,把分组与聚合一次性讲透。
一、分组与聚合本质定义
分组与聚合的本质一句话概括:groupBy(分组)按指定的 key 将数据分成多个组,聚合函数(Aggregate Function)对每个组内的多行数据进行计算,输出一行聚合结果,是关系代数中最核心的多对一转换操作(γ操作)。
下面这张图把分组与聚合的本质定义、四步执行原理、内置聚合函数六大分类放在一起展示。
从关系代数的角度看,groupBy + 聚合对应分组聚合操作(γ),是统计分析、数据汇总、指标计算的基础。几乎所有的数据分析作业都会用到分组聚合。
1.1 批处理 vs 流处理
批处理的分组聚合是一次性的——所有数据到达后,按 key 分组,对每组应用聚合函数,输出最终结果。批处理聚合不需要维护状态(处理完就释放),也不会产生回撤消息。
流处理的分组聚合是持续的——数据持续到达,每来一条数据就更新对应 key 的聚合状态,输出更新后的结果(可能产生回撤消息)。流处理聚合需要维护状态,这是与批处理最大的区别,也是流处理聚合容易出问题的根源。
1.2 有状态算子
分组聚合是有状态算子,需要维护每个 key 的聚合累加器(Accumulator)状态。状态大小 = key 基数 × 每个 key 的累加器大小。如果 key 基数大(如用户 ID、订单 ID),状态会持续增长,必须设置状态 TTL(table.exec.state.ttl)防止 OOM。
这一点非常重要——很多同学的流处理聚合作业运行一段时间后 OOM,根因就是没有设置状态 TTL,状态无限增长。
二、groupBy 分组聚合四步原理
从输入数据到聚合结果,完整的流程分为四步:
- 数据按 key 分区:按 groupBy 的 key 对数据进行 hash 分区,相同 key 的数据路由到同一个并行实例
- 状态查找与更新:每个并行实例维护 key → 聚合累加器的映射,每来一条数据查找对应 key 的累加器并更新(调用 accumulate())
- 聚合结果计算:从更新后的累加器中提取聚合值(sum/count/avg 等),生成输出行(调用 getValue())
- 结果输出(可能回撤):流处理输出更新后的结果,可能产生回撤消息(Retract)或覆盖消息(Upsert)
关键理解:第 2 步的状态更新是分组聚合的核心。每个 key 维护一个累加器,累加器是一个可变对象,存储聚合的中间状态。比如 sum 的累加器就是一个数值(当前累加和),count 的累加器是一个计数器,avg 的累加器是(sum, count)两个值。
自定义聚合函数的核心就是实现累加器的创建(createAccumulator)、更新(accumulate)、结果提取(getValue)和合并(merge)四个方法。
三、内置聚合函数六大分类
Flink SQL 提供了丰富的内置聚合函数,可以分为六大类:
3.1 常规聚合
最常用的聚合函数,所有统计分析的基础:
SUM:求和COUNT:计数AVG:平均值MIN:最小值MAX:最大值
3.2 去重聚合
对去重后的值进行聚合,UV 统计等场景常用:
COUNT(DISTINCT):去重计数SUM(DISTINCT):去重求和AVG(DISTINCT):去重平均
注意:去重聚合需要维护去重状态(存储所有不重复的值),状态开销大。流处理中推荐用APPROX_COUNT_DISTINCT(近似去重,基于 HyperLogLog,状态极小)替代精确COUNT(DISTINCT)。
3.3 统计聚合
统计分析相关的聚合函数:
STDDEV/STDDEV_POP/STDDEV_SAMP:标准差(总体/样本)VAR_POP/VAR_SAMP:方差(总体/样本)
3.4 字符串聚合
将多行字符串合并为一行:
STRING_AGG:字符串拼接(可指定分隔符)LISTAGG:列表聚合
注意:流处理字符串聚合可能产生大状态(存储所有字符串),需设置 TTL。
3.5 布尔聚合
布尔值相关的聚合:
EVERY:全部为真SOME:至少一个为真BOOL_AND/BOOL_OR:布尔与/或
3.6 自定义聚合
内置函数不支持的场景,用自定义聚合函数:
AggregateFunction:自定义聚合函数(多行转一行)TableAggregateFunction:自定义表聚合函数(多行转多行,如 TopN)
适用场景:中位数、百分位、复杂统计指标、TopN 等。
四、状态管理与回撤机制
流处理聚合与批处理最大的区别就是状态管理和回撤机制。理解这两个概念是掌握流处理聚合的关键。
下面这张图把状态管理与回撤机制、两阶段聚合、分组集放在一起展示。
4.1 聚合状态管理
每个并行实例维护key → 累加器(Accumulator)的映射。每来一条数据,按 key 查找累加器,调用accumulate()更新。状态大小 = key 基数 × 每个 key 的累加器大小。
状态清理:设置table.exec.state.ttl,过期 key 的状态自动清理。这是流处理聚合的必配项,key 基数大时不设置 TTL 必然 OOM。
状态后端:小状态用HashMapStateBackend(内存,速度快),大状态用EmbeddedRocksDBStateBackend(磁盘,支持 TB 级状态)。聚合 key 基数大时推荐用 RocksDB。
4.2 Retract 回撤模式
每次聚合结果更新时,输出两条消息——-(旧结果)(回撤旧值)和+(新结果)(发送新值)。下游算子需要处理回撤消息,如 sink 需要支持撤回(retract sink)。
缺点:输出消息量翻倍(每条更新两条消息),网络和下游处理开销大。
触发条件:聚合函数不支持merge(),或 groupBy 的 key 不能作为唯一主键。
4.3 Upsert 覆盖模式
每次聚合结果更新时,只输出一条带主键的更新消息,下游按主键覆盖旧结果。主键来源是 groupBy 的 key(自然成为聚合结果的唯一主键)。
优点:输出消息量减半(每条更新一条消息),效率更高。
触发条件:聚合结果有唯一主键(groupBy key),且聚合函数支持merge()。Flink 优先选择 upsert 模式。
最佳实践:优先使用 upsert 模式(groupBy key 作为主键),选择支持 upsert 的 sink(JDBC upsert、HBase、Redis),避免 retract 模式的双倍消息量。
4.4 状态一致性
聚合状态参与 Checkpoint,故障恢复后状态一致,回撤消息不丢失。迟到数据到达后,更新状态,输出新的回撤/更新消息。
注意:状态 TTL 过期后,迟到数据会被当作新 key 处理,可能导致结果不准确。需要根据业务的迟到数据容忍度设置合理的 TTL。
五、两阶段聚合(Local-Global Aggregation)
两阶段聚合是解决数据倾斜和提升聚合性能的核心优化,分为本地预聚合和全局聚合两个阶段。
5.1 执行流程
- 阶段 1:Local Aggregation(本地预聚合)——每个并行实例先对本地数据按 key 做部分聚合,减少需要 shuffle 的数据量
- Shuffle:按 key 重新分区——本地预聚合结果按 key hash 分区,相同 key 路由到同一个全局聚合实例
- 阶段 2:Global Aggregation(全局聚合)——对所有本地预聚合结果进行 merge,输出最终聚合结果
5.2 解决数据倾斜的原理
如果某个 key 的数据量特别大(热点 key),单阶段聚合时该 key 的所有数据都路由到同一个并行实例,导致该实例处理压力过大(数据倾斜)。
两阶段聚合时,热点 key 的数据先在多个本地实例预聚合,每个本地实例只输出一条部分聚合结果,全局聚合实例只需要 merge 少量部分结果,大幅减轻热点 key 的压力。
5.3 配置
table.optimizer.agg-phase-strategy: TWO_PHASE(开启两阶段聚合)。可选值:
AUTO:自动选择TWO_PHASE:强制两阶段ONE_PHASE:单阶段(关闭)
生产环境推荐设置为TWO_PHASE或AUTO。
5.4 限制
两阶段聚合要求聚合函数实现merge()方法(将两个累加器合并)。内置聚合函数(sum/count/avg/min/max)都支持 merge,自定义聚合函数需要手动实现 merge()。不实现 merge 则无法使用两阶段聚合,数据倾斜时性能差。
六、微批预聚合(MiniBatch)
微批预聚合是时机上的优化——攒一批数据再处理,而不是来一条处理一条。
6.1 原理
开启 MiniBatch 后,聚合算子不会每条数据都更新状态和输出结果,而是攒一批数据(达到时间阈值或条数阈值),对这批数据先做本地预聚合,再批量更新状态和输出结果。
6.2 效果
- 减少状态访问开销(批量更新状态,减少状态读写次数)
- 减少输出消息量(批量输出,减少回撤消息数量)
- 提升吞吐(尤其是高吞吐、低延迟要求不高的场景)
6.3 配置
table.exec.mini-batch.enabled: true:开启微批预聚合table.exec.mini-batch.allow-latency: 5s:微批最大延迟(攒够时间就触发)table.exec.mini-batch.size: 5000:微批最大条数(攒够条数就触发)
allow-latency 和 size 先到先触发。延迟敏感场景调小 allow-latency,数据量大可调大 size。
6.4 与两阶段聚合的关系
两阶段聚合是结构上的优化(Local + Global),微批预聚合是时机上的优化(攒一批再处理)。两者可以同时使用,效果叠加——微批攒一批数据,本地预聚合减少 shuffle 量,全局聚合输出最终结果。生产环境推荐同时开启。
七、分组集(GROUPING SETS / ROLLUP / CUBE)
分组集是多维聚合的强大工具,一次查询输出多个维度的聚合结果,避免多次查询。
7.1 GROUPING SETS
显式指定多个分组维度组合,一次查询输出所有组合的聚合结果。结果数等于指定的维度组合数。
适用场景:需要自定义维度组合的多维分析。用GROUPING()函数标识当前行属于哪个维度组合。
7.2 ROLLUP
生成从详细到汇总的层级维度组合,适合有层级关系的维度。n 个维度生成 n+1 个组合。
例如ROLLUP(年,季,月)生成 (年,季,月)、(年,季)、(年)、() 四个组合。
适用场景:时间维度、地理维度等有层级关系的多维分析。注意维度顺序很重要,ROLLUP(a,b) ≠ ROLLUP(b,a)。
7.3 CUBE
生成所有可能的维度组合(笛卡尔积),最全面的多维聚合。n 个维度生成 2^n 个组合。
例如CUBE(地区,产品)生成 (地区,产品)、(地区)、(产品)、() 四个组合。
适用场景:需要全面交叉分析的多维报表。注意:维度多时组合数爆炸(10 维 = 1024 组合),慎用。
7.4 流处理分组集注意事项
流处理分组集需要维护所有维度组合的聚合状态,状态大小 = 所有维度组合的 key 基数之和。维度多时状态膨胀严重,必须设置状态 TTL。
CUBE 的组合数是 2^n,流处理中慎用(3 维以上就可能有状态问题)。推荐用 GROUPING SETS 显式指定需要的组合,避免不必要的维度组合。
分组集中的全聚合(() 空维度)会将所有数据路由到同一个并行实例,造成严重的数据倾斜。解决方案:开启两阶段聚合 + 微批预聚合,或单独执行全聚合。
八、自定义聚合函数(AggregateFunction)
内置聚合函数不支持的场景(中位数、百分位、复杂统计指标),需要自定义聚合函数。自定义聚合函数的核心是实现四个方法:
8.1 createAccumulator()
创建累加器(Accumulator),存储聚合的中间状态。每个新 key 第一次到达时调用一次。
累加器应该是可变对象(如 List、Map、自定义类),避免每次更新都创建新对象。累加器的大小决定状态大小,尽量精简。
8.2 accumulate(acc, value…)
每来一条数据,更新累加器。每条数据到达时调用,是调用最频繁的方法,性能至关重要。
方法名必须是 accumulate,第一个参数是累加器,后面的参数是输入值。可以重载多个 accumulate 方法,支持不同的输入类型。更新操作应该是原地修改累加器,避免创建新对象。
避免在 accumulate 中做昂贵操作(如排序、网络请求)。
8.3 getValue(acc)
从累加器中提取最终聚合结果。每次需要输出聚合结果时调用。
返回值类型是聚合函数的输出类型。getValue 应该是纯计算,不修改累加器。如果计算成本高(如中位数需要排序),可以考虑在 accumulate 中维护排序好的数据结构,getValue 直接返回。
8.4 merge(acc, accs)
合并多个累加器,用于两阶段聚合和会话窗口。第一个参数是目标累加器,后面是待合并的累加器迭代器。
不实现 merge 则无法使用两阶段聚合,数据倾斜时性能差。生产环境自定义聚合必须实现 merge。
可选方法:retract(acc, value)(回撤,用于 Retract 模式)、resetAccumulator(acc)(重置累加器)。
九、完整代码实现
下面是一个完整的分组聚合使用示例,包含创建环境、DDL 建表、groupBy 分组聚合、两阶段聚合配置、自定义聚合函数、输出执行。
下面这张图把开发流程、自定义聚合函数开发要点、常见配置、六个常见坑、八条最佳实践放在一起展示。
importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.table.api.Table;importorg.apache.flink.table.api.TableResult;importorg.apache.flink.table.api.bridge.java.StreamTableEnvironment;importorg.apache.flink.table.functions.AggregateFunction;importjava.time.ZoneId;importstaticorg.apache.flink.table.api.Expressions.*;publicclassTableAggregateExample{// 自定义聚合函数:计算中位数publicstaticclassMedianFunctionextendsAggregateFunction<Double,MedianAccumulator>{@OverridepublicMedianAccumulatorcreateAccumulator(){returnnewMedianAccumulator();}@Overridepublicvoidaccumulate(MedianAccumulatoracc,Doublevalue){if(value!=null){acc.values.add(value);}}@OverridepublicDoublegetValue(MedianAccumulatoracc){if(acc.values.isEmpty())returnnull;java.util.Collections.sort(acc.values);intsize=acc.values.size();if(size%2==0){return(acc.values.get(size/2-1)+acc.values.get(size/2))/2.0;}else{returnacc.values.get(size/2);}}@Overridepublicvoidmerge(MedianAccumulatoracc,Iterable<MedianAccumulator>it){for(MedianAccumulatorother:it){acc.values.addAll(other.values);}}}// 中位数累加器publicstaticclassMedianAccumulator{publicjava.util.List<Double>values=newjava.util.ArrayList<>();}publicstaticvoidmain(String[]args)throwsException{// 1. 创建执行环境StreamExecutionEnvironmentenv=StreamExecutionEnvironment.getExecutionEnvironment();env.enableCheckpointing(60000);StreamTableEnvironmenttableEnv=StreamTableEnvironment.create(env);tableEnv.getConfig().setLocalTimeZone(ZoneId.of("Asia/Shanghai"));// 2. 开启两阶段聚合和微批预聚合tableEnv.getConfig().set("table.optimizer.agg-phase-strategy","TWO_PHASE");tableEnv.getConfig().set("table.exec.mini-batch.enabled","true");tableEnv.getConfig().set("table.exec.mini-batch.allow-latency","5s");tableEnv.getConfig().set("table.exec.mini-batch.size","5000");// 3. 设置状态TTL(防止聚合状态无限增长)tableEnv.getConfig().set("table.exec.state.ttl","1h");// 4. 注册自定义聚合函数tableEnv.createTemporarySystemFunction("median",MedianFunction.class);// 5. DDL创建Kafka源表tableEnv.executeSql(""" CREATE TABLE orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, city STRING, order_time TIMESTAMP(3), WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'orders', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json' ) """);// 6. DDL创建MySQL结果表(主键支持Upsert模式)tableEnv.executeSql(""" CREATE TABLE user_stat ( user_id BIGINT, total_amount DECIMAL(10,2), order_count BIGINT, avg_amount DECIMAL(10,2), median_amount DOUBLE, PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://localhost:3306/flink_db', 'table-name' = 'user_stat', 'username' = 'root', 'password' = '123456' ) """);// 7. groupBy分组聚合:按用户统计多维度指标TableuserStat=tableEnv.from("orders").filter($("status").isEqual(lit("PAID"))).groupBy($("user_id")).select($("user_id"),$("amount").sum().as("total_amount"),$("order_id").count().as("order_count"),$("amount").avg().as("avg_amount"),call("median",$("amount")).as("median_amount"));// 8. 查看执行计划(确认两阶段聚合和Upsert模式)Stringplan=userStat.explain();System.out.println("=== 执行计划 ===");System.out.println(plan);// 9. 写入结果表并执行TableResultresult=userStat.executeInsert("user_stat");// 10. 等待作业完成result.await();}}代码关键点:
- 两阶段聚合配置:
table.optimizer.agg-phase-strategy=TWO_PHASE,解决数据倾斜。 - 微批预聚合配置:
table.exec.mini-batch.enabled=true+allow-latency=5s+size=5000,减少状态访问开销。 - 状态 TTL:
table.exec.state.ttl=1h,防止聚合状态无限增长导致 OOM。 - 自定义聚合函数:MedianFunction 计算中位数,实现 createAccumulator/accumulate/getValue/merge 四个方法,merge 是两阶段聚合的必要条件。
- groupBy 分组聚合:按用户分组,计算总金额、订单数、平均金额、中位数四个指标,聚合后用
as()命名结果列。 - Upsert 模式:结果表定义主键(user_id),聚合结果按主键覆盖,优先使用 Upsert 模式,效率更高。
- explain 检查:执行
userStat.explain()查看执行计划,确认两阶段聚合和输出模式。
十、常见配置
生产环境聚合相关的八项关键配置:
| 配置项 | 推荐值 | 说明与影响 |
|---|---|---|
| table.optimizer.agg-phase-strategy | TWO_PHASE | 两阶段聚合策略,解决数据倾斜,需聚合函数实现 merge() |
| table.exec.mini-batch.enabled | true | 开启微批预聚合,攒一批再处理,减少状态访问开销 |
| table.exec.mini-batch.allow-latency | 5s | 微批最大延迟,与 size 配合平衡延迟和吞吐 |
| table.exec.mini-batch.size | 5000 | 微批最大条数,与 allow-latency 配合,先到先触发 |
| table.exec.state.ttl | 1h~24h | 状态过期时间,防止聚合状态无限增长导致 OOM |
| state.backend | rocksdb | 状态后端,大状态用 RocksDB(磁盘),小状态用 HashMap(内存) |
| table.exec.resource.default-parallelism | 按集群 | Table 作业默认并行度,聚合并行度影响数据分布和状态大小 |
| pipeline.name | 作业名 | 作业名称,显示在 Flink Web UI,生产环境必须设置 |
十一、六个常见坑
11.1 坑一:未设置状态 TTL 导致状态无限增长
现象:作业运行一段时间后 Checkpoint 越来越大,TaskManager OOM,作业失败。
根因:流处理聚合维护每个 key 的累加器状态,key 基数大时状态持续增长,不会自动清理。
解决方案:设置table.exec.state.ttl(1h~24h),过期 key 的状态自动清理。大状态用 RocksDB 状态后端。
11.2 坑二:数据倾斜导致单个并行实例过载
现象:某个并行实例 CPU/内存使用率远高于其他实例,作业整体吞吐上不去,延迟高。
根因:某个 key 的数据量特别大(热点 key),单阶段聚合时该 key 的所有数据都路由到同一个实例。
解决方案:① 开启两阶段聚合(TWO_PHASE);② 开启微批预聚合(MiniBatch);③ 热点 key 加盐(随机前缀打散);④ 单独处理热点 key。
11.3 坑三:自定义聚合函数未实现 merge() 无法两阶段聚合
现象:自定义聚合函数的数据倾斜严重,开启两阶段聚合后没有效果。
根因:自定义聚合函数没有实现 merge() 方法,优化器无法将其用于两阶段聚合的全局 merge 阶段。
解决方案:实现 merge(acc, accs) 方法,将多个累加器合并为一个。merge 是两阶段聚合的必要条件,生产环境自定义聚合必须实现。
11.4 坑四:COUNT(DISTINCT) 状态膨胀
现象:使用 COUNT(DISTINCT) 的作业状态特别大,Checkpoint 慢,OOM。
根因:去重聚合需要维护每个 key 的去重集合(存储所有不重复的值),状态大小 = key 基数 × 每个 key 的去重值数量。
解决方案:① 用 APPROX_COUNT_DISTINCT(近似去重,基于 HyperLogLog,状态极小)替代精确 COUNT(DISTINCT);② 设置状态 TTL;③ 大状态用 RocksDB。
11.5 坑五:流处理聚合输出回撤消息下游不支持
现象:聚合结果写入 sink 时报错,或结果不正确(旧值没有被撤回)。
根因:流处理聚合可能输出回撤消息(Retract 模式),下游 sink 不支持回撤处理,导致旧值没有被撤回。
解决方案:① 优先使用 upsert 模式(groupBy key 作为主键,sink 支持 upsert);② 选择支持 retract 的 sink;③ 用 explain() 确认输出模式。
11.6 坑六:分组集 CUBE 维度过多状态爆炸
现象:使用 CUBE 的作业状态特别大,OOM,作业失败。
根因:CUBE 生成所有维度组合(2^n 个),每个组合都维护独立的聚合状态,维度多时状态爆炸。
解决方案:① 用 GROUPING SETS 显式指定需要的组合,避免不必要的组合;② 用 ROLLUP 替代 CUBE(层级组合,n+1 个);③ 流处理中 CUBE 维度不超过 3 个;④ 设置状态 TTL。
十二、八条最佳实践 Checklist
上线前逐条检查:
必须设置状态 TTL:
table.exec.state.ttl设置 1h~24h,防止聚合状态无限增长导致 OOM,key 基数大时必须设置。开启两阶段聚合:
table.optimizer.agg-phase-strategy=TWO_PHASE,解决数据倾斜,自定义聚合函数必须实现 merge()。开启微批预聚合:
table.exec.mini-batch.enabled=true+allow-latency=5s+size=5000,减少状态访问开销,提升吞吐。大状态用 RocksDB:聚合 key 基数大、状态大时用 RocksDB 状态后端(磁盘存储),小状态用 HashMap(内存,速度快)。
优先用 upsert 输出模式:groupBy key 作为主键,选择支持 upsert 的 sink(JDBC/HBase/Redis),避免 retract 模式的双倍消息量。
精确去重改用近似去重:UV 统计等场景用 APPROX_COUNT_DISTINCT 替代 COUNT(DISTINCT),基于 HyperLogLog,状态极小,误差可控。
分组集用 GROUPING SETS:多维聚合用 GROUPING SETS 显式指定需要的组合,避免 CUBE 的 2^n 组合爆炸,流处理 CUBE 维度不超过 3 个。
上线前 explain 检查:执行 explain() 检查执行计划,确认两阶段聚合、输出模式(upsert/retract)、算子顺序、并行度是否符合预期。
十三、总结与下一篇预告
Table 分组与聚合原理及代码实现要点回顾:
第一,本质:groupBy(分组)按 key 将数据分成多个组,聚合函数对每个组内的多行数据进行计算,输出一行聚合结果,是关系代数中最核心的多对一转换操作(γ操作)。批处理一次性聚合,流处理持续聚合(维护状态,可能回撤)。
第二,四步执行原理:数据按 key 分区 → 状态查找与更新(accumulate)→ 聚合结果计算(getValue)→ 结果输出(可能回撤 Retract 或覆盖 Upsert)。分组聚合是有状态算子,状态大小 = key 基数 × 累加器大小。
第三,内置聚合函数六大分类:常规聚合(SUM/COUNT/AVG/MIN/MAX)、去重聚合(COUNT(DISTINCT),推荐用 APPROX_COUNT_DISTINCT)、统计聚合(STDDEV/VAR)、字符串聚合(STRING_AGG)、布尔聚合(EVERY/SOME)、自定义聚合(AggregateFunction/TableAggregateFunction)。
第四,状态管理与回撤机制:聚合状态是 key → 累加器映射,设置状态 TTL 自动清理,大状态用 RocksDB。两种输出模式:Retract 回撤(输出 -(旧)+(新) 两条消息,开销大)和 Upsert 覆盖(输出带主键更新消息,效率高)。优先使用 Upsert 模式。
第五,两阶段聚合:Local 本地预聚合 → Shuffle → Global 全局聚合,解决数据倾斜。要求聚合函数实现 merge()。配置table.optimizer.agg-phase-strategy=TWO_PHASE。
第六,微批预聚合:攒一批数据再处理,减少状态访问开销和输出消息量。配置table.exec.mini-batch.enabled=true+allow-latency=5s+size=5000。与两阶段聚合同时使用效果叠加。
第七,分组集:GROUPING SETS(自定义组合)、ROLLUP(层级组合 n+1)、CUBE(全组合 2^n,慎用)。流处理 CUBE 维度不超过 3 个,推荐用 GROUPING SETS。
第八,自定义聚合函数:四个核心方法——createAccumulator(创建累加器)、accumulate(更新累加器,调用最频繁)、getValue(提取结果)、merge(合并累加器,两阶段聚合必要条件)。
第九,六个常见坑:未设置状态 TTL 导致 OOM、数据倾斜单个实例过载、自定义聚合未实现 merge()、COUNT(DISTINCT) 状态膨胀、回撤消息下游不支持、CUBE 维度过多状态爆炸。
第十,八条最佳实践:必须设置状态 TTL、开启两阶段聚合、开启微批预聚合、大状态用 RocksDB、优先用 upsert 输出模式、精确去重改用近似去重、分组集用 GROUPING SETS、上线前 explain 检查。
分组与聚合是 Flink SQL 中最常用也最容易出性能问题的算子。理解了状态管理、回撤机制、两阶段聚合、微批预聚合这些核心概念,就能写出高性能、稳定的流处理聚合作业,避免状态膨胀和数据倾斜这些常见的生产事故。