可延迟算子

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

传感器问题

传感器占满 worker 槽位

 

  • Poke 模式:传感器全程占用 worker 槽位
  • 长时运行的传感器会饿死其他任务
使用 Airflow 构建数据流水线

Reschedule 模式

Poke 模式(默认)

from airflow.providers.standard.sensors.filesystem \
    import FileSensor

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="poke",
    poke_interval=30,
)
  • 槽位全程被占用

Reschedule 模式

from airflow.providers.standard.sensors.filesystem \
    import FileSensor

wait = FileSensor(
    task_id="wait",
    filepath="/data/report.csv",
    mode="reschedule",
    poke_interval=300,
)
  • 检查间隙释放槽位
使用 Airflow 构建数据流水线

可延迟模式

由 Triggerer 处理的延迟传感器

  • 任务交由 Triggerer 进程处理
  • Triggerer 使用 异步 I/O 同时等待数百个事件
  • 等待期间不占用任何 worker 槽位
使用 Airflow 构建数据流水线

启用可延迟

from airflow.providers.standard.sensors.filesystem import FileSensor

wait_for_data = FileSensor(
    task_id="wait_for_data",
    filepath="/data/incoming/report.csv",

deferrable=True,
poke_interval=30, timeout=3600, )

 

  • poke_intervaltimeout 仍然生效
使用 Airflow 构建数据流水线

各模式的使用时机

模式 最适用 worker 槽位
poke 短等待(< 1 分钟) 全程占用
reschedule 中等等待,无 Triggerer 检查间隙释放
deferrable 长等待,传感器多 从不占用

 

  • 若已运行 Triggerer,deferrable 几乎总是首选默认
使用 Airflow 构建数据流水线

全局可延迟设置

AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE=True

 

  • 所有兼容的算子会自动使用可延迟模式
  • 无需修改 DAG 代码
  • 如需覆盖,可在算子上设置 deferrable=False
使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...