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

Airflow中pandas读取xlsx触发BadZipFile报错排查

Airflow DAG 偶发BadZipFile读取xlsx异常修复方案

核心根因

这个偶发报错和openpyxl版本、文件结构损坏无关,本质是竞态条件+代码鲁棒性不足导致的:

  • xlsx格式本质是zip压缩包,你当前的逻辑是按文件创建时间排序后直接取第一个匹配日期的文件读取,但RPA生成xlsx时会先创建文件占位、再逐步写入内容,只要DAG在RPA写完文件前触发读取,就会拿到zip头不完整的半成文件,触发BadZipFile: File is not a zip file报错,这也是问题偶发的核心原因。
  • 现有逻辑没有过滤RPA生成的~$开头的临时锁文件、没有做文件合法性校验,很容易把非目标文件读进来。
  • 贴出的代码存在变量名笔误:定义的变量是lastfile,读取时传入的是last_file,如果不是贴代码时的手误,线上还会偶发NameError。

可直接落地的修复方案

修复逻辑

  • 第一步:提前过滤目录、临时文件、非xlsx后缀的无效文件,只保留目标日期范围内的候选文件
  • 第二步:新增文件完整性校验:先通过zipfile.is_zipfile()判断文件是不是完整的xlsx包,再做读取操作
  • 第三步:增加重试机制,碰到未写完的文件等待几秒后重试,当前文件损坏时自动按创建时间顺序尝试下一个候选文件
  • 第四步:增加兜底异常,所有候选文件都读取失败时抛出明确的错误信息,方便排查

修复后代码

import os
import glob
import time
import zipfile
from datetime import datetime
import pandas as pd

# 可根据RPA写文件的实际速度调整参数
MAX_RETRY = 3
RETRY_WAIT = 2  # 单次重试等待秒数
def is_file_completed(file_path):
    """校验文件是否写入完成、是合法的xlsx文件"""
    if not os.path.isfile(file_path):
        return False
    # 连续两次校验文件大小一致,排除正在写入/传输的文件
    size1 = os.path.getsize(file_path)
    time.sleep(1)
    size2 = os.path.getsize(file_path)
    if size1 != size2 or size1 == 0:
        return False
    # 校验是合法zip格式(xlsx本质为zip包)
    return zipfile.is_zipfile(file_path)

files_path = os.path.join(local_path, '*')
# 按创建时间倒序排列文件
files = sorted(glob.iglob(files_path), key=os.path.getctime, reverse=True)
candidate_files = []

# 筛选符合日期要求的有效候选文件
for file in files:
    file_name = os.path.basename(file)
    # 跳过目录、临时锁文件、隐藏文件、非xlsx文件
    if os.path.isdir(file) or file_name.startswith(('~$', '.')) or not file_name.lower().endswith('.xlsx'):
        continue
    # 匹配目标日期
    file_create_date = datetime.fromtimestamp(os.path.getctime(file)).strftime('%Y-%m-%d')
    if file_create_date == dateday:
        candidate_files.append(file)

df = None
target_file = None
# 从新到旧遍历候选文件,找到第一个可正常读取的文件
for file in candidate_files:
    read_ok = False
    for retry_cnt in range(MAX_RETRY):
        try:
            if not is_file_completed(file):
                time.sleep(RETRY_WAIT)
                continue
            df = pd.read_excel(file, engine='openpyxl')
            target_file = file
            read_ok = True
            break
        except zipfile.BadZipFile:
            # 文件未写完,等待后重试
            time.sleep(RETRY_WAIT)
        except Exception:
            # 其他读取错误直接跳过当前文件,尝试下一个
            break
    if read_ok:
        break

# 兜底校验
if df is None:
    raise RuntimeError(f"路径{local_path}下未找到日期为{dateday}的可读取有效xlsx文件")

# 后续业务逻辑处理df即可

额外优化建议

  • 如果RPA和Airflow通过NFS/SMB等共享存储传输文件,建议把DAG的触发时间延后1-2分钟,给文件落盘、传输留足缓冲时间,能大幅降低竞态概率
  • 检查RPA端的文件保存逻辑,确认写入xlsx后正常关闭文件句柄、释放文件锁,不要在文件打开状态下就触发DAG调度
  • 如果环境允许,RPA写完文件后可以生成一个和xlsx同名的.ok标识文件,DAG侧只扫描存在对应.ok标识的xlsx文件,从根源上避免读到半成文件

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 23:24:26