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状态识别逻辑:
- Shuffle Tracking的残留问题:场景1中
repartition操作会触发Shuffle,虽然开启了shuffleTracking.enabled,但PySpark UDF执行后,Executor上的Shuffle元数据或资源引用可能未被及时清理。DynamicAllocation会认为Executor仍持有Shuffle资源,不会触发回收。 - GC延迟导致状态未更新:自定义UDF中的
time.sleep可能导致Executor上的对象未被及时GC,结合periodicGC.interval=30s的间隔,Idle状态无法被及时上报给Driver,延迟了回收逻辑触发。 - Parquet场景的无残留特性:场景2的
count()操作不会产生持久化的Shuffle资源,任务完成后Executor无残留状态,DynamicAllocation能快速识别Idle状态并执行回收。
解决方案
- 添加Shuffle跟踪超时配置:
该配置会强制清理超时的Shuffle跟踪数据,让Executor被正确标记为Idle。spark.dynamicAllocation.shuffleTracking.timeout=60s - 优化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
相关产品推荐
相关产品推荐

