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

Airflow自定义异常类:如何实现任务失败自动重试并避免异常前缀?

解决Airflow自定义异常的重试与前缀问题

你遇到的问题核心是:

  • 继承AirflowException的自定义异常,Airflow默认判定为「明确业务失败」,不会触发重试
  • 继承普通Exception的自定义异常,会触发重试,但因Airflow无法序列化该异常,导致错误信息出现unusual_prefix_前缀

以下两种方案可以解决这个问题:

方案一:继承AirflowException,显式指定任务重试该异常

通过任务的retry_on_exception参数,让Airflow识别你的自定义异常需要重试:

from airflow.exceptions import AirflowException
from airflow.operators.python import PythonOperator
from datetime import timedelta

class AirflowCustomException(AirflowException):
    """自定义可重试业务异常"""
    pass

def _raise_custom_exception() -> None:
    raise AirflowCustomException("Value not found")

raise_custom_exception = PythonOperator(
    task_id='raise_custom_exception',
    python_callable=_raise_custom_exception,
    retries=3,  # 设置重试次数
    retry_delay=timedelta(seconds=10),
    # 仅当捕获到自定义异常时触发重试
    retry_on_exception=lambda e: isinstance(e, AirflowCustomException),
    dag=dag
)

方案二:继承普通Exception,注册到Airflow序列化系统

将自定义异常添加到Airflow的可序列化异常列表中,避免被包装成带前缀的异常:

from airflow.operators.python import PythonOperator
from airflow.serialization.serialized_objects import KNOWN_SERIALIZABLE_EXCEPTIONS
from datetime import timedelta

class MyCustomException(Exception):
    """自定义可重试异常"""
    pass

# 注册异常,让Airflow可以正确序列化传递
KNOWN_SERIALIZABLE_EXCEPTIONS.add(MyCustomException)

def _raise_custom_exception() -> None:
    raise MyCustomException("Value not found")

raise_custom_exception = PythonOperator(
    task_id='raise_custom_exception',
    python_callable=_raise_custom_exception,
    retries=3,
    retry_delay=timedelta(seconds=10),
    dag=dag
)

方案选择

  • 若自定义异常仅针对单个/少数任务使用,方案一更灵活
  • 若自定义异常需要在全局多个DAG/任务中复用,方案二更适合

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 09:51:22