Spark SQL数据倾斜调优:Glue环境下COUNT DISTINCT性能优化问询
Spark SQL 数据倾斜优化进阶方案(COUNT DISTINCT场景)
场景与问题
在Glue g8x实例、32个Worker环境下,执行以下Spark SQL处理约3亿条记录,最终生成70万条结果,但运行耗时极长:
WITH first_cte AS ( SELECT s.two_id, s.d_month, COUNT(*) AS total, COUNT(DISTINCT one_id) AS unique_one_id_count FROM iceberg_db.tab s GROUP BY s.two_id, s.d_month ) select * from first_cte
核心问题是数据倾斜:部分two_id对应的one_id数量极大,其余two_id的one_id数量极少,导致多数任务等待少数任务完成。
已尝试的优化与排除原因
- Storage Partition Join(SPJ)、Broadcast:与Join逻辑相关,当前场景不适用
- explode/COLLECT_SET/flatten:需将所有唯一
one_id存入内存,内存压力过大,放弃 - Salting(加盐):针对COUNT DISTINCT使用RAND会导致计数变为近似值,无法满足精度要求
已采用分阶段聚合方案,性能有明显提升,但部分two_id仍存在数据倾斜问题:
WITH first_cte AS ( SELECT s.one_id, s.two_id, s.d_month, COUNT(*) AS total FROM iceberg_db.tab s GROUP BY s.one_id, s.two_id, s.d_month ), second_cte AS ( SELECT s.two_id, s.d_month, SUM(s.total) AS total, COUNT(*) AS unique_one_id_count FROM first_cte s GROUP BY s.two_id, s.d_month ) select * from second_cte
进阶优化方案
1. 热点two_id拆分处理
先识别出数据量极大的热点two_id,对热点和非热点数据分开处理,避免热点任务阻塞整体流程:
-- 第一步:识别热点two_id(阈值根据实际数据调整) WITH hot_two_ids AS ( SELECT two_id FROM ( SELECT two_id, COUNT(DISTINCT one_id) AS cnt FROM iceberg_db.tab GROUP BY two_id ) t WHERE cnt > 100000 ), -- 处理非热点数据,沿用原分阶段聚合逻辑 non_hot_first AS ( SELECT one_id, two_id, d_month, COUNT(*) AS total FROM iceberg_db.tab WHERE two_id NOT IN (SELECT two_id FROM hot_two_ids) GROUP BY one_id, two_id, d_month ), non_hot_second AS ( SELECT two_id, d_month, SUM(total) AS total, COUNT(*) AS unique_one_id_count FROM non_hot_first GROUP BY two_id, d_month ), -- 处理热点数据:对one_id哈希加盐拆分,分散计算压力 hot_data AS ( SELECT one_id, two_id, d_month, COUNT(*) AS total, MOD(ABS(HASH(one_id)), 10) AS salt -- 拆分为10组,可根据数据量调整 FROM iceberg_db.tab WHERE two_id IN (SELECT two_id FROM hot_two_ids) GROUP BY one_id, two_id, d_month, MOD(ABS(HASH(one_id)), 10) ), hot_first AS ( SELECT two_id, d_month, salt, SUM(total) AS total, COUNT(*) AS unique_one_id_count FROM hot_data GROUP BY two_id, d_month, salt ), hot_second AS ( SELECT two_id, d_month, SUM(total) AS total, SUM(unique_one_id_count) AS unique_one_id_count FROM hot_first GROUP BY two_id, d_month ) -- 合并热点与非热点结果 SELECT * FROM non_hot_second UNION ALL SELECT * FROM hot_second
该方案通过加盐将热点two_id下的one_id分散到多个子任务,既保证计数准确,又避免单个任务处理过多数据。
2. 调整Spark执行参数
- 开启自适应执行:
spark.sql.adaptive.enabled=true,让Spark根据实际数据量自动调整分区数和任务资源 - 调整Shuffle分区数:根据Worker数量设置
spark.sql.shuffle.partitions=256(比如32个Worker,每个Worker处理8个分区) - 针对倾斜优化:设置
spark.sql.adaptive.skewedPartitionThresholdInBytes=1024m,让Spark自动识别并拆分倾斜分区 - Glue环境专属:调整
--conf spark.glue.execution.partition.maxRecordsPerPartition=1000000,控制单个分区的最大记录数
3. 利用Iceberg表特性优化
- 确保查询利用Iceberg表的分区(如
d_month或two_id分区),减少扫描数据量 - 定期对Iceberg表执行
OPTIMIZE操作,合并小文件,降低读取开销 - 给
two_id和one_id建立Iceberg二级索引,加速分组过滤逻辑
内容的提问来源于stack exchange,提问作者Nir
相关产品推荐
相关产品推荐

