Airflow 可延后执行的 Operator 与 Trigger

Airflow 可延后执行的 Operator 与 Trigger

标准 Operator 和 Sensor 在运行期间始终占用一个工作槽位,即使只是在等待。例如只有 100 个槽位,而 100 个 DAG 都停在正在运行但空闲的 Sensor 上,就无法运行其他任务,尽管整个集群几乎没有实际工作。Sensor 的 reschedule 模式通过定期执行缓解了这一问题,但只能根据时间恢复,不能灵活响应其他条件。

可延后执行的 Operator(Deferrable Operator)可以在只需等待时挂起自己,释放工作节点执行其他任务。任务延后后,执行转移到 triggerer,由指定的 Trigger 轮询或等待。条件满足时,Trigger 发出信号,让 Operator 恢复。延后期间任务不占工作槽位;默认也不占 pool 槽位,如果希望继续占用,可以调整相应 pool。

Trigger 是短小的异步 Python 代码,可在同一个 Python 进程中高效并存,由 Airflow 的 triggerer 组件运行。

Airflow 3.2 还支持原生 Python 异步任务,它们在一个工作槽位中并发执行 I/O。延后 Operator 在等待外部事件时释放工作槽位;异步任务则保持任务进程运行,通过共享事件循环复用操作。选择方法见 Deferred vs Async Operators。如需与使用任务状态存储、可在崩溃后恢复的同步 Operator 比较,请参阅 Airflow 的 Resumable Tasks 文档。

整个过程如下:

  1. 正在运行的任务实例需要等待其他操作或条件,于是绑定 Trigger 并延后自身,释放工作节点。
  2. Airflow 注册新的 Trigger 实例,由某个 triggerer 进程接手。
  3. Trigger 运行直至产生事件,调度器重新调度其来源任务。
  4. 调度器将任务加入队列,在工作节点恢复执行。

DAG 作者可以使用现成的可延后 Operator,也可以自己编写;自行实现时需要遵守设计约束。

使用可延后 Operator

使用 Airflow 提供的可延后 Operator,例如 TimeSensor,只需:

  • 除常规 scheduler 外,至少运行一个 triggerer 进程。
  • 在 DAG 中使用支持延后的 Operator 或 Sensor。

Airflow 自动处理延后过程。

人机协作(Human-in-the-loop)Operator 不使用延后机制或 triggerer,而是在调度器管理的 awaiting_input 状态中等待。因此仅等待人工输入、不使用可延后 Operator 的部署,不需要运行 triggerer。

升级现有 DAG 时,可以换用 API 兼容的 Sensor 变体,通常无需其他修改。

不能在自定义 PythonOperator 或 TaskFlow Python 函数内部使用延后功能;它只适用于传统的类式 Operator。

编写可延后 Operator

需要考虑以下几点:

  • 必须通过 Trigger 延后自己,可以使用 Airflow 内置 Trigger,也可以自行实现。
  • 延后时 Operator 停止运行并从工作节点移除,状态不会自动保存。可以通过指定恢复方法及传递关键字参数保存所需状态。
  • 可以延后多次,在主要工作前后延后,或仅在某些条件下延后,例如外部系统无法立即回答时;控制权在实现者手中。
  • 任何 Operator 都可以延后,不需要类级特殊标记,也不限于 Sensor。
  • 同时支持延后和非延后模式时,建议在 __init__ 中加入 deferrable: bool = conf.getboolean("operators", "default_deferrable", fallback=False),据此选择执行模式。支持该约定的 Operator 和 Sensor 可以通过配置节 [operators] 中的 default_deferrable 统一设置默认值。原文说明段将配置节写成了单数 operator,下方代码使用复数 operators。

下面的 Sensor 支持两种模式:

import time
from datetime import timedelta
from typing import Any

from airflow.sdk import BaseSensorOperator, Context, conf
from airflow.providers.standard.triggers.temporal import TimeDeltaTrigger


