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 文档。
整个过程如下:
- 正在运行的任务实例需要等待其他操作或条件,于是绑定 Trigger 并延后自身,释放工作节点。
- Airflow 注册新的 Trigger 实例,由某个 triggerer 进程接手。
- Trigger 运行直至产生事件,调度器重新调度其来源任务。
- 调度器将任务加入队列,在工作节点恢复执行。
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 组合。
启用队列分配:
- 将
[triggerer] queues_enabled设为true,让任务延后时把所分配的任务队列传给新注册 Trigger。 - 在相应 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/)开发的软件。











暂无评论内容