拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink SQL 之 Table 分组与聚合原理及代码实现:从状态管理到两阶段聚合与数据倾斜处理

Flink SQL 之 Table 分组与聚合原理及代码实现:从状态管理到两阶段聚合与数据倾斜处理 上一篇讲了 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状态会持续增长必须设置状态 TTLtable.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 的 sinkJDBC upsert、HBase、Redis避免 retract 模式的双倍消息量。4.4 状态一致性聚合状态参与 Checkpoint故障恢复后状态一致回撤消息不丢失。迟到数据到达后更新状态输出新的回撤/更新消息。注意状态 TTL 过期后迟到数据会被当作新 key 处理可能导致结果不准确。需要根据业务的迟到数据容忍度设置合理的 TTL。五、两阶段聚合Local-Global Aggregation两阶段聚合是解决数据倾斜和提升聚合性能的核心优化分为本地预聚合和全局聚合两个阶段。5.1 执行流程阶段 1Local Aggregation本地预聚合——每个并行实例先对本地数据按 key 做部分聚合减少需要 shuffle 的数据量Shuffle按 key 重新分区——本地预聚合结果按 key hash 分区相同 key 路由到同一个全局聚合实例阶段 2Global 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 个维度生成 n1 个组合。例如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{// 自定义聚合函数计算中位数publicstaticclassMedianFunctionextendsAggregateFunctionDouble,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);intsizeacc.values.size();if(size%20){return(acc.values.get(size/2-1)acc.values.get(size/2))/2.0;}else{returnacc.values.get(size/2);}}Overridepublicvoidmerge(MedianAccumulatoracc,IterableMedianAccumulatorit){for(MedianAccumulatorother:it){acc.values.addAll(other.values);}}}// 中位数累加器publicstaticclassMedianAccumulator{publicjava.util.ListDoublevaluesnewjava.util.ArrayList();}publicstaticvoidmain(String[]args)throwsException{// 1. 创建执行环境StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();env.enableCheckpointing(60000);StreamTableEnvironmenttableEnvStreamTableEnvironment.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分组聚合按用户统计多维度指标TableuserStattableEnv.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模式StringplanuserStat.explain();System.out.println( 执行计划 );System.out.println(plan);// 9. 写入结果表并执行TableResultresultuserStat.executeInsert(user_stat);// 10. 等待作业完成result.await();}}代码关键点两阶段聚合配置table.optimizer.agg-phase-strategyTWO_PHASE解决数据倾斜。微批预聚合配置table.exec.mini-batch.enabledtrueallow-latency5ssize5000减少状态访问开销。状态 TTLtable.exec.state.ttl1h防止聚合状态无限增长导致 OOM。自定义聚合函数MedianFunction 计算中位数实现 createAccumulator/accumulate/getValue/merge 四个方法merge 是两阶段聚合的必要条件。groupBy 分组聚合按用户分组计算总金额、订单数、平均金额、中位数四个指标聚合后用as()命名结果列。Upsert 模式结果表定义主键user_id聚合结果按主键覆盖优先使用 Upsert 模式效率更高。explain 检查执行userStat.explain()查看执行计划确认两阶段聚合和输出模式。十、常见配置生产环境聚合相关的八项关键配置配置项推荐值说明与影响table.optimizer.agg-phase-strategyTWO_PHASE两阶段聚合策略解决数据倾斜需聚合函数实现 merge()table.exec.mini-batch.enabledtrue开启微批预聚合攒一批再处理减少状态访问开销table.exec.mini-batch.allow-latency5s微批最大延迟与 size 配合平衡延迟和吞吐table.exec.mini-batch.size5000微批最大条数与 allow-latency 配合先到先触发table.exec.state.ttl1h~24h状态过期时间防止聚合状态无限增长导致 OOMstate.backendrocksdb状态后端大状态用 RocksDB磁盘小状态用 HashMap内存table.exec.resource.default-parallelism按集群Table 作业默认并行度聚合并行度影响数据分布和状态大小pipeline.name作业名作业名称显示在 Flink Web UI生产环境必须设置十一、六个常见坑11.1 坑一未设置状态 TTL 导致状态无限增长现象作业运行一段时间后 Checkpoint 越来越大TaskManager OOM作业失败。根因流处理聚合维护每个 key 的累加器状态key 基数大时状态持续增长不会自动清理。解决方案设置table.exec.state.ttl1h~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层级组合n1 个③ 流处理中 CUBE 维度不超过 3 个④ 设置状态 TTL。十二、八条最佳实践 Checklist上线前逐条检查必须设置状态 TTLtable.exec.state.ttl设置 1h~24h防止聚合状态无限增长导致 OOMkey 基数大时必须设置。开启两阶段聚合table.optimizer.agg-phase-strategyTWO_PHASE解决数据倾斜自定义聚合函数必须实现 merge()。开启微批预聚合table.exec.mini-batch.enabledtrueallow-latency5ssize5000减少状态访问开销提升吞吐。大状态用 RocksDB聚合 key 基数大、状态大时用 RocksDB 状态后端磁盘存储小状态用 HashMap内存速度快。优先用 upsert 输出模式groupBy key 作为主键选择支持 upsert 的 sinkJDBC/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-strategyTWO_PHASE。第六微批预聚合攒一批数据再处理减少状态访问开销和输出消息量。配置table.exec.mini-batch.enabledtrueallow-latency5ssize5000。与两阶段聚合同时使用效果叠加。第七分组集GROUPING SETS自定义组合、ROLLUP层级组合 n1、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 中最常用也最容易出性能问题的算子。理解了状态管理、回撤机制、两阶段聚合、微批预聚合这些核心概念就能写出高性能、稳定的流处理聚合作业避免状态膨胀和数据倾斜这些常见的生产事故。
返回列表