Databricks中PySpark循环停滞且无Spark任务运行问题排查
问题背景
在Databricks中运行for循环生成模拟增量数据时,初始迭代速度快,随后逐渐变慢直至完全停滞。已排除数据量增加、未清理垃圾变量、RAM/CPU/DISK/NETWORK满载等常见原因,停滞时SparkUI无正在处理的任务,但Notebook单元格仍显示运行中。
需求是基于基础DataFrame input_df,生成400个Schema相同的CSV文件(仅created_time逐天变化),模拟每日增量加载场景,代码如下:
another=input_df for i in range(400): # 获取id列的最大最小值,用于生成新数据的id aggs=another.agg(max('id'),min('id')) max_id=aggs.collect()[0][0] min_id=aggs.collect()[0][1] # 修改id和created_time列,模拟新数据 another=another.withColumn('id',col('id')+max_id-min_id+1).\ withColumn('created_time',date_add(col('created_time'),1)).\ withColumn('created_time',date_format(col("created_time"), "yyyy-MM-dd'T'HH:mm:ss.SSSZ")) # 生成文件名对应的日期 date_procesed=datetime.strptime('20220112','%Y%m%d') + timedelta(days=i+1) date_procesed=date_procesed.strftime('%Y%m%d') print(date_procesed) # 写入单CSV文件到S3 another.coalesce(1).write.option('header','true').csv('dbfs:/tmp/wiki/transaction/'+date_procesed)
停滞现象
- 循环执行约11次(对应40个已完成Spark任务,S3文件已生成)后停滞
- SparkUI中无正在运行的任务,仿佛未创建新任务
- GangliaUI显示仅Driver在工作,但CPU/RAM/NETWORK/DISK均未满载,且近一小时无波动
- Driver存在大量未终止进程,怀疑Spark未关闭端口/连接或未标记进程完成
待解答问题
- 是否存在RAM/CPU/DISK/NETWORK之外的资源阻塞?
- 为何for循环仍显示运行但SparkUI无任务?
- 该for循环失效的原因是什么?
解答
1. 存在其他资源阻塞
是,主要是Spark作业元数据累积、Driver端任务队列阻塞以及文件句柄泄漏:
- 每次执行
agg和collect都会生成作业元数据,循环迭代中这些元数据会在Driver端持续累积,占用Driver的堆外内存或元数据存储资源,即使RAM未满载,也会导致元数据处理逻辑卡顿。 - 频繁调用
coalesce(1).write会创建大量小文件,每个写操作都会占用文件句柄,若Driver未及时释放,会达到系统文件句柄上限,导致后续写操作无法发起。
2. for循环显示运行但SparkUI无任务的原因
循环停滞在Driver端的非Spark作业逻辑中,而非Spark任务执行阶段:
- 当Driver的元数据累积到一定程度,处理
another.agg(...)或collect()时,会陷入本地计算的卡顿(比如元数据遍历、对象序列化/反序列化),这部分逻辑属于Driver本地代码,不会触发Spark作业,所以SparkUI无任务,但Notebook单元格仍显示运行(因为Driver进程还在执行本地代码)。 - 另外,若文件句柄耗尽,Driver尝试发起写操作时会被系统阻塞,这部分也属于本地IO阻塞,不会生成Spark任务。
3. 循环失效的核心原因
(1)DataLineage无限累积
每次循环中another = another.withColumn(...)都会在原有DataFrame的血统(Lineage)上叠加新的转换操作,导致Lineage链越来越长。当Lineage过长时,Driver解析执行计划会变得异常缓慢,甚至陷入死循环或内存溢出前的卡顿。即使每次只修改少量列,Spark仍需要维护整个转换链的元数据,迭代次数越多,解析成本越高。
(2)频繁collect()操作拖垮Driver
每次循环都调用aggs.collect(),这会将聚合结果拉取到Driver端。虽然单次数据量小,但400次迭代的累积操作会让Driver的GC压力剧增,同时每次collect都会触发Spark作业,生成的作业元数据不断累积,占用Driver的元数据存储资源,最终导致Driver处理逻辑卡顿。
(3)coalesce(1)的低效与单点瓶颈
coalesce(1)会将所有数据shuffle到单个Executor节点,每次写操作都需要该节点处理全量数据,频繁的单节点IO会导致该节点的磁盘IO队列阻塞,同时Driver需要等待该节点完成写操作才能继续循环。即使Ganglia显示整体资源未满载,单个节点的IO瓶颈也会拖慢整个循环,甚至导致后续写操作无法发起。
(4)未释放无用资源
每次写操作后,未对another DataFrame进行缓存清理或Lineage截断,导致Driver端的对象引用无法被GC回收,内存中累积大量无用的DataFrame元数据和对象,进一步加重Driver的负担。
内容的提问来源于stack exchange,提问作者Eugenio.Gastelum96