class WaitOneHourSensor(BaseSensorOperator):
    def __init__(
        self, deferrable: bool = conf.getboolean("operators", "default_deferrable", fallback=False), **kwargs
    ) -> None:
        super().__init__(**kwargs)
        self.deferrable = deferrable

    def execute(self, context: Context) -> None:
        if self.deferrable:
            self.defer(
                trigger=TimeDeltaTrigger(timedelta(hours=1)),
                method_name="execute_complete",
            )
        else:
            time.sleep(3600)

    def execute_complete(
        self,
        context: Context,
        event: dict[str, Any] | None = None,
    ) -> None:
        # We have no more work to do here. Mark as complete.
        return

编写 Trigger

Trigger 继承 BaseTrigger,实现三个方法:

  • __init__:接收实例化它的 Operator 传入的参数。从 2.10.0 起,任务可以直接从预定义 Trigger 开始执行;使用此功能时所有参数必须可序列化。
  • run:异步方法,以异步生成器方式运行逻辑,并产生一个或多个 TriggerEvent。
  • serialize:返回用于重建 Trigger 的元组,包含类导入路径和传给 __init__ 的关键字参数。

还有两个可选生命周期钩子:

  • cleanup:无论 run 因成功、超时、triggerer 关闭或用户终止而退出,都会调用。适合释放本地资源,如连接和临时文件。
  • on_kill:仅在用户明确终止延后任务时调用,例如标记失败、清除或标记成功。可以在这里取消不应继续运行的外部工作,例如 BigQuery 作业或 Databricks 运行。它不会因 triggerer 重启或重新分配而调用,因此不会在滚动部署时误取消正在执行的外部工作。

以下是非常简化的 DateTimeTrigger 结构:

import asyncio

from airflow.triggers.base import BaseTrigger, TriggerEvent
from airflow.sdk.timezone import utcnow


class DateTimeTrigger(BaseTrigger):
    def __init__(self, moment):
        super().__init__()
        self.moment = moment

    def serialize(self):
        return ("airflow.providers.standard.triggers.temporal.DateTimeTrigger", {"moment": self.moment})

    async def run(self):
        while self.moment > utcnow():
            await asyncio.sleep(1)
        yield TriggerEvent(self.moment)

这个示例体现了几个要点:

  • __init__ 与 serialize 成对设计。Operator 提交延后请求时先实例化 Trigger,再序列化,最后在负责它的 triggerer 上重建。
  • run 必须是 async def。使用 asyncio.sleep,不能用会阻塞整个进程的普通 time.sleep。
  • 事件载荷中包含 self.moment,便于同一个 Trigger 在多台主机冗余运行时去重。

只要遵守约束,Trigger 可以很简单,也可以很复杂;它支持高可用,并自动分配到各 triggerer 主机。应尽量避免持久状态,所需信息都通过 __init__ 提供,以便自由序列化和迁移。

不熟悉 Python 异步编程时,编写 run() 尤其要谨慎。阻塞操作如果没有正确 await,可能阻塞整个进程。Airflow 会尝试检测并在 triggerer 日志中警告。开发时可设置 PYTHONASYNCIODEBUG=1,启用更多检查。文件系统调用也要小心,网络文件系统可能产生阻塞。

设计约束如下:

  • run 必须使用 asyncio 异步执行,遇到阻塞操作时正确等待。
  • 必须用 yield 产生 TriggerEvent,不能直接返回它。在产生第一个事件前返回,或抛出异常,都会使等待它的任务实例失败。
  • 应假定 Trigger 可能运行多次。例如网络分区时,Airflow 可能在另一台机器重新启动同一 Trigger。因此要考虑副作用,不宜随意在 Trigger 中插入数据库记录。
  • 若设计成产生多个事件,每个事件必须包含可去重载荷;原文说明当前尚不支持多事件用途。只产生一个事件且无需向 Operator 返回信息时,载荷可为 None。
  • Trigger 可能突然从一个服务迁移到另一个服务,例如部署或网络分区时。需要释放资源可实现 cleanup;仅用户主动取消时才要终止外部工作,则实现 on_kill。
  • 修改 Trigger 后,需要重启 triggerer 才能生效。
  • Trigger 不能来自 DAG bundle;放在 sys.path 的其他可导入位置即可,因为 triggerer 运行 Trigger 时不初始化 bundle。

