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

SparkSQL小大表左连接性能优化及资源配置问题咨询

SparkSQL任务性能优化与资源配置问题分析

我来帮你梳理下这个SparkSQL任务的性能瓶颈和优化方向,先从当前的任务信息入手:

一、当前任务基本信息

资源配置命令

spark-sql --master yarn --driver-memory 10G --executor-memory 20G --executor-cores 20 --num-executors 20

注:未按比例申请资源,敬请谅解

执行的SQL语句

select 'ADDRESS',count(a.pid) 
FROM (
    SELECT pa.pid 
    FROM dmgr.ex_p10ids_address pa 
    WHERE pa.pt IN ('20200227') 
      AND pa.src_sys = 'APP0001' 
      AND pa.endtime = '99991231999' 
      AND pa.idtype NOT IN ('00') 
      AND certificate_type(pa.idtype, 'P10IDS') <> '0'
) a 
LEFT JOIN (
    SELECT pr.apid, pr.pid 
    FROM p10ids_riskcon pr 
    WHERE pr.classcode NOT IN ('26371100', '26371200', '26371300', '13770100', '26376000')
) b 
ON a.pid = b.apid OR a.pid = b.pid;

任务耗时

6小时33分15秒

二、源表数据量级情况

小表(dmgr.ex_p10ids_address)

查询语句:

select count(1) from (
    SELECT count(pa.pid) 
    FROM dmgr.ex_p10ids_address pa 
    WHERE pa.pt IN ('20200227') 
      AND pa.src_sys = 'APP0001' 
      AND pa.endtime = '99991231999' 
      AND pa.idtype NOT IN ('00') 
      AND certificate_type(pa.idtype, 'P10IDS') <> '0'
    -- group by pa.pid 
) t;
  • 未分组结果:46644
  • 分组后结果:45094

大表(p10ids_riskcon)

查询语句:

select count(1) from( 
    SELECT pr.apid, pr.pid 
    FROM p10ids_riskcon pr 
    WHERE pr.classcode NOT IN ('26371100', '26371200', '26371300', '13770100', '26376000')
    -- group by pr.apid, pr.pid 
) t ;
  • 未分组结果:1493862737(约15亿)
  • 分组后结果:489730113(约4.9亿)

三、当前资源配置存在的问题

  1. Executor核心数与内存不匹配:你设置了--executor-cores 20和--executor-memory 20G,平均每个核心仅分配1G内存,这会导致两个问题:
    • 内存紧张,频繁触发GC(垃圾回收),大量时间浪费在GC停顿上;
    • Hadoop/Yarn的最佳实践中,每个核心建议分配2-4G内存,这样CPU和内存才能协同高效工作。
  2. Driver内存过高:--driver-memory 10G对于这个任务来说明显过剩,Driver主要负责SQL解析、任务调度,不需要处理大量数据,通常4-6G足够,浪费的内存可以分配给Executor提升性能。
  3. 资源分配未考虑集群实际容量:20个Executor每个20核,总共需要400个集群核心,如果你的Yarn集群没有这么多空闲资源,会导致资源排队,部分Executor无法正常启动,反而拖慢任务。

四、缩短任务运行时间的优化建议

1. 调整资源配置(优先优化)

建议按照“内存-核心”比例调整,比如:

spark-sql --master yarn \
  --driver-memory 6G \
  --executor-memory 24G \
  --executor-cores 8 \
  --num-executors 15 \
  --conf spark.dynamicAllocation.enabled=true \
  --conf spark.executor.extraJavaOptions="-XX:+UseG1GC"
  • 每个Executor分配8核24G,平均每核3G,符合最佳实践;
  • 开启动态资源分配,让Spark根据任务阶段自动调整Executor数量;
  • 使用G1GC垃圾回收器,减少GC停顿时间。

2. SQL语句优化(核心优化点)

(1)优化JOIN的OR条件

原SQL中ON a.pid = b.apid OR a.pid = b.pid的OR条件会让Spark无法使用高效的哈希JOIN,只能采用性能较差的排序合并JOIN,建议拆分为两个JOIN后UNION ALL:

select 'ADDRESS', count(distinct pid) as cnt
from (
    SELECT a.pid
    FROM (小表子查询) a
    LEFT JOIN (大表子查询) b ON a.pid = b.apid
    WHERE b.apid is not null
    UNION ALL
    SELECT a.pid
    FROM (小表子查询) a
    LEFT JOIN (大表子查询) b ON a.pid = b.pid
    WHERE b.pid is not null
) t;

注:如果你的原需求是统计小表中所有pid(无论是否匹配大表),那原SQL的count(a.pid)其实等于小表的pid数量,这时候JOIN是多余的,需要确认需求是否正确。

(2)广播小表

小表仅4.5万条数据,完全可以广播到所有Executor节点,避免shuffle操作,在JOIN时添加广播提示:

select 'ADDRESS',count(a.pid) 
FROM /*+ BROADCAST(a) */ (
    -- 小表子查询
) a 
LEFT JOIN (
    -- 大表子查询
) b 
ON ...;

(3)提前对大表去重

大表未分组有15亿数据,分组后仅4.9亿,说明存在大量重复的apid+pid组合,提前去重可以减少JOIN的数据量:

SELECT DISTINCT pr.apid, pr.pid 
FROM p10ids_riskcon pr 
WHERE pr.classcode NOT IN ('26371100', '26371200', '26371300', '13770100', '26376000')

(4)优化UDF性能

certificate_type(pa.idtype, 'P10IDS')如果是自定义UDF,尽量替换为Spark内置函数,或者优化UDF的实现(比如避免在UDF中做复杂逻辑、减少对象创建),UDF的性能通常比内置函数差很多。

3. 数据存储优化

  • 检查大表p10ids_riskcon是否采用列式存储(Parquet/ORC),列式存储能大幅减少扫描的数据量,比行式存储(比如TextFile)快数倍;
  • 如果大表经常按classcode过滤,可以考虑按classcode分区;或者按apid/pid分桶,JOIN时可以减少shuffle数据。

4. 其他优化

  • 开启Spark的AQE(自适应执行计划):--conf spark.sql.adaptive.enabled=true,Spark会根据运行时的数据量自动调整执行计划,比如动态调整shuffle分区数、自动切换JOIN策略;
  • 检查小表dmgr.ex_p10ids_address是否按pt分区,WHERE pt IN ('20200227')如果命中分区,能直接跳过其他分区的扫描,减少数据读取量。

内容的提问来源于stack exchange,提问作者王 怘

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 18:52:39