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亿)
三、当前资源配置存在的问题
- Executor核心数与内存不匹配:你设置了
--executor-cores 20和--executor-memory 20G,平均每个核心仅分配1G内存,这会导致两个问题:- 内存紧张,频繁触发GC(垃圾回收),大量时间浪费在GC停顿上;
- Hadoop/Yarn的最佳实践中,每个核心建议分配2-4G内存,这样CPU和内存才能协同高效工作。
- Driver内存过高:
--driver-memory 10G对于这个任务来说明显过剩,Driver主要负责SQL解析、任务调度,不需要处理大量数据,通常4-6G足够,浪费的内存可以分配给Executor提升性能。 - 资源分配未考虑集群实际容量: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,提问作者王 怘
相关产品推荐
相关产品推荐

