应避免group by中多个count(distinct)或grouping sets触发expand导致数据膨胀,改用approx_count_distinct()、拆分查询或去除grouping sets;合理设置shuffle分区数(按100–200mb/分区估算);经营类报表优先采用预聚合表,并确保刷新逻辑与口径严格一致。

GROUP BY触发EXPAND操作导致数据膨胀10倍以上
当SQL中出现多个COUNT(DISTINCT)或GROUPING SETS,Spark会悄悄插入EXPAND物理算子——它不是优化,而是把单行输入“复制”成多行输出,只为满足不同聚合路径的计算需求。你看到的执行计划里Input Rows是730亿,但EXPAND之后变成3000亿,网络和磁盘压力直接翻倍。
避免方式很直接:
- 用
approx_count_distinct()替代COUNT(DISTINCT),误差可控且不触发EXPAND - 把一个含多个
COUNT(DISTINCT)的GROUP BY拆成多个独立查询,再用JOIN合并结果 - 确认是否真需要
GROUPING SETS;若只是分组统计,去掉后整个执行计划会退化为标准HashAggregate
shuffle分区数设为200会让小数据变慢、大数据OOM
spark.sql.shuffle.partitions默认200,对百亿级数据来说,每个分区平均要塞5亿行,Executor内存很容易撑爆;但对千万级数据,200个分区又造成大量空跑Task和调度开销。
合理设置的关键是让每个分区数据量落在100–200MB之间:
- 先估算输入总大小(如Parquet文件总字节数),除以200MB,向上取整得到目标分区数
- 在SQL前加
spark.conf.set("spark.sql.shuffle.partitions", 800)动态调整,不要写死在集群配置里 - 如果聚合后数据量骤减(比如从100亿行聚合成10万行),可考虑在
GROUP BY后接repartition(50),避免下游Stage继承过大分区数
预聚合表比实时GROUP BY更适合经营类报表
经营日报/门店排名这类场景,本质是“稳定口径+高频查询”,不是“任意维度即席分析”。每次跑GROUP BY都重算3亿行原始数据,既浪费资源,又容易因统计信息过期导致执行计划劣化。
更稳的做法是建一张预聚合表:
- 按天/小时粒度,用
INSERT OVERWRITE把明细层聚合到宽表(如sales_daily_by_store) - 聚合字段只保留报表真正需要的:不要
SELECT *后再GROUP BY,先WHERE过滤+SELECT必要列 - 给预聚合表加
CLUSTERED BY (store_id, dt),后续按门店查日报时能跳过大部分文件
预聚合不是偷懒,是把计算成本从“每次查询”转移到“数据写入时”,而写入通常是低峰期、可错峰、可重试的。
内置函数比UDF快,但AVG和COUNT(DISTINCT)仍需警惕
用sum()、count()这些内置聚合函数,Catalyst能做局部聚合(partial_agg)+合并(merge),网络传输量极小;但一旦写成UDF,就失去所有优化能力,还带序列化开销。
不过两个坑依然存在:
-
AVG(x)在分布式下等价于sum(x)/count(x),但如果x有NULL,部分引擎会把NULL当成0参与sum,结果失真——务必显式写sum(COALESCE(x, 0))/count(x) -
COUNT(DISTINCT)仍是高风险操作,即使没EXPAND,也会强制全量shuffle;优先用HLL(HyperLogLog)近似算法,或提前在ETL层用collect_set()去重
最常被忽略的一点:预聚合表的刷新逻辑必须和报表口径完全一致。今天少算了一个WHERE status = 'paid',明天所有日报指标就系统性偏低——这种错误不会报错,只会静默污染业务判断。