原文说明,目前 Trigger 用于恢复延后任务,因此只处理第一个事件,任务随后恢复;未来计划支持由 Trigger 启动 DAG,届时多事件 Trigger 会更有用。

Trigger 中的敏感信息

从 Airflow 2.9.0 起,Trigger 的关键字参数会先序列化、加密,再写入数据库。因此传入的敏感信息以密文存储,读取时解密。

发起延后

在 Operator 任意位置调用 self.defer(trigger, method_name, kwargs, timeout),会抛出 Airflow 专用异常。参数如下:

  • trigger:要等待的 Trigger 实例,会序列化存入数据库。
  • method_name:恢复时调用的 Operator 方法名。
  • kwargs:传给恢复方法的额外关键字参数,可选,默认 {}。
  • timeout:可选 timedelta,超过此时长则本次延后和任务实例失败。默认 None,表示不超时。

Sensor 发起延后的简单示例:

from __future__ import annotations

from datetime import timedelta
from typing import Any

from airflow.sdk import BaseSensorOperator, Context
from airflow.providers.standard.triggers.temporal import TimeDeltaTrigger


class WaitOneHourSensor(BaseSensorOperator):
    def execute(self, context: Context) -> None:
        self.defer(trigger=TimeDeltaTrigger(timedelta(hours=1)), method_name="execute_complete")

    def execute_complete(self, context: Context, event: dict[str, Any] | None = None) -> None:
        # We have no more work to do here. Mark as complete.
        return

调用延后后,Operator 在该处停止,从当前工作节点移除。局部变量和运行时写入 self 的属性都不会保留。恢复时得到的是新的 Operator 实例,只能通过 method_name 和 kwargs 把旧实例状态传给新实例。

恢复时,Airflow 还会向指定方法传入 context 和 event。event 包含触发恢复的事件载荷,可能是状态码、结果 URL,也可能只是时间戳。无论是否使用,恢复方法都必须接受这两个关键字参数。

首次 execute() 或后续由 method_name 指定的方法正常返回时,Operator 被视为完成。上面的 WaitOneHourSensor 只是 Trigger 的薄封装:延后等待事件,再进入立即返回的恢复方法,将 Sensor 标记成功。

self.defer 抛出 TaskDeferred,因此即使在 execute() 深层调用中也能生效。也可以手动抛出接受相同参数的 TaskDeferred。

Operator 的 execution_timeout 按总运行时间计算,不是各次延后之间单段执行时间。因此即使刚恢复几秒,也可能在延后状态中或恢复后因总时限到达而失败。

多次延后

假设 Operator 遍历长度不定的项目列表,并为每项延后处理,例如向数据库提交多条查询,或处理多个文件。

可以把 method_name 设为 execute,只使用一个入口,但该方法必须接受可选的 event 参数。原文示意代码如下,其中省略号表示待实现的业务逻辑:

import asyncio

from airflow.sdk import BaseOperator
from airflow.triggers.base import BaseTrigger, TriggerEvent


class MyItemTrigger(BaseTrigger):
    def __init__(self, item):
        super().__init__()
        self.item = item

    def serialize(self):
        return (self.__class__.__module__ + "." + self.__class__.__name__, {"item": self.item})

    async def run(self):
        result = None
        try:
            # Somehow process the item to calculate the result
            ...
            yield TriggerEvent({"result": result})
        except Exception as e:
            yield TriggerEvent({"error": str(e)})


class MyItemsOperator(BaseOperator):
    def __init__(self, items, **kwargs):
        super().__init__(**kwargs)
        self.items = items

    def execute(self, context, current_item_index=0, event=None):
        last_result = None
        if event is not None:
            # execute method was deferred
            if "error" in event:
                raise Exception(event["error"])
            last_result = event["result"]
            current_item_index += 1

        try:
            current_item = self.items[current_item_index]
        except IndexError:
            return last_result

        self.defer(
            trigger=MyItemTrigger(item),
            method_name="execute",  # The trigger will call this same method again
            kwargs={"current_item_index": current_item_index},
        )

