Airflow 必知必会(1)

举报
剑指南天 发表于 2026/09/12 14:48:52 2026/09/12
【摘要】 Airflow 是使用最广泛的开源工作流编排平台之一,理解 Airflow 的重要特性对于创建好的工作流至关重要。

1.概述

Airflow 是使用最广泛的开源工作流编排平台之一,理解 Airflow 的重要特性对于创建好的工作流至关重要。

2. Xcoms

默认情况下,任务实例是完全隔离的,并且可能运行在不同的机器上,所以任务之间无法传递数据。Xcoms 是 Airflow 的一个工作流中任务实例之间传递小批量元数据的功能(Cross-communications 的缩写)。

2.1 PythonOperator

2.1.1 Push 的两种方式

显式 push

@task
def push(ti=None):
    """Pushes an XCom without a specific target"""
    ti.xcom_push(key="value from pusher 1", value=value_1)

返回值 push
@task
def push_by_returning():
    """Pushes an XCom without a specific target, just by returning it"""
    return value_2

2.1.2 Pull 的三种方式

@task
def pull_value_from_bash_push(ti=None):
    # 显式取 Task 返回值
    bash_pushed_via_return_value = ti.xcom_pull(key="return_value", task_ids="push_by_returning")
    # 显式取值
    bash_manually_pushed_value = ti.xcom_pull(key="value from pusher 1", task_ids="push")

@task
def puller(pulled_value_2):
    print(pulled_value_2)

# TaskFlow API
puller(push_by_returning()) << push()

2.2 BashOperator

2.2.1 Push的两种方式

2.2.2 针对显式声明和返回值 Pull 的四种方式


2.3 Xcoms 后端可以配置外部存储,比如大 XCom(比如想传几百 KB 的 JSON)可以配对象存储,不塞元数据库

2.4 自定义 Xcoms 后端

3. 变量和宏变量

3.1 Variable

Variable 是 Airflow 里存储全局配置键值对的机制,数据存在元数据库里,所有 DAG 都能读。适合放开关、阈值、环境标识等小配置。

3.1.1 创建 Variable

3.1.1.1 Web UI 创建

3.1.1.2 CLI 创建

airflow variables set key value

3.1.1.3 代码创建

3.1.2 读取 Variable(注意:把 Variable.get()放在 task 函数内部(运行期才执行),而不是 DAG 文件顶层)

3.1.2.1 TaskFlow API 里读取 Variable

3.1.2.2 Jinja 模板里读取 Variable

{{ var.value.get('my.var', 'fallback') }} 

{{ var.json.get('my.dict.var', {'key1': 'val1'}) }}

3.1.2.3 CLI 读取 Variable

airflow variables get key

airflow variables list

3.2 宏变量

3.2.1 宏变量的分类

3.2.2 在 Operator 的 template_fields里用 Jinja

3.2.3 在 Python 函数里通过 context字典读取

3.2.4 自定义宏变量

3.2.4.1 DAG级别 user_defined_macros 和 user_defined_filters

3.2.4.2 全局 Plugin 注册

a. 配置 aiflow.cfg 文件配置 [core] plugins_folder,插件的位置

b. 编写 py 代码文件配置自定义宏

c. 将 py 代码文件放置在 plugins_folder 下,然后重启 airflow

d. Web UI 查看 plugin

e. 在 BashOperator 中使用 macros 变量

4. Jinja模板

Airflow 的模板系统完全基于 Jinja2。DAG 文件解析时只建结构,task 执行前 Airflow 会把 Operator 里标记为"模板字段"的参数交给 Jinja 渲染,把 {{ ds }}、{{ var.value.xxx }}这些占位符替换成实际值。

4.1 Jinja语法

a. 变量输出

