Pyspark/Databricks无原因两次执行for循环问题如何解决?
问题根因
该异常不是循环主动重启,是Spark的惰性求值机制+任务自动重试导致的:
- 代码中
do_stuff_to_file_and_save_resulting_dataframe_to_curated_folder涉及的Spark转换操作(如spark.read.parquet、数据清洗逻辑)属于惰性执行,仅会生成执行计划,直到触发行动操作(如写文件、count、collect)才会在executor端分布式执行 - driver端的for循环执行完成后,所有原文件已经被移动到归档目录,此时如果executor端的任意任务执行失败(如网络波动、资源抢占),Spark会按照默认配置自动重试任务,重试时会重新执行读原文件的逻辑,但原文件已被移除,就会抛出文件不存在的异常
- 直接遍历
dbutils.fs.ls返回的FileInfo对象,没有提前固化路径为静态字符串,也可能导致Spark执行时动态解析路径获取到过期信息
修复方案
可以按优先级选择以下方案:
- 提前固化文件路径为静态字符串列表,避免后续动态解析
FileInfo对象带来的异常
def f(files): # 提前提取所有文件的路径为静态字符串,切断和原FileInfo对象的关联 file_paths = [file.path for file in dbutils.fs.ls(files)] for i, file_path in enumerate(file_paths): print(f"\n Trying to read file = {file_path}, loop index = {i} \n") do_stuff_to_file_and_save_resulting_dataframe_to_curated_folder(file_path, curated_folder) archive_file_by_moving_to_a_different_folder(file_path, destination_folder)
- 调整文件归档时机,确认当前文件的所有处理、写操作完全执行完成后,再执行移动操作。可以在处理完每个文件后增加轻量的行动操作,确认executor端任务全部执行完成:
def do_stuff_to_file_and_save_resulting_dataframe_to_curated_folder(file_path, curated_folder): df = spark.read.parquet(file_path) # 原有数据处理逻辑 processed_df = df.filter(xxx) # 写结果到目标目录 processed_df.write.mode("append").parquet(curated_folder) # 新增轻量行动操作,强制确认写操作完成后再执行后续归档逻辑 spark.read.parquet(curated_folder).limit(1).count()
- 若不需要任务自动重试,可调整Spark配置降低重试次数,避免无意义的重试:
# 代码开头添加配置,将任务最大重试次数设为1 spark.conf.set("spark.task.maxFailures", "1")
- 稳定性要求高的场景下,可以先处理所有文件,确认全部处理成功后再批量归档原文件,彻底避免处理过程中移走原文件的问题:
def f(files): file_paths = [file.path for file in dbutils.fs.ls(files)] # 第一步:处理所有文件 for i, file_path in enumerate(file_paths): print(f"\n Trying to read file = {file_path}, loop index = {i} \n") do_stuff_to_file_and_save_resulting_dataframe_to_curated_folder(file_path, curated_folder) # 第二步:全部处理成功后再批量归档 for file_path in file_paths: archive_file_by_moving_to_a_different_folder(file_path, destination_folder)
内容的提问来源于stack exchange,提问作者Mateusz Makowski
相关产品推荐
相关产品推荐

