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

Airflow任务重试后误标记为成功及BQ配额超限问题求助

问题分析与解决方案

你的核心问题出在writing_to_bq函数的异常处理逻辑上:当捕获到BQ配额超限的403错误时,函数仅记录了重试日志,既没有在重试次数耗尽后抛出异常,也没有向Airflow传递任务失败的信号,导致Airflow误判任务执行成功,进而触发后续流程。

具体修复步骤

1. 修正异常处理逻辑,重试耗尽后主动抛出异常

当前函数在每次捕获异常后仅打印日志,循环结束后无论是否成功写入BQ,函数都会正常返回。需要在重试次数用完仍失败时,重新抛出异常,让Airflow捕获并标记任务失败。

2. 统一参数名(避免调用与定义不匹配)

调用函数时传的是retry_number=3,但函数定义的参数是retry,存在参数名不匹配的问题,需要统一。

3. (可选)针对性捕获BQ配额错误

可以只捕获BQ相关的特定异常,避免捕获无关异常导致的误处理(比如DataFrame格式错误等)。

修改后的代码示例

# 调用时统一参数名
writing_to_bq(df=ent1, table=table, retry=3) 

def writing_to_bq(df, table, retry):  
    for attempt in range(retry + 1):  
        try:  
            df.to_gbq(table, if_exists="append")
            log.info(f"BQ写入成功,共尝试{attempt+1}次")
            return  # 写入成功后直接返回
        except Exception as e:
            error_msg = str(e)
            if "append/update limit exceeded" in error_msg:
                # 识别到BQ配额超限错误
                if attempt == retry:
                    # 重试次数耗尽,抛出异常让Airflow标记任务失败
                    log.error(f"BQ写入失败,已耗尽{retry+1}次重试机会,错误信息:{error_msg}")
                    raise e
                else:
                    log.info(f"BQ写入失败,将进行第{attempt+2}次重试,错误信息:{error_msg}")
            else:
                # 非配额类错误,直接抛出,不进行重试(根据需求调整)
                log.error(f"BQ写入遇到非配额错误,直接终止任务,错误信息:{error_msg}")
                raise e

额外说明

  • 当函数抛出异常后,Airflow会根据你DAG中设置的重试次数(5次)自动触发任务重试,直到重试次数耗尽后标记任务失败。
  • 若需要针对BQ配额错误设置更长的重试间隔(比如等待1小时后重试,避开配额刷新时间),可以在Airflow任务的retry_delay参数中配置,不推荐在函数内部添加time.sleep(),会占用worker资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 07:21:24