此原始示例的 MyItemTrigger(item) 使用了未定义的 item;依据前文赋值,应核对是否应为 current_item 后再使用。这里保留原代码,避免把未经运行的修正伪称为已验证示例。

从任务开始时就延后

此功能在 2.10.0 引入。若希望任务直接进入 triggerer、不先进入工作节点,可将类属性 start_from_trigger 设为 True,并添加 StartTriggerArgs 类型的 start_trigger_args:

  • trigger_cls:Trigger 类的可导入路径。
  • trigger_kwargs:初始化 Trigger 的参数,全部必须可由 Airflow 序列化,这是主要限制。
  • next_method:恢复时调用的 Operator 方法。
  • next_kwargs:传给恢复方法的额外参数。
  • timeout:可选超时时长,默认 None;超时会让延后和任务实例失败。

Sensor 需要将 TimeDeltaTrigger 的路径作为 trigger_cls:

from __future__ import annotations

from datetime import timedelta
from typing import Any

from airflow.sdk import BaseSensorOperator, Context, StartTriggerArgs


class WaitOneHourSensor(BaseSensorOperator):
    start_trigger_args = StartTriggerArgs(
        trigger_cls="airflow.providers.standard.triggers.temporal.TimeDeltaTrigger",
        trigger_kwargs={"moment": timedelta(hours=1)},
        next_method="execute_complete",
        next_kwargs=None,
        timeout=None,
    )
    start_from_trigger = True

    def execute_complete(self, context: Context, event: dict[str, Any] | None = None) -> None:
        # We have no more work to do here. Mark as complete.
        return

start_from_trigger 和 trigger_kwargs 也可以在实例级修改,增加配置灵活性:

from __future__ import annotations

from datetime import timedelta
from typing import Any

from airflow.sdk import BaseSensorOperator, Context, StartTriggerArgs


class WaitHoursSensor(BaseSensorOperator):
    start_trigger_args = StartTriggerArgs(
        trigger_cls="airflow.providers.standard.triggers.temporal.TimeDeltaTrigger",
        trigger_kwargs={"moment": timedelta(hours=1)},
        next_method="execute_complete",
        next_kwargs=None,
        timeout=None,
    )
    start_from_trigger = True

    def __init__(self, *args: list[Any], **kwargs: dict[str, Any]) -> None:
        super().__init__(*args, **kwargs)
        self.start_trigger_args.trigger_kwargs = {"hours": 2}
        self.start_from_trigger = True

    def execute_complete(self, context: Context, event: dict[str, Any] | None = None) -> None:
        # We have no more work to do here. Mark as complete.
        return

映射任务的初始化发生在调度器提交给执行器之后,因此该功能对动态任务映射的支持有限,与常规用法不同。需要在 __init__ 中定义 start_from_trigger 和/或 trigger_kwargs,不必同时定义,但参数名称必须完全一致。例如使用 t_kwargs 再赋给 self.start_trigger_args.trigger_kwargs 不会生效。

映射 start_from_trigger=True 的任务时,整个 __init__ 被跳过。调度器使用 partial 和 expand 提供的这两个参数,未提供时回退到类属性,决定把任务交给执行器还是 triggerer、以及如何提交。此阶段不会解析 XCom 值。

Trigger 完成后,任务可能回到工作节点执行 next_method,也可能直接结束。若回到工作节点,__init__ 参数仍会在执行恢复方法前生效,但不会影响已完成的 Trigger 执行。

from __future__ import annotations

from datetime import timedelta
from typing import Any

from airflow.sdk import BaseSensorOperator, Context, StartTriggerArgs


class WaitHoursSensor(BaseSensorOperator):
    start_trigger_args = StartTriggerArgs(
        trigger_cls="airflow.providers.standard.triggers.temporal.TimeDeltaTrigger",
        trigger_kwargs={"moment": timedelta(hours=1)},
        next_method="execute_complete",
        next_kwargs=None,
        timeout=None,
    )
    start_from_trigger = True

    def __init__(
        self,
        *args: list[Any],
        trigger_kwargs: dict[str, Any] | None,
        start_from_trigger: bool,
        **kwargs: dict[str, Any],
    ) -> None:
        # This whole method will be skipped during dynamic task mapping.

        super().__init__(*args, **kwargs)
        self.start_trigger_args.trigger_kwargs = trigger_kwargs
        self.start_from_trigger = start_from_trigger

    def execute_complete(self, context: Context, event: dict[str, Any] | None = None) -> None:
        # We have no more work to do here. Mark as complete.
        return

