如何让Airflow DAG在Python代码抛出异常时标记为失败?
如何让Airflow DAG在Python代码抛出异常时标记为失败?
嘿,我来帮你搞定这个头疼的问题!明明代码里抛了异常,Airflow却把DAG标成成功,只在日志里打个INFO,这肯定是代码逻辑或者异常处理出了问题,咱们一步步拆解修复:
首先看你代码里的核心问题
1. all_accounts_failed的逻辑完全走偏了
你现在的代码里,只有df1非空时才把all_accounts_failed设为False,但如果所有API都失败,df1是空的,这个变量根本没被赋值,直接进if all_accounts_failed:要么报错,要么因为变量未定义根本触发不了raise。正确的姿势是一开始就把这个变量初始化为True,只要有一个API请求成功(也就是df1新增了数据),就把它改成False。
2. 异常抛出的逻辑有冗余且位置不对
你写的raise Exception(...)后面还跟了个多余的raise,这行代码永远不会执行,纯纯多余。另外,如果你的else块里的try没有对应的except,或者外层有其他try-except把异常吞了,Airflow就捕获不到异常,自然会认为任务成功。
3. 循环里的break逻辑不符合你的需求
你现在只要遇到一个非200的响应就break,直接终止循环,这会导致后面的ID都不尝试了,根本没法判断是不是所有API都失败,顶多只能判断第一个就失败了。
给你修正后的完整代码示例
# 先初始化必要的变量,避免未定义报错 all_accounts_failed = True # 默认所有请求都失败 df1 = pd.DataFrame() # 初始化空DataFrame # 假设你有要遍历的ID列表,比如ids = [1, 2, 3, ...] for id in ids: response = requests.get(url) data = response.json() if response.status_code == 200: # 只要有一个请求成功,就把标记改成False all_accounts_failed = False # 拼接数据 df1 = pd.concat([df1, pd.json_normalize(data)]) print(f"成功获取ID {id} 的数据") else: print(f"ID {id} 请求失败,状态码:{response.status_code},原因:{response.reason}") # 这里不要加break,要遍历完所有ID才能判断是否全失败 # 遍历完所有ID后,判断是否全失败 if all_accounts_failed: # 直接抛出异常,Airflow会自动捕获并标记任务为失败 raise Exception("所有API请求都失败了!") # 下面是正常的业务逻辑,不需要套多余的try(除非你有特定的局部异常处理需求) # rest of the code
最后再提醒几个容易踩的坑
- 绝对不要在任务函数的外层加大的
try-except把异常吞掉!比如如果你的整个任务函数被try: ... except: print(...)包裹,异常就被悄悄处理了,Airflow根本看不到,自然会认为任务成功。 - 确保你用的是
PythonOperator(或者其他支持执行Python函数的Operator),这类Operator默认会把函数里未捕获的异常视为任务失败,只要你不手动吞异常,它就会正确标记任务状态。 - 之前你代码里
if not df1.empty: all_accounts_failed = False的逻辑也有问题,因为df1为空时这个变量根本没被赋值,会触发NameError,所以一定要先初始化变量。
按照这个逻辑改完后,只要所有API都失败,代码就会抛出异常,Airflow会立刻把这个任务标记为失败,整个DAG也会跟着显示失败状态啦!
备注:内容来源于stack exchange,提问作者S_2103
相关产品推荐
相关产品推荐

