Spark-SQL查询执行失败求助:结果序列化大小超driver配置限制
嗨,我来帮你搞定这个问题。你遇到的错误是因为Driver端接收的任务结果总大小超过了配置上限——33个任务的序列化结果加起来有1046.8MB,超过了spark.driver.maxResultSize默认的1024MB限制。下面我分步骤给你解决方法和调优建议:
一、快速让查询跑起来的临时方案
如果只是想先看到结果,可以临时调大Driver的结果大小限制:
- 提交任务时添加配置参数:
或者在Spark SQL控制台先执行:--conf spark.driver.maxResultSize=2g
注意:别把这个值设得太夸张,要是超过Driver所在机器的内存,反而会导致Driver内存溢出。如果你的结果集本身就特别大,这只是临时 workaround,更推荐从查询本身优化。SET spark.driver.maxResultSize=2g;
二、核心调优:减少返回给Driver的数据量
你之前缓存contrib表没用,是因为这个问题不是计算阶段重复读表导致的,而是最终要返回给Driver的结果集太大,装不下。所以核心思路是要么缩小结果的大小,要么避免把结果拉到Driver端。
1. 先修正查询里的语法错误
看你的原查询,flg IN ('001', '004' and covr_fr > '2018-01-01' )这里有语法问题——IN的括号没闭合,还多了个and,应该改成:
select * from contrib where ind_cd = 'BP' and flg IN ('001', '004') and covr_fr > '2018-01-01'
另外,你的子查询里漏了聚合需要用到的fc_am和mc_am字段,外层聚合时会找不到这两个字段,后面我会在修正后的完整查询里补上。
2. 提前过滤无用数据,减少聚合压力
你的contrib表过滤后还有6.8亿条数据,先看看能不能在聚合前过滤掉更多无效行:
- 如果业务允许,能不能缩小
covr_fr的范围?比如只查2018年之后的某几年数据,而不是全部? - 可以加个条件过滤掉
fc_am和mc_am都为0的行,这些行对最终的聚合结果没有贡献:select * from contrib where ind_cd = 'BP' and flg IN ('001', '004') and covr_fr > '2018-01-01' and (fc_am != 0 or mc_am != 0)
3. 优化聚合逻辑,缩小结果集
你是按policy, order, covr_fr, covr_to分组,返回聚合后差值大于0的结果:
- 检查分组字段是否有冗余?比如
covr_fr和covr_to是不是可以按月份聚合?如果业务允许,把时间粒度调大,能大幅减少分组的数量。 - 如果不需要精确到每一天的
covr_fr和covr_to,可以用date_trunc('month', covr_fr)这类函数把日期转成月份,再分组。
4. 改用更高效的序列化方式
Spark默认用Java序列化,换成Kryo序列化能显著减少数据序列化后的大小(通常是Java序列化的1/5左右):
- 提交任务时加配置:
或者在Spark SQL里设置:--conf spark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.kryo.registrationRequired=falseSET spark.serializer=org.apache.spark.serializer.KryoSerializer; SET spark.kryo.registrationRequired=false;
5. 避免把大结果集拉到Driver端(最彻底的解决方法)
如果你的最终结果集确实很大(比如数百万行以上),别直接让Driver接收结果。把结果写到外部存储(比如HDFS、Hive表、数据库)里,修正后的完整查询如下:
INSERT INTO your_result_table -- 替换成你的目标表名 select policy, `order`, covr_fr, covr_to, sum(fc_am) - SUM(mc_am) As amount from ( select coalesce(cntr.vre2, cntr.cacont_acc) as policy, cntr.`order`, date_dim_get_sk.date_dim_sk as pst_dt_sk, cntr.covr_fr, cntr.covr_to, cntr.fc_am, cntr.mc_am -- 补上聚合需要的字段 from ( select * from contrib where ind_cd = 'BP' and flg IN ('001', '004') and covr_fr > '2018-01-01' and (fc_am != 0 or mc_am != 0) ) cntr JOIN date_dim ON date_dim.dt = cntr.pstng_dt JOIN date_dim_get_sk ON date_dim_get_sk.dt = date_dim.dt_lst_dayofmon ) t GROUP BY policy, `order`, covr_fr, covr_to HAVING sum(fc_am) - SUM(mc_am) > 0;
这样结果直接写入存储,不会经过Driver,从根本上解决结果集过大的问题。
三、其他性能调优建议
- 优化Join性能:确保
date_dim和date_dim_get_sk的dt字段是分区字段或者有索引(如果是Parquet/ORC这类列式存储),减少Join时的Shuffle数据量。 - 分桶表优化:如果
contrib表经常被按ind_cd、flg过滤,可以把它做成分桶表,按ind_cd或flg分桶,能大幅提升过滤和Join的速度。 - 调整Shuffle分区数:Spark默认的Shuffle分区数是200,对于6.8亿条数据来说太少了,可以调大到600或1000:
这样每个Shuffle分区的数据量不会太大,避免单个任务处理过多数据。SET spark.sql.shuffle.partitions=800;
先按上面的步骤修正语法错误,然后优先尝试把结果写入外部存储,这是最彻底解决问题的方法。如果还有疑问,可以告诉我更多业务细节,我再帮你调整。
内容的提问来源于stack exchange,提问作者marie20

