Apache Airflow 2.4.0 包含超过 650 条“面向用户”的提交(不包括对 providers 或 chart 的提交),总计超过 870 条。内容包括 46 项新功能、39 项改进、52 条 bug 修复以及若干文档变更。
详情:
📦 PyPI: https://pypi.ac.cn/project/apache-airflow/2.4.0/
📚 文档: https://airflow.org.cn/docs/apache-airflow/2.4.0/
🛠️ 发布说明: https://airflow.org.cn/docs/apache-airflow/2.4.0/release_notes.html
🐳 Docker 镜像: docker pull apache/airflow:2.4.0
🚏 约束条件: https://github.com/apache/airflow/tree/constraints-2.4.0
数据感知调度 (AIP-48)
这次改动很大。Airflow 现在能够基于其他任务更新的数据集来调度 DAG。
这到底意味着什么?这是一个很棒的新功能,允许 DAG 作者创建更小、更独立的 DAG,并将它们链式组合成更大的基于数据的工作流。如果您目前正在使用 ExternalTaskSensor 或 TriggerDagRunOperator,应该了解一下数据集——在大多数情况下,您可以用它们来替代这些,并加速调度!
不过话不多说,来看一个简短的示例。首先我们写一个简单的 DAG,其中有一个任务叫 my_task,它会生成一个名为 my-dataset 的数据集。
from airflow import Dataset
dataset = Dataset(uri='my-dataset')
with DAG(dag_id='producer', ...)
@task(outlets=[dataset])
def my_task():
...
数据集是通过 URI 定义的。现在,我们可以创建第二个 DAG(consumer),它会在该数据集发生变化时被调度。
from airflow import Dataset
dataset = Dataset(uri='my-dataset')
with DAG(dag_id='dataset-consumer', schedule=[dataset]):
...
有了这两个 DAG,一旦 my_task 完成,Airflow 将为 dataset-consumer 工作流创建一次 DAG 运行。
我们知道目前的实现并不能满足所有用户对数据集的使用场景,未来的次要版本(2.5、2.6 等)中我们会在此基础上进行扩展和改进。
数据集代表抽象的数据集概念,当前(此版本)并不具备直接读写能力——我们在本次发行版中添加了将来要构建的基础特性,并且我们致力于通过更小的版本迭代更快地把新功能交到用户手中!
想了解更多关于数据集的信息,请参阅 数据感知调度文档。其中包括数据集如何通过 URI 标识、如何依赖多个数据集,以及如何理解数据集(提示:不要在数据集中加入“日期分区”,它的层级要高于此)。
使用新引入的 ExternalPythonOperator 更轻松地管理冲突的 Python 依赖
虽然我们希望所有 Python 库都能和谐共存,但现实并非如此,有时在 Airflow 环境中安装多个 Python 库会产生冲突——我们现在经常听到 dbt-core 的冲突。
为了解决这个问题,我们引入了 @task.external_python(以及对应的 ExternalPythonOperator),它允许您在预配置的虚拟环境,甚至是不同的 Python 版本中运行 Python 函数作为 Airflow 任务。例如
@task.external_python(python='/opt/venvs/task_deps/bin/python')
def my_task(data_interval_start, data_interval_env)
print(f'Looking at data between {data_interval_start} and {data_interval_end}')
...
根据您访问的上下文变量,虚拟环境中需要安装的内容会有细微差别,请务必阅读 ExternalPythonOperator 使用指南
对动态任务映射的更多改进 (AIP-42)
您提出需求,我们倾听。动态任务映射现在支持
expand_kwargs:为非 TaskFlow 运算符分配多个参数。zip:在不产生笛卡尔积的情况下组合多个对象。map:在任务运行前对参数进行转换。
欲了解动态任务映射的更多信息,请参阅文档新章节:映射数据的转换、上游数据的组合(即“zip”)、以及为非 TaskFlow 运算符分配多个参数。
在上下文管理器中自动注册 DAG(无需再写 as dag:)
这是一项小幅度的使用体验改进,我不想承认有多少次忘记了 as dag:,甚至出现了重复的 as dag:。
with DAG(dag_id="example") as dag:
...
@dag
def dag_maker():
...
dag2 = dag_maker()
可以变成
with DAG(dag_id="example"):
...
@dag
def my_dag():
...
my_dag()
如果出于任何原因想关闭此行为,可在 DAG 上设置 auto_register=False
# This dag will not be picked up by Airflow as it's not assigned to a variable
with DAG(dag_id="example", auto_register=False):
...
其他改进
超过 650 次提交,完整的功能、修复和变更列表太长,无法在此全部列出(请查看发布说明获取完整列表),但一些值得注意或有趣的小功能包括
- 主页自动刷新
- 新增
@task.short_circuitTaskFlow 装饰器 - 在 CLI 中新增删除角色命令
- 为
ExternalTaskSensor添加对TaskGroup的支持 - 新增
@task.kubernetesTaskFlow 装饰器 - 实验性
parsing_context,用于在工作节点上优化动态 DAG 处理 - 合并为单一的
schedule参数 - 在 管理员 → 配置 中允许显示非敏感配置项(而不是全部或全部不显示)
- 运算符名称与类分离(使用 TaskFlow 时不再出现
_PythonDecoratedOperator)
贡献者
感谢所有为此版本做出贡献的人员,包括 Andrey Anshin、Ash Berlin‑Taylor、Bartłomiej Hirsz、Brent Bovenzi、Chenglong Yan、D. Ferruzzi、Daniel Standish、Drew Hubl、Elad Kalif、Ephraim Anierobi、Jarek Potiuk、Jed Cunningham、Josh Fell、Mark Norman Francis、Niko、Tzu‑ping Chung、Vincent、Wojciech Januszek、chethanuk‑plutoflume、pierrejeambrun,以及所有其他 152 位提交者!正是有了你们,Airflow 才能成为如此成功的项目!
分享