Airflow DAG遇403手动抛异常仍不失败的解决方法问询
解决Airflow DAG捕获异常后任务仍成功的问题
你的核心问题是异常被完全捕获后仅记录日志,Airflow无法感知错误导致任务继续执行。以下是几种直接有效的解决方法:
方法一:捕获异常后重新抛出
在日志记录后重新抛出异常,让Airflow检测到未处理的错误并标记任务失败:
try: response = requests.get(url, params=params, headers=headers) if response.status_code == 403: raise Exception(f"Data Unavailable for metric: {m}, account: {account}") except Exception as e: logging.error(f"Exception during API request: {e} for metric: {m}, account: {account}") raise # 重新抛出异常,触发任务失败逻辑
说明:raise会将捕获的异常向上传递,Airflow任务执行器识别到未处理异常后,会立即终止任务并标记为失败。这里建议把logging.info改为logging.error,符合错误日志的级别规范。
方法二:使用Airflow专属异常类
Airflow提供了AirflowException类,专门用于标记任务失败,抛出该异常会更明确地触发失败逻辑:
from airflow.exceptions import AirflowException try: response = requests.get(url, params=params, headers=headers) if response.status_code == 403: raise AirflowException(f"Data Unavailable for metric: {m}, account: {account}") except AirflowException as e: logging.error(f"Exception during API request: {e} for metric: {m}, account: {account}") raise
说明:AirflowException是Airflow官方定义的任务失败专属异常,执行器会优先识别这类异常并终止任务,比通用Exception更贴合Airflow运行逻辑。
方法三:利用requests内置状态码异常
使用response.raise_for_status()自动将4xx/5xx状态码转换为异常,无需手动判断状态码:
import requests from airflow.exceptions import AirflowException try: response = requests.get(url, params=params, headers=headers) response.raise_for_status() # 自动抛出4xx/5xx对应的HTTPError except requests.exceptions.HTTPError as e: if response.status_code == 403: err_msg = f"Data Unavailable for metric: {m}, account: {account}" logging.error(err_msg) raise AirflowException(err_msg) else: logging.error(f"HTTP error occurred: {e}") raise except Exception as e: logging.error(f"Unexpected error: {e}") raise AirflowException(f"Unexpected error: {e}")
说明:这种方式能覆盖所有HTTP错误场景,同时针对403做自定义处理,确保任务在遇到问题时立即终止。
内容的提问来源于stack exchange,提问作者x89
相关产品推荐
相关产品推荐

