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

PySpark 3.3.2 on Kubernetes动态分配Executor未释放求助

PySpark on Kubernetes中Executor任务完成后未释放的问题分析

环境与配置

使用PySpark 3.3.2部署在Kubernetes环境,核心配置如下:

spark.dynamicAllocation.enabled=true
spark.dynamicAllocation.shuffleTracking.enabled=true
spark.dynamicAllocation.executorIdleTimeout=30s
spark.dynamicAllocation.cachedExecutorIdleTimeout=40s
spark.executor.instances=0
spark.dynamicAllocation.minExecutors=0
spark.dynamicAllocation.maxExecutors=50
spark.cleaner.periodicGC.interval=30s
spark.master=k8s://https://<...>:6443

测试场景对比

场景1:自定义UDF任务(Executor未释放)

执行以下代码后,Executor完成任务后始终无法被回收:

import time
import random
from pyspark.sql import functions as F

instances = 50
df = spark.createDataFrame([{"a": i} for i in range(instances * 3)])
df = df.repartition(instances)

@F.udf()
def udf():
    time.sleep(random.randint(2, 10))
    return random.randint(1, 42)

print(df.withColumn('random', udf()).agg(F.sum('random').alias('sum')).collect()[0].sum)
time.sleep(999) # executors never get removed

场景2:Parquet读取任务(Executor正常释放)

执行读取Parquet并计数的代码后,Executor会在阶段结束后20-30秒内正常释放:

df = spark.read.parquet("s3a://<directory with 1TB of parquet files>")
print(df.count())
time.sleep(999) # this is fine, all executors get removed within 20-30s after stage end

原因分析

这不属于完全的异常行为,核心差异源于Shuffle资源的残留状态和Executor的Idle状态识别逻辑:

  1. Shuffle Tracking的残留问题:场景1中repartition操作会触发Shuffle,虽然开启了shuffleTracking.enabled,但PySpark UDF执行后,Executor上的Shuffle元数据或资源引用可能未被及时清理。DynamicAllocation会认为Executor仍持有Shuffle资源,不会触发回收。
  2. GC延迟导致状态未更新:自定义UDF中的time.sleep可能导致Executor上的对象未被及时GC,结合periodicGC.interval=30s的间隔,Idle状态无法被及时上报给Driver,延迟了回收逻辑触发。
  3. Parquet场景的无残留特性:场景2的count()操作不会产生持久化的Shuffle资源,任务完成后Executor无残留状态,DynamicAllocation能快速识别Idle状态并执行回收。

解决方案

  • 添加Shuffle跟踪超时配置:
    spark.dynamicAllocation.shuffleTracking.timeout=60s
    
    该配置会强制清理超时的Shuffle跟踪数据,让Executor被正确标记为Idle。
  • 优化Executor GC配置:在spark.executor.extraJavaOptions中添加GC参数,加快垃圾回收:
    spark.executor.extraJavaOptions=-XX:+UseG1GC -XX:MaxGCPauseMillis=200
    
  • 避免Driver阻塞:测试时不要用time.sleep(999)长时间阻塞Driver,Driver阻塞会影响状态同步。若需观察自动回收,确保Driver处于正常运行状态;任务完成后可主动调用spark.stop()触发资源回收。
  • 调整Idle超时时间:适当延长spark.dynamicAllocation.executorIdleTimeout至60s,排除K8s节点状态上报延迟的影响。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 16:16:01