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

Spark作业仅1个Executor运行1个任务的问题求助

问题分析与解决方案

核心原因排查

  • 首先确认load_dt字段的基数(不同值的数量):如果该字段的唯一值数量极少(比如仅1个),Spark执行group by load_dt时,shuffle后只会生成对应数量的分区,自然仅1个任务运行,其余Executor处于空闲状态。
  • 验证方式:执行以下SQL查看基数:
    select count(distinct load_dt) from df_missing
    

针对性解决方案

1. 调整Shuffle分区数(若基数足够)

Spark执行聚合类操作时,默认通过spark.sql.shuffle.partitions参数控制shuffle后的分区数(默认值通常为200)。你之前对df_missing执行的repartition(100)仅作用于该DataFrame的读取阶段,不影响后续group by的shuffle分区设置。

  • 解决步骤:在聚合SQL前添加参数配置:
    spark.sql("SET spark.sql.shuffle.partitions=100")
    
    之后再执行你的分组计数SQL,Spark会按设置的分区数并行处理shuffle任务。

2. 处理低基数场景(若基数极少)

如果load_dt的唯一值确实很少,直接group by无法并行,可通过二次聚合拆分计算压力:

  • 第一步按load_dt+哈希拆分键分组,并行计算局部计数:
    create or replace temp view df_partial_cnt as
    select load_dt, hash(tid) % 100 as part, count(*) as partial_cnt
    from df_missing
    group by load_dt, part
    
  • 第二步对局部计数求和得到最终结果:
    select load_dt, sum(partial_cnt) as TID_DEL
    from df_partial_cnt
    group by load_dt
    
    这种方式让第一步分组计算并行执行(100个任务),仅第二步全局聚合为单任务,可大幅缩短总耗时。

3. 验证分区生效情况

确认df_missing的分区数是否符合预期,避免repartition未生效:

print(df_missing.rdd.getNumPartitions())

若输出不是100,可尝试将repartition(100)替换为coalesce(100, shuffle=True)强制触发shuffle,确保分区设置生效。


内容的提问来源于stack exchange,提问作者A Raks

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 13:05:27