Apache Airflow Tasks 深度指南:从任务依赖到重试策略与超时控制
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
本篇技术指南以 Apache Airflow 官方核心概念文档 tasks.rst 为骨架,系统讲解 Airflow 中"任务(Task)"这一基本执行单元:如何定义任务、声明上下游依赖、理解 Task Instance 生命周期与状态机、配置超时与重试策略,以及如何借助特殊异常、executor_config 等机制精确控制任务行为。读完本文,你将掌握编写健壮、可控、可运维 DAG 任务所需的全部核心知识,并了解其底层实现依据(对应源码见 state.py、retry_policy.py)。
什么是 Task:Airflow 中的基本执行单元
在 Airflow 中,Task 是最基本的执行单元。Task 被组织进 DAG(有向无环图)中,并通过在任务之间设置 upstream(上游)与 downstream(下游)依赖来表达它们应当执行的顺序。
Airflow 中一共有三种基本的 Task:
- Operators(操作符):预定义好的任务模板,可以快速拼接出 DAG 的大部分组成部分,例如执行 Bash 命令、传输文件、调用云服务 API 等,详见 operators.rst。
- Sensors(传感器):Operator 的一个特殊子类,其全部职责就是等待某个外部事件发生(例如等待某个文件出现、等待某个 API 可用),详见 sensors.rst。
@task装饰的 TaskFlow 任务:把一个自定义的 Python 函数包装成 Task,详见 taskflow.rst。
从实现角度看,这三种形式在内部全部都是BaseOperator的子类,因此 Task 与 Operator 的概念在某种程度上可以互换。但更准确的思考方式是:Operator 和 Sensor 是"模板"(template),当你在 DAG 文件中调用它们一次,就产生了一个 Task(任务实例化的定义)。
关系(Relationships):如何声明任务依赖
使用 Task 的关键在于定义它们彼此之间的关系,即依赖(dependencies),Airflow 中称为upstream(上游)与downstream(下游)任务。通常的做法是:先声明所有 Task,再声明它们的依赖关系。
注意:所谓 upstream 任务,是指直接排在另一个任务前面的那个任务(旧文档中曾称之为 parent task)。需要留意的是,这一概念并不描述任务层级中更高层级的任务(即并非该任务的间接祖先)。downstream 任务的定义同理,它必须是另一个任务的直接后继。
两种声明依赖的方式
Airflow 提供了两种声明依赖的语法:
方式一:位运算操作符>>与<<(推荐)
first_task >> second_task >> [third_task, fourth_task]方式二:显式方法set_upstream与set_downstream
first_task.set_downstream(second_task) third_task.set_upstream(second_task)这两种写法实现的效果完全相同,但官方推荐优先使用位运算操作符,因为它在大多数场景下可读性更强——箭头方向直观地表达了数据/执行的流动方向。
默认依赖语义与高级控制
默认情况下,一个 Task 会在其所有upstream 任务成功后运行。但 Airflow 提供了大量方式修改这一行为,例如:
- 引入分支(branching);
- 只等待部分 upstream 任务(例如通过
TriggerRule.ONE_SUCCESS、ALL_DONE等触发规则); - 根据当前运行在历史中的位置(如是否为回填、是否为补数据运行)改变行为。
这些高级控制方式参见 dags.rst 中的控制流章节(trigger rules)以及 backfill.rst。
此外需要注意:Task 默认不向彼此传递信息,彼此完全独立运行。如果需要在任务之间传递数据,应当使用 XComs(跨任务通信机制),详见 xcoms.rst。
Task Instance:任务的一次具体运行
正如 DAG 每次运行会被实例化为Dag Run,DAG 下的每个 Task 也会被实例化为Task Instance(任务实例)。
一个 Task Instance 是"该任务在给定 DAG(因而也是给定数据区间 data interval)下的一次具体运行"。它同时也是**拥有状态(state)**的任务表示,反映它处于生命周期的哪个阶段。当任何自定义 Task(Operator)运行时,它都会获得一份 Task Instance 的拷贝;除了可以检查任务元数据外,Task Instance 还携带诸如 XComs 读写之类的方法。
Task Instance 的完整状态集
Task Instance 可能的状态如下(对应源码枚举定义见 state.py 中的TaskInstanceState,其中IntermediateTIState表示尚未进入终态/运行态的中间状态,TerminalTIState表示终态):
| 状态 | 含义 |
|---|---|
none | 任务尚未被排队执行(其依赖尚未满足),由调度器创建但尚未运行时使用None表示 |
scheduled | 调度器已判定任务的依赖满足,应当运行(IntermediateTIState.SCHEDULED) |
queued | 任务已分配给某个 Executor,正在等待 worker(IntermediateTIState.QUEUED) |
running | 任务正在 worker 上运行(或在 local/synchronous executor 上运行) |
success | 任务运行完毕且无错误(终态TerminalTIState.SUCCESS) |
restarting | 任务在运行时被外部请求重启(例如运行时被 clear)(IntermediateTIState.RESTARTING) |
failed | 任务执行期间出错而运行失败(终态TerminalTIState.FAILED) |
skipped | 任务因分支(branching)、LatestOnly 等机制被跳过(终态TerminalTIState.SKIPPED) |
upstream_failed | 某个上游任务失败,且触发规则(Trigger Rule)要求等待它(终态TerminalTIState.UPSTREAM_FAILED) |
up_for_retry | 任务失败,但还有重试次数,将被重新调度(IntermediateTIState.UP_FOR_RETRY) |
up_for_reschedule | 任务是一个处于reschedule模式的 Sensor(IntermediateTIState.UP_FOR_RESCHEDULE) |
deferred | 任务已**延迟(deferred)**到一个触发器(trigger)等待异步事件(IntermediateTIState.DEFERRED),详见 deferring.rst |
awaiting_input | 任务是一个Human-in-the-loop(人机交互)任务,等待人工响应;由调度器管理,既不占用 worker 槽位也不占用 triggerer(IntermediateTIState.AWAITING_INPUT),详见 hitl.rst |
removed | 自运行开始后任务已从 DAG 中消失(终态TerminalTIState.REMOVED) |
在理想情况下,任务应当按如下路径流转:none→scheduled→queued→running→success(见上图生命周期示意)。
关系术语:upstream/downstream 与 previous/next 的区别
对于任意一个 Task Instance,它与其它实例之间存在两类关系:
第一类:upstream 与 downstream 任务
task1 >> task2 >> task3当 DAG 运行时,会为这些彼此互为上下游的任务创建实例,这些实例共享同一个数据区间。
第二类:同一任务在不同数据区间上的实例
这些实例来自同一 DAG 的其他运行,Airflow 称之为previous(上一个)与next(下一个)——这与 upstream/downstream 是完全不同的关系维度!
注意:一些较老版本的 Airflow 文档可能仍用 "previous" 表示 "upstream"。如果发现这类用法,可以协助社区修正文档。
超时控制(Timeouts)
如果希望给任务设置最大运行时长,可以设置任务的execution_timeout属性,值为一个datetime.timedelta。该设置适用于所有Airflow 任务(包括 Sensor)。execution_timeout控制的是每一次执行所允许的最大时间;一旦超时,任务会超时并抛出AirflowTaskTimeout异常。
此外,Sensor 还有一个额外的timeout参数,仅对reschedule模式下的 Sensor 有效。timeout控制的是"Sensor 最终成功所允许的最大总时间";一旦超过,将抛出AirflowSensorTimeout,Sensor立即失败且不再重试。
综合示例:SFTPSensor 的超时与重试组合
以下SFTPSensor示例清晰地说明了这两层超时的配合(代码取自原文档,参数逐条解释如下):
sensor = SFTPSensor( task_id="sensor", path="/root/test", execution_timeout=timedelta(seconds=60), timeout=3600, retries=2, mode="reschedule", )该 Sensor 处于reschedule模式(即周期性地执行、重新调度,直到成功为止),各参数行为如下:
- 每次探测(poke)SFTP 服务器最多允许 60 秒(
execution_timeout)。 - 如果一次探测超过 60 秒,抛出
AirflowTaskTimeout,此时允许重试,最多重试 2 次(retries)。 - 从第一次执行开始,到最终成功(即文件
root/test出现)为止,总共最多允许 3600 秒(timeout)。如果 3600 秒内文件始终未出现,抛出AirflowSensorTimeout,此时不再重试。 - 如果在 3600 秒窗口内因其他原因失败(如网络中断),仍可最多重试 2 次(
retries)。重试不会重置timeout——它依然总共只有 3600 秒用于成功。
SLAs 的历史变更
Airflow 2 中基于 SLA 的功能在Airflow 3.0 中已被移除,并在 Airflow 3.1 中被Deadlines Alerts(截止时间告警)取代。当前项目中如需使用类似能力,请参考 deadline-alerts.rst。
特殊异常(Special Exceptions):从任务代码内部控制状态
如果希望从自定义 Task/Operator 代码内部控制任务的最终状态,Airflow 提供了两个可主动抛出的特殊异常(定义于 exceptions.py,导出名AirflowSkipException、AirflowFailException):
AirflowSkipException:将当前任务标记为skipped(跳过)。AirflowFailException:将当前任务标记为failed(失败),并且忽略剩余的任何重试次数。
这两个异常非常适合代码对自身环境有额外认知、希望"更快地失败/跳过"的场景,例如:
- 已知本次没有可用数据时,直接跳过任务;
- 检测到 API Key 无效时快速失败(因为重试不会修复密钥问题,没必要浪费重试次数)。
重试策略(Retry Policies):按异常类型精细化控制重试
默认情况下,Airflow 以固定的次数和固定的延迟重试失败任务,而不区分错误类型。Retry Policy(重试策略)允许你以任务或 Operator 上的一个参数,按异常类型配置逐异常(per-exception)的重试行为,无需修改任务代码。
定义一个策略并应用到任务
策略由"将异常类型映射到动作的规则"组成,完整可运行的示例见 example_retry_policy.py。策略定义与使用如下:
# 定义策略 from airflow.sdk import DAG, ExceptionRetryPolicy, RetryAction, RetryRule, task API_RETRY_POLICY = ExceptionRetryPolicy( rules=[ RetryRule( exception="requests.exceptions.HTTPError", action=RetryAction.RETRY, retry_delay=timedelta(minutes=5), reason="Rate limit, backing off", ), RetryRule( exception="google.auth.exceptions.RefreshError", action=RetryAction.FAIL, reason="Auth failure, not retryable", ), RetryRule( exception=ConnectionError, action=RetryAction.RETRY, retry_delay=timedelta(seconds=30), ), ], ) # 应用到任务 with DAG( dag_id="example_retry_policy", schedule=None, catchup=False, tags=["example", "retry_policy"], ): @task(retries=5, retry_delay=timedelta(minutes=1), retry_policy=API_RETRY_POLICY) def call_external_api(): import requests response = requests.get("https://api.example.com/data") response.raise_for_status() return response.json() call_external_api()工作原理:策略在 worker 进程中求值
策略运行在任务 worker 进程中(绝不在调度器中),位于"捕获异常"与"决定任务下一状态"之间。每次策略决策都会记录在任务日志中,格式为:
Retry policy decision action=<action> reason=<reason>当任务失败时,策略对异常求值并返回三种动作之一(对应源码 retry_policy.py 中的RetryAction枚举):
- RETRY:重试任务,可选择自定义延迟以覆盖
retry_delay。重试仍受任务的retries计数约束——策略可以让任务更早失败,但不能超出配置的最大次数。 - FAIL:立即失败,跳过剩余的重试。
- DEFAULT:回退到标准重试逻辑(即
retries次数与retry_delay)。
规则按顺序求值,第一个匹配的规则生效;若没有规则匹配,策略返回DEFAULT(即标准重试行为)。
异常匹配规则
异常类型既可以指定为 Python 类,也可以指定为点分导入路径字符串(例如"requests.exceptions.HTTPError")。字符串路径会在DAG 解析时进行校验:不含点号的路径会立即抛出ValueError,无法解析的路径会产生警告(见 retry_policy.py 中RetryRule.__post_init__的校验逻辑)。
默认情况下,规则使用isinstance匹配,因此针对OSError的规则也会匹配ConnectionError(其子类)。可以通过match_subclasses=False改为精确类型匹配:
RetryRule(exception=OSError, match_subclasses=False) # 仅匹配 OSError 本身,不含子类另外,RetryRule.exception还支持传入列表,此时该规则对列表中任一异常匹配即生效(源码 retry_policy.py 中的RetryRule文档字符串有明确示例)。
与既有参数的组合行为
当任务设置了retry_policy时,各既有参数的行为如下:
| 参数 | 设置retry_policy后的行为 |
|---|---|
retries | 仍然是最大重试次数。策略可以更早失败,但不能超过该上限。 |
retry_delay/retry_exponential_backoff/max_retry_delay | 当策略返回 DEFAULT 或RetryDecision.retry_delay为 None 时使用。 |
on_retry_callback | 在所有重试(包括策略驱动的重试)上都会触发。 |
AirflowFailException | 始终具有最高优先级。该异常从不经过策略求值(直接失败)。 |
注:
AirflowSensorTimeout同样总是使任务立即失败,重试策略也不会被调用(见 retry_policy.py 中evaluate的文档说明)。
复用策略:跨 DAG 共享
策略可以定义一次,然后通过default_args或共享模块在多个 DAG 间复用:
# policies.py -- 在任何 DAG 中导入 STANDARD_RETRY_POLICY = ExceptionRetryPolicy( rules=[ RetryRule(exception="requests.exceptions.HTTPError", action=RetryAction.FAIL), RetryRule(exception=ConnectionError, retry_delay=timedelta(seconds=10)), ], )动态任务映射(Mapped Tasks)下的策略
策略与动态任务映射(dynamic task mapping)通过.partial()配合使用。策略按每个映射出的任务实例分别生效——例如 10 个映射实例中的第 2 个命中 FAIL,其余 9 个仍独立继续:
@task.partial(retry_policy=my_policy).expand(input=[1, 2, 3]) def my_mapped_task(input): ...策略在任务级别通过.partial()设置,所有映射实例共享同一个策略;.expand()不支持按索引变化。但策略的evaluate()方法会接收到异常、try_number和完整 context,因此如果需要,可以在策略内部实现按索引(per-index)的分支逻辑。
自定义重试策略:子类化RetryPolicy
对于高级场景,可以子类化airflow.sdk.definitions.retry_policy.RetryPolicy并实现evaluate()。当需要检查异常属性(状态码、响应头、响应体)而声明式的ExceptionRetryPolicy规则无法覆盖时,子类化是正确的选择。evaluate()方法接收异常、尝试次数、最大尝试次数以及完整的 Airflow context(dag_run、params等),并返回一个RetryDecision(该数据类及便捷构造方法fail/retry/default见 retry_policy.py)。
模式一:按 HTTP 状态码路由重试决策(并响应 429 的Retry-After头):
from datetime import timedelta import requests from airflow.sdk import RetryDecision, RetryPolicy class HTTPStatusRetryPolicy(RetryPolicy): """按 HTTP 状态码路由重试决策,429 时遵循 Retry-After 头。""" def evaluate(self, exception, try_number, max_tries, context=None): if isinstance(exception, requests.HTTPError) and exception.response is not None: status = exception.response.status_code if status == 429: # 被限流 -- 遵循 Retry-After 头 retry_after = int(exception.response.headers.get("Retry-After", 60)) return RetryDecision.retry(retry_delay=timedelta(seconds=retry_after)) if 500 <= status < 600: # 服务器错误 -- 值得重试 return RetryDecision.retry() if 400 <= status < 500: # 客户端错误 -- 不可重试 return RetryDecision.fail(reason=f"HTTP {status}") return RetryDecision.default()模式二:利用运行 context 决策(例如回填时快速失败):
from airflow.sdk import RetryDecision, RetryPolicy class BackfillAwareRetryPolicy(RetryPolicy): """回填期间快速失败,让历史错误立即暴露出来。""" def evaluate(self, exception, try_number, max_tries, context=None): if context and context["dag_run"].run_type == "backfill": return RetryDecision.fail(reason="Backfill run -- not retrying") return RetryDecision.default()提示:自定义策略需要支持 DAG 序列化(
serialize()/deserialize()),内置的ExceptionRetryPolicy已提供相应实现,子类化时需自行实现(参见 retry_policy.py 的接口定义)。
Task Instance 心跳超时(Heartbeat Timeout):处理"僵尸任务"
没有哪个系统是完美运行的,Task Instance 偶尔也会"死亡"。Task Instance 可能停留在running状态,但其关联的 job 已经不活跃了(例如该 TaskInstance 的 worker 内存耗尽被杀)。这类任务在旧版中被称为"僵尸任务(zombie tasks)"。Airflow 会定期发现它们、进行清理,并将 TaskInstance 标记为失败;如果还有可用重试次数,则触发重试。
TaskInstance 心跳超时的常见诱因包括:
- Airflow worker 内存耗尽被 OOMKilled;
- worker 未通过存活探针(liveness probe),导致系统(如 Kubernetes)重启了 worker;
- 系统(如 Kubernetes)缩容,将 worker 从一个节点迁移到另一个节点。
本地复现心跳超时(供开发/测试)
如果需要在本地产复现心跳超时,可按下述步骤操作:
步骤 1:设置环境变量(也可以等价地修改airflow.cfg中的对应配置项):
export AIRFLOW__SCHEDULER__TASK_INSTANCE_HEARTBEAT_SEC=600 export AIRFLOW__SCHEDULER__TASK_INSTANCE_HEARTBEAT_TIMEOUT=2 export AIRFLOW__SCHEDULER__TASK_INSTANCE_HEARTBEAT_TIMEOUT_DETECTION_INTERVAL=5三个变量分别含义为:TASK_INSTANCE_HEARTBEAT_SEC为任务心跳间隔(这里设成 600 秒,模拟长时间无心跳);TASK_INSTANCE_HEARTBEAT_TIMEOUT为心跳超时阈值(设为 2 秒,让检测快速触发);TASK_INSTANCE_HEARTBEAT_TIMEOUT_DETECTION_INTERVAL为检测间隔(每 5 秒扫描一次)。
步骤 2:准备一个耗时约 10 分钟的任务,例如:
from airflow.sdk import dag from airflow.providers.standard.operators.bash import BashOperator from datetime import datetime @dag(start_date=datetime(2021, 1, 1), schedule="@once", catchup=False) def sleep_dag(): t1 = BashOperator( task_id="sleep_10_minutes", bash_command="sleep 600", ) sleep_dag()运行上述 DAG 并等待一段时间后,TaskInstance 将在约<task_instance_heartbeat_timeout>秒后被标记为失败。
Executor 级配置(executor_config):按任务定制执行环境
部分 Executor 允许可选的按任务配置(per-task configuration)。例如KubernetesExecutor允许为某个任务指定运行它的镜像。这是通过 Task 或 Operator 的executor_config参数实现的。以下示例为将在KubernetesExecutor上运行的任务设置 Docker 镜像:
MyOperator(..., executor_config={ "KubernetesExecutor": {"image": "myCustomDockerImage"} } )executor_config中可设置的项因 Executor 而异,请阅读各个 Executor 的专属文档(见 executor/index.rst)了解可配置项。
小结与最佳实践
本文围绕 Airflow Task 的完整生命周期梳理了以下核心知识点:
- Task 的三种形式:Operator(模板)、Sensor(等待外部事件)、
@task(TaskFlow 函数),底层均为BaseOperator子类; - 依赖声明:推荐使用
>>/<<位运算操作符,依赖关系基于"直接上游/下游"; - Task Instance 状态机:从
none到scheduled→queued→running→success的理想路径,以及up_for_retry、deferred、awaiting_input、removed等全部 15 种状态(源码依据见 state.py); - 双层超时:
execution_timeout(单次执行上限,超时抛AirflowTaskTimeout可重试)与 Sensor 的timeout(总成功时限,超时抛AirflowSensorTimeout立即失败); - 特殊异常:
AirflowSkipException(跳过)与AirflowFailException(失败且不重试); - Retry Policy:以
ExceptionRetryPolicy+RetryRule声明式配置按异常重试,支持 RETRY / FAIL / DEFAULT 三动作、isinstance/精确类型匹配、字符串导入路径校验、映射任务支持及自定义RetryPolicy子类(示例见 example_retry_policy.py); - 心跳超时:僵尸任务的检测与清理机制及其本地复现方法;
- executor_config:按任务定制执行环境(如 Kubernetes 镜像)。
在实际编写 DAG 时,建议优先采用位运算操作符声明依赖、为所有外部调用类任务显式设置execution_timeout、把可重试/不可重试的异常区分写入 Retry Policy,从而让任务具备可预测、可自愈、可运维的高质量行为。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考