You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 20:03:23