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,并将它们链式组合成更大的基于数据的工作流。如果您目前正在使用 ExternalTaskSensorTriggerDagRunOperator,应该了解一下数据集——在大多数情况下,您可以用它们来替代这些,并加速调度!

不过话不多说,来看一个简短的示例。首先我们写一个简单的 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_circuit TaskFlow 装饰器
  • 在 CLI 中新增删除角色命令
  • ExternalTaskSensor 添加对 TaskGroup 的支持
  • 新增 @task.kubernetes TaskFlow 装饰器
  • 实验性 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 才能成为如此成功的项目!

分享

阅读更多

Apache Airflow 2.5.0:滴答声

Ash Berlin‑Taylor

我们自豪地宣布 Apache Airflow 2.5.0 已发布,带来了众多提升用户体验的改动。

Apache Airflow 2.0 来了!

Ash Berlin‑Taylor

我们自豪地宣布 Apache Airflow 2.0.0 已正式发布。