Airflow 必知必会(1)
1.概述
Airflow 是使用最广泛的开源工作流编排平台之一,理解 Airflow 的重要特性对于创建好的工作流至关重要。

2. Xcoms
默认情况下,任务实例是完全隔离的,并且可能运行在不同的机器上,所以任务之间无法传递数据。Xcoms 是 Airflow 的一个工作流中任务实例之间传递小批量元数据的功能(Cross-communications 的缩写)。
2.1 PythonOperator
2.1.1 Push 的两种方式
显式 push
返回值 push
2.1.2 Pull 的三种方式
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 |
用途 |
|---|---|---|
|
数据库 |
|
PostgreSQL 查询/写入 |
|
|
MySQL 查询/写入 |
|
|
|
SQLite 操作 |
|
|
|
Snowflake 数据仓库 |
|
|
|
Google BigQuery |
|
|
|
Amazon Redshift |
|
|
云存储 |
|
AWS S3 上传/下载/列举 |
|
|
Google Cloud Storage |
|
|
|
SFTP 文件传输 |
|
|
消息队列 |
|
Kafka 生产/消费消息 |
|
HTTP |
|
调用 REST API |
|
自定义 |
继承 |
对接内部系统 |
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 |
用途 |
|---|---|---|
|
|
amazon |
等 S3 文件出现 |
|
|
amazon |
等 S3 某个前缀下有文件 |
|
|
standard |
等本地文件 |
|
|
各数据库 provider |
等 SQL 查询结果非空 |
|
|
http |
等 API 返回特定状态 |
|
|
standard |
等另一个 DAG/Task 完成 |
|
|
standard |
等到某个时间点 |
|
|
cncf.kubernetes |
等 Pod 到某状态 |
|
|
snowflake |
等 Snowflake query 完成 |
8.2 deferrable 模式
基于 Triggerer 组件做异步等待。Sensor 把检查逻辑交给 Triggerer(单线程事件循环),自己释放 worker;条件满足后 Triggerer 通知 Scheduler 恢复 task。
8.3 实践
8.4 自定义 Sensor 核心是继承 BaseSensorOperator并实现 poke;想要异步等待就配合 Trigger 做 deferrable。
9. Dynamic Dag
在 Airflow 中,Dynamic DAG(动态 DAG) 指的是在 DAG 文件被调度器解析时,通过 Python 代码动态决定 DAG 的结构、任务数量或 DAG 本身的数量,而不是手写死每一个 task_id。它解决的核心痛点是:避免重复代码、根据外部配置/数据自动生成流水线。

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 能够成为最受欢迎的工作流编排工具的原因之一。
- 点赞
- 收藏
- 关注作者
评论(0)