{{ ds }}                    {# 输出变量 #}
{{ dag.dag_id }}            {# 对象属性 #}
{{ var.value.env_name }}    {# 读 Variable #}
{{ conn.smtp_default.host }} {# 读 Connection #}

b. 过滤器

{{ ds | replace("-", "_") }}
{{ logical_date | ds }}
{{ data_interval_start | ts }}
{{ my_list | join(",") }}
{{ "hello" | upper }}

c. 条件/循环

{% if params.debug %}echo "debug on"{% endif %}
{% for t in params.tables %}echo "{{ t }}"{% endfor %}

d. 函数调用

{{ macros.ds_add(ds, 7) }}
{{ macros.datetime.now() }}
{{ s3_partition(ds) }}          {# 自定义宏 #}

4.2 Operator 通过 template_fields 类属性声明哪些字段走 Jinja

5. Provider

Provider 就是 Airflow 的"外部系统对接插件包"——把连接某个系统(AWS、GCP、数据库、SMTP…)所需的 Hook、Operator、Sensor 等打包成独立 pip 包,和 Airflow Core 分开安装、分开升级。

6. Operator

Operator 是一个预定义好的“任务模板”,它封装了一类具体工作的执行逻辑,你只需声明式地传入参数,就能在 DAG 里生成一个任务。DAG 管流程,Operator 管动作,Task 是 Operator 在 DAG 里的实例。

6.1 Operator的类型

a. 动作型(Action)

PythonOperator:跑 Python 函数

BashOperator:跑 shell 命令

SQLExecuteQueryOperator/ SnowflakeOperator:执行 SQL

SimpleHttpOperator:调 HTTP 接口

DockerOperator/ KubernetesPodOperator:跑容器

b. 传输型

把数据从 A 系统搬到 B 系统,如 S3ToRedshiftOperator、GCSToBigQueryOperator

c. 传感型

特殊 Operator,一直等到某个条件满足才放行(如 FileSensor等文件出现、S3KeySensor、ExternalTaskSensor、HttpSensor)

d. 流程控制

EmptyOperator(原 DummyOperator)、BranchPythonOperator

6.2 自定义 Operator

6.2.1 在 Airflow 中自定义 Operator 核心是继承 BaseOperator,重写 __init__和 execute方法,可选 pre_execute() 和 post_execute() 最后将其放入 PYTHONPATH 即可被 DAG 导入。

6.2.2 注意事项

6.2.2.1 __init__保持轻量

调度器每个解析周期都会实例化一次 Operator,严禁在 __init__里查库、发请求、建连接,重活一律放 execute()。

6.2.2.2 模板字段直接赋值

模板字段必须在构造里直接 self.x = x,不能 self.foo = foo.lower(),也不能只透传 **kwargs不赋值,否则会被 pre-commit validate-operators-init拦下。

6.2.2.3 外部系统交互走 Hook

连外部系统(DB、API、云存储)时,通信层抽成 Hook,Operator 里只调 Hook,保证复用与连接信息管理。

6.2.3 模板集成

6.2.3.1 集成Jinja 模板

声明 template_fields 让参数支持 {{ ds }}等运行时变量,Jinja 渲染的是属性而非 __init__的 args

声明 template_ext 让文件路径自动变成文件内容模板。当同时满足以下三个条件时,Airflow 会读取文件内容作为模板,而不是把路径字符串当模板:

a. 字段名在 template_fields中

b. 该字段的值是一个字符串,且看起来像文件路径(不以空格、换行开头,且文件存在)

c. 文件的扩展名在 template_ext列表中

6.2.3.2 集成 Hook

6.3 自定义 Operator 文件放置位置

Airflow 默认把 AIRFLOW_HOME下的 dags/、plugins/、config/加入 PYTHONPATH。可以把 Operator 文件放置在 plugins 文件夹下或者自定义文件夹加入 PYTHONPATH。

6.3.1 创建 Operator 文件,并放置在 plugins 文件夹下

6.3.2 重启 airflow 服务

6.3.3 在 Dag 中导入自定义的 Operator

6.3.4 测试结果

7. Hook

Hook 是 Airflow 里对外部系统连接的抽象封装——它把认证、建连、重试、关闭等脏活全包了,让你(或 Operator/Sensor)只用关心"发什么 SQL / 调什么 API",不用管"怎么连、怎么认证"。

7.1 常见的 Hook

类别

Hook

用途

数据库

PostgresHook

PostgreSQL 查询/写入

MySqlHook

MySQL 查询/写入

SqliteHook

SQLite 操作

SnowflakeHook

Snowflake 数据仓库

BigQueryHook

Google BigQuery

RedshiftHook

Amazon Redshift

云存储

S3Hook

AWS S3 上传/下载/列举

GCSHook

Google Cloud Storage

SFTPHook

SFTP 文件传输

消息队列

KafkaHook

Kafka 生产/消费消息

HTTP

HttpHook

调用 REST API

自定义

继承 BaseHook

对接内部系统

7.2 Hook可以利用 airflow 的 connections 配置信息。配置一次,多有 Dag 共享。

7.3 在 Airflow 中自定义 Hook 核心是继承 BaseHook。唯一强制重写的是 get_conn()。可以把 Operator 文件放置在 plugins 文件夹下或者自定义文件夹加入 PYTHONPATH。

8. Sensor

Sensor 是 Airflow 里一种专门用来"等"的 Operator——它不干活,只检查某个条件是否满足,满足了才让下游 task 继续跑。

8.1 常见的sensor

Sensor

Provider

用途

S3KeySensor

amazon

等 S3 文件出现

S3PrefixSensor

amazon

等 S3 某个前缀下有文件

FileSensor

standard

等本地文件

SqlSensor

各数据库 provider

等 SQL 查询结果非空

HttpSensor

http

等 API 返回特定状态

ExternalTaskSensor

standard

等另一个 DAG/Task 完成

TimeSensor/ TimeDeltaSensor

standard

等到某个时间点

KubernetesPodSensor

cncf.kubernetes

等 Pod 到某状态

SnowflakeSensor

snowflake

等 Snowflake query 完成

8.2 deferrable 模式

基于 Triggerer​ 组件做异步等待。Sensor 把检查逻辑交给 Triggerer(单线程事件循环),自己释放 worker;条件满足后 Triggerer 通知 Scheduler 恢复 task。

8.3 实践

    t2a = TimeSensor(
        task_id="timeout_after_second_date_in_the_future_async",
        # timeout是从第一次 poke开始算起,Sensor 最多等待的总时长(秒)。超过这个时间条件还没满足,Sensor 就会判定为超时。
        timeout=1,
        # soft_fail=True时,Sensor 等待超时后会被标记为 SKIPPED(跳过),而不是 FAILED(失败)。
        soft_fail=True,
        target_time=(datetime.datetime.now(tz=datetime.timezone.utc) + datetime.timedelta(seconds=30)).time(),
        # 等待期交 Triggerer 的 asyncio 事件循环,完全不占 worker 槽,调度器压力低,需部署 airflow triggerer
        deferrable=True,
    )

8.4 自定义 Sensor 核心是继承 BaseSensorOperator并实现 poke;想要异步等待就配合 Trigger 做 deferrable。

9. Dynamic Dag

在 Airflow 中,Dynamic DAG(动态 DAG)​ 指的是在 DAG 文件被调度器解析时,通过 Python 代码动态决定 DAG 的结构、任务数量或 DAG 本身的数量,而不是手写死每一个 task_id。它解决的核心痛点是:避免重复代码、根据外部配置/数据自动生成流水线。

from datetime import datetime

from airflow.sdk import DAG, task
from airflow.providers.standard.operators.bash import BashOperator
from airflow.providers.standard.operators.empty import EmptyOperator

# A Dag represents a workflow, a collection of tasks
with DAG(dag_id="aaaa_demo", start_date=datetime(2022, 1, 1), schedule="0 0 * * *") as dag:
    # Tasks are represented as operators
    start = EmptyOperator(task_id="start")
    pre = start
    for i in range(5):
        t = BashOperator(task_id=f"task_{i}", bash_command=f"echo hello {i}")
        pre >> t
        pre = t
    end = EmptyOperator(task_id="end")
    pre >> end

10. Dag dependency

在 Airflow 中,DAG Dependency(DAG 间依赖)​ 是指一个 DAG 的执行需要等待另一个 DAG 完成,或由另一个 DAG 的数据产出触发。

10.1 Asset(数据驱动)

from airflow.sdk import Asset, dag, task
from datetime import datetime

orders = Asset("s3://warehouse/orders/daily/")

@dag(dag_id="produce_orders", schedule="@daily", start_date=datetime(2026, 1, 1), catchup=False)
def produce():
    @task(outlets=[orders]
    def build():
        print("write data")
    build()

@dag(dag_id="consume_orders", schedule=[orders], start_date=datetime(2026, 1, 1), catchup=False)
def consume():
    @task
    def transform():
        print("consume data")
    transform()
produce(); consume()

注意事项:

a. 只有带 outlets的任务成功,才会发 AssetEvent,下游才触发。

b. 多条件组合或者使用 & | 表达式

c. AssetAlias 解耦

生产者写 AssetAlias,消费者监听别名

from airflow.sdk import AssetAlias
orders_alias = AssetAlias("orders-daily")
# 生产者 outlets=[orders_alias];消费者 schedule=orders_alias

10.2 TriggerDagRunOperator(编排触发)

trigger = TriggerDagRunOperator(
        task_id="trigger_downstream",
        trigger_dag_id="child_report",       # 要触发的 DAG ID
        wait_for_completion=False,            # 是否等下游 DAG 跑完
        deferrable=False,                     # True 则用 async 等待
        conf={"source_date": "{{ ds }}"},     # 传给下游 DAG 的 conf
        reset_dag_run=True,                   # 如果目标 DAG Run 已存在,重置
    )

11.  密文的保存

Airflow 存密文分两层:Fernet 加密元库(单机/起步用)和外部 Secrets Backend(生产用,密码不落元库)。DAG 里永远只写 conn_id/ Variable.get,不写明文。

12. Plugins

13. template_searchpath

template_searchpath是 DAG 级别的参数,用于告诉 Airflow 的 Jinja 引擎去哪里找外部模板文件(如 .sql、.sh),避免把长脚本硬写在 Python 里。如果你的 sales_report.sql和 DAG 文件在同一个目录,不配 template_searchpath也能找到。

14. TriggerRule

TriggerRule 是 Airflow 中决定下游任务在上游(直接连接)处于哪些状态时才会被trigger的规则,默认是 all_success。

15. 总结

Airflow 提供的丰富组件和功能是 Airflow 能够成为最受欢迎的工作流编排工具的原因之一。

【版权声明】本文为华为云社区用户原创内容,未经允许不得转载,如需转载请自行联系原作者进行授权。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。