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

Python中如何处理FileNotFoundException?Avro文件聚合场景排障

解决Spark读取空文件夹时的FileNotFoundException捕获问题

问题根源

你当前捕获的是Python标准库的FileNotFoundException,但Spark抛出的是PySpark封装的Java异常(实际类型为pyspark.sql.utils.AnalysisException),所以原try-except块无法命中该异常。此外你的代码里还有一处语法错误:.option("ignoreExtension":"true")应该改为.option("ignoreExtension", "true")(用逗号而非冒号)。

正确处理方案

方案1:捕获PySpark的AnalysisException

导入PySpark的异常类,通过判断异常消息中的特征文本,精准处理“无Avro文件”的情况:

from pyspark.sql.utils import AnalysisException
import os
import datetime
import subprocess

def aggregate_data(path_to_source, input_format, output_format, output_table_format='avro'):
    all_input_dates = [str(fileStatus.getPath()).split('/')[-1] for fileStatus in
                        lake.listStatus(fs.Path(path_to_source))
                        if 'latest' not in str(fileStatus.getPath())]
    
    # 修复原代码集合运算的括号位置错误
    valid_dates = [i for i in all_input_dates if validate_date(i, input_format)]
    today_output = datetime.date.today().strftime(output_format)
    all_output_dates = list(set(datetime.datetime.strptime(i, input_format).strftime(output_format)
                        for i in valid_dates) - {today_output})
                        
    print("***** {0} *****".format(path_to_source))
    
    try:
        for _partition in all_output_dates:
            src_path = os.path.join(path_to_source, '{0}*'.format(_partition))
            # 修复option参数的语法错误
            src_df = spark.read.option("ignoreExtension", "true")\
                        .format(output_table_format)\
                        .load(src_path)
            
            agg_df = src_df.repartition(estimate_part_num(src_df))\
                        .write.mode('overwrite')\
                        .format(output_table_format)\
                        .save(os.path.join(path_to_source, _partition))
            
            command_exist = ['hadoop', 'fs', '-rm', '-r', os.path.join(path_to_source, '{0}[.]*'.format(_partition))]
            subprocess.call(command_exist)
    except AnalysisException as e:
        # 通过消息特征判断是否为目标异常
        if "No avro files found" in str(e):
            print(f"**** 分区 {_partition} 无Avro文件,跳过处理 ****")
        else:
            # 非目标异常重新抛出,不吞掉其他错误
            raise e

方案2:读取前先检查文件存在性

在尝试读取前,先通过HDFS API检查目标路径下是否有Avro文件,从根源避免异常触发:

import os
import datetime
import subprocess

def aggregate_data(path_to_source, input_format, output_format, output_table_format='avro'):
    all_input_dates = [str(fileStatus.getPath()).split('/')[-1] for fileStatus in
                        lake.listStatus(fs.Path(path_to_source))
                        if 'latest' not in str(fileStatus.getPath())]
    
    valid_dates = [i for i in all_input_dates if validate_date(i, input_format)]
    today_output = datetime.date.today().strftime(output_format)
    all_output_dates = list(set(datetime.datetime.strptime(i, input_format).strftime(output_format)
                        for i in valid_dates) - {today_output})
                        
    print("***** {0} *****".format(path_to_source))
    
    for _partition in all_output_dates:
        src_path = os.path.join(path_to_source, '{0}*'.format(_partition))
        # 检查路径下是否存在Avro文件
        file_statuses = []
        try:
            file_statuses = lake.globStatus(fs.Path(src_path + ".avro"))
        except Exception:
            pass
        
        if not file_statuses:
            print(f"**** 分区 {_partition} 无Avro文件,跳过处理 ****")
            continue
            
        src_df = spark.read.option("ignoreExtension", "true")\
                    .format(output_table_format)\
                    .load(src_path)
        
        agg_df = src_df.repartition(estimate_part_num(src_df))\
                    .write.mode('overwrite')\
                    .format(output_table_format)\
                    .save(os.path.join(path_to_source, _partition))
        
        command_exist = ['hadoop', 'fs', '-rm', '-r', os.path.join(path_to_source, '{0}[.]*'.format(_partition))]
        subprocess.call(command_exist)

关键说明

  • PySpark的IO类异常大多封装在AnalysisException中,需要针对性捕获。
  • 优先通过预检查避免异常,比事后捕获更高效且逻辑更清晰。
  • 原代码中的集合运算括号位置错误,已在修复示例中修正。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 08:40:12