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
相关产品推荐
相关产品推荐