下面展开为两个任务,hours 分别为 1 和 2:

WaitHoursSensor.partial(task_id="wait_for_n_hours", start_from_trigger=True).expand(
    trigger_kwargs=[{"hours": 1}, {"hours": 2}]
)

这些片段按原文保留;其中 TimeDeltaTrigger 的参数名在 moment 与 hours 之间变化,使用时需与安装的 Standard provider 版本签名核对。

直接从 Trigger 结束延后任务

此功能在 2.10.0 引入。设置实例属性 end_from_trigger,可让任务直接在 triggerer 结束,无需重新进入工作节点,节省启动资源。

Trigger 可以把执行交回工作节点,也可以直接结束任务。直接结束时,method_name 不再起作用,可设为 None;否则应指定恢复方法。

class WaitFiveHourSensorAsync(BaseSensorOperator):
    # this sensor always exits from trigger.
    def __init__(self, **kwargs) -> None:
        super().__init__(**kwargs)
        self.end_from_trigger = True

    def execute(self, context: Context) -> NoReturn:
        self.defer(
            method_name=None,
            trigger=WaitFiveHourTrigger(duration=timedelta(hours=5), end_from_trigger=self.end_from_trigger),
        )

TaskSuccessEvent 和 TaskFailureEvent 可直接结束任务实例,按 task_instance_state 设置状态,并在适用时推送 XCom。示例如下:

class WaitFiveHourTrigger(BaseTrigger):
    def __init__(self, duration: timedelta, *, end_from_trigger: bool = False):
        super().__init__()
        self.duration = duration
        self.end_from_trigger = end_from_trigger

    def serialize(self) -> tuple[str, dict[str, Any]]:
        return (
            "your_module.WaitFiveHourTrigger",
            {"duration": self.duration, "end_from_trigger": self.end_from_trigger},
        )

    async def run(self) -> AsyncIterator[TriggerEvent]:
        await asyncio.sleep(self.duration.total_seconds())
        if self.end_from_trigger:
            yield TaskSuccessEvent()
        else:
            yield TriggerEvent({"duration": self.duration})

end_from_trigger=True 时,Trigger 产生 TaskSuccessEvent 直接结束任务;否则使用 Operator 指定的方法恢复。

原文限制:可延后 Operator 集成了 listeners 时,不支持直接从 Trigger 结束。当前同时设置 end_from_trigger=True 与 listeners,会在 DAG 解析时抛出异常。编写自定义 Trigger 时,若插件添加了 listeners,不应直接结束任务。如果作者把这一属性改成别的名称,DAG 解析可能不再报错,但依赖该任务的 listeners 仍无法工作。原文说明此限制将在后续版本解决。

高可用

Trigger 支持高可用架构。可在多台主机运行多个 triggerer,它们与 scheduler 类似,通过锁和高可用机制自动协作。

根据工作量,一台 triggerer 主机可容纳数百到数万个 Trigger。默认每个进程最多尝试同时运行 1,000 个,可用 --capacity 调整。总需求超过所有 triggerer 总容量时,部分 Trigger 要等其他 Trigger 完成才会运行。

Airflow 尽量保证一个 Trigger 同时只在一处运行,并监测 triggerer 心跳。进程死亡,或与数据库所在网络断开时,其 Trigger 自动重新调度到其他主机。重新调度前会等待 2.1 * triggerer.job_heartbeat_sec 秒,看原主机是否恢复。

因此同一 Trigger 偶尔可能同时在多处运行。这是设计中考虑的正常情况。Airflow 会对重复事件去重,对 Operator 透明。

每增加一个 triggerer,就会增加一条到数据库的持久连接。

平衡高可用 triggerer 的负载

此功能在 3.2.0 引入。每轮只选取 [triggerer] max_trigger_to_select_per_loop 个 Trigger,避免其他进程分配不到工作。建议该值显著低于 [triggerer] capacity。默认值分别为 50 和 1000。

根据基准测试,默认设置下,两个 triggerer 仍可在一秒内领取 1,000 个 Trigger,且负载近乎均匀。

可以创建大量 Trigger,例如触发包含很多可延后任务的 DAG,观察各进程负载分布及全部领取完成所需时间,为自己的部署选择合适数值。

按 Trigger 控制主机分配

此功能在 3.2.0 引入。在一些场景中,希望将 Trigger 限定到某组主机,例如多租户系统为各团队独立运行 triggerer,或者不同主机组具有不同云权限、负责不同操作。

使用 Multi-Team 模式时,--team-name 为任务创建、事件驱动及回调三类 Trigger 提供原生团队范围分配。下文的 --queues 是较早的队列机制,也可按需要与 --team-name 组合。

启用队列分配:

  1. 将 [triggerer] queues_enabled 设为 true,让任务延后时把所分配的任务队列传给新注册 Trigger。
  2. 在相应 triggerer 的启动命令中添加 --queues=<逗号分隔的任务队列名称>,使它只领取来自这些队列的任务所创建的 Trigger。

例如两台主机 X、Y 使用以下命令:

# triggerer "X" startup command
airflow triggerer --queues=alice,bob
# triggerer "Y" startup command
airflow triggerer --queues=test_q

X 仅运行来自任务队列 alice 或 bob 的 Trigger;Y 仅运行来自 test_q 的 Trigger。

队列分配限制

此功能仅适用于具有任务 queue 概念的执行器,例如 CeleryExecutor;目前也仅支持由任务的 defer 方法创建的 Trigger。

Trigger 类型 支持队列 启用 queues_enabled 后如何分配
任务创建的 Trigger 是 其 --queues 包含该任务队列的任意 triggerer
事件驱动 Trigger 否 未设置 --queues 的任意 triggerer
异步回调创建的 Trigger 否 未设置 --queues 的任意 triggerer

如果任务 Trigger 使用队列,同时还使用事件驱动或回调 Trigger,必须至少运行一个不带 --queues 的 triggerer,让后两类仍能被执行。

要启用 Trigger 队列,至少一个 triggerer 启动命令必须设置 --queues,各进程可配置不同队列。若名称不对应任何任务队列,它不会运行任何 Trigger。启用该机制后,不带 --queues 的进程只消费事件驱动和回调 Trigger。原文此处有单数 --queue 的笔误,命令示例使用 --queues。

假设所有任务队列为 team_A 或 team_B,以下配置覆盖全部 Trigger 类型:

# triggerer "A" startup command, only consumes triggers registered by tasks in queue "team_A"
airflow triggerer --queues=team_A
# triggerer "B" startup command, only consumes triggers registered by tasks in queue "team_B"
airflow triggerer --queues=team_B
# triggerer "C" startup command, consumes only event-based triggers and callback-based triggers.
airflow triggerer

Sensor 中 mode='reschedule' 与 deferrable=True 的区别

Sensor 等待特定条件满足后才允许下游任务继续。mode='reschedule' 是 BaseSensorOperator 的内置参数,让条件未满足的 Sensor 重新调度。deferrable=True 则是部分 Operator 采用的约定,并不是所有 Operator 的内置参数或模式;具体延后行为取决于实现。

mode='reschedule' deferrable=True
不断重新调度,直到条件满足 空闲时暂停,条件变化后恢复
因重复执行而使用较多资源 等待时释放工作槽位,资源使用较少
适合随时间变化的条件,如文件创建 适合外部事件或资源,如 API 响应
内置重新调度功能 需要实现延后和处理外部变化的逻辑

原文:Deferrable Operators & Triggers。本文源代码快照取自 Apache Airflow 3.3.2。版权归 Apache Software Foundation 及贡献者所有;原文件使用 Apache License 2.0,并要求保留相关 NOTICE。原文中的示例和版本限定按上述快照保留。

原始 NOTICE:Apache Airflow,Copyright 2016–2026 The Apache Software Foundation。本作品包含 Apache Software Foundation(https://www.apache.org/)开发的软件。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容