进阶 dagster.io 2026-10-10 01:48:56 · 6 阅读

第20章 第 20 章:使用 Dagster 与 Airflow 构建数据管道

用 Dagster 重写 Airflow 教程

Dagster 从 Airflow 吸取了大量经验,致力于让数据应用的开发和维护更轻松。你可以在 Dagster 里构建很多有趣的东西(欢迎浏览我们的[示例](https://docs.dagster.io/examples)),但为了展示 Dagster 和 Airflow 在构建方式上的差异,最好用一个简单的例子来说明。

我们将在 Dagster 中实现 [Airflow 的入门教程](https://airflow.apache.org/docs/apache-airflow/stable/tutorial/fundamentals.html)。在逐步完成的过程中,我们会指出 Dagster 开发方式上的不同之处,以及如何看待数据任务。读完之后,你应该能体会到 Dagster 是如何帮助你快速上手构建的。

Pipeline 初始化

Airflow

Airflow 教程的第一步是初始化一个 DAG。Airflow 的 pipeline 在任务层面构建,再组织成 DAG。一个 Airflow 部署由多个 DAG 组成,每个任务都必须关联到某个 DAG:

from airflow.models.dag import DAG

with DAG( "tutorial", # These args will get passed on to each operator # You can override them on a per-task basis during operator initialization default_args={ "depends_on_past": False, "email": ["airflow@example.com"], "email_on_failure": False, "email_on_retry": False, "retries": 1, "retry_delay": timedelta(minutes=5), # 'queue': 'bash_queue', # 'pool': 'backfill', # 'priority_weight': 10, # 'end_date': datetime(2016, 1, 1), # 'wait_for_downstream': False, # 'sla': timedelta(hours=2), # 'execution_timeout': timedelta(seconds=300), # 'on_failure_callback': some_function, # or list of functions # 'on_success_callback': some_other_function, # or list of functions # 'on_retry_callback': another_function, # or list of functions # 'sla_miss_callback': yet_another_function, # or list of functions # 'on_skipped_callback': another_function, #or list of functions # 'trigger_rule': 'all_success' }, description="A simple tutorial DAG", schedule=timedelta(days=1), start_date=datetime(2021, 1, 1), catchup=False, tags=["example"], ) as dag:

Dagster

Dagster 不需要 DAG 或任何更高层的抽象。它不关注某个独立的 DAG 及其任务,而是把一切都视为资产。

资产是你数据平台的构建单元,对应数据栈中的各种操作。资产可以是云存储中的文件、数据库中的表,也可以是 ML 模型。以资产为中心的思路更符合现代数据栈的演进方式——一个资产往往有多个上游和下游依赖。

Dagster 不会强迫你把某个资产绑定到一个特定的 DAG 或 pipeline。资产及其关系会随时间自然生长,由 Dagster 负责管理这些资产之间的关系。

节点

Airflow

DAG 里包含的是 operator,也就是 Airflow 实际执行的工作。要使用 operator,必须先把它初始化为任务。本教程使用 `BashOperator` 来执行 shell 命令。Airflow 提供多种类型的 operator,各有自己的参数和用法(尽管它们都继承自 `BashOperator`):

from airflow.operators.bash import BashOperator

t1 = BashOperator( task_id="print_date", bash_command="date", )

t2 = BashOperator( task_id="sleep", depends_on_past=False, bash_command="sleep 5", retries=3, )

Dagster

前面说过,Dagster 把工作视为资产。Dagster 还主张数据工程应该像软件工程一样,所以定义资产的过程更 Pythonic、更像写普通函数(Airflow 中与资产最接近的是 TasksFlow API,但相比 Dagster 仍有局限)。

在 Dagster 中定义工作、创建执行 bash 命令的资产有几种方式。我们可以用 Pipes 执行 Python 以外的语言,同时保留 Dagster 的编排能力。但为了与 Airflow 教程保持一致,这里直接用 `subprocess` 执行所需命令:

import subprocess import dagster as dg

@dg.asset def print_date(): subprocess.run(["date"])

@dg.asset( retry_policy=dg.RetryPolicy(max_retries=3) ) def sleep(): subprocess.run(["sleep", "5"])

这看起来就是标准的 Python 代码,唯一 Dagster 特有的部分是 `dg.asset` 装饰器,它把函数变成资产。这个装饰器还能设置一些执行相关的参数,比如重试策略。总体而言,编写 Dagster 资产需要的领域特定知识更少。

模板

Airflow

Airflow 教程的最后一个任务专门用来展示 Airflow 对 Jinja 模板的支持:

import textwrap

templated_command = textwrap.dedent( """ {% for i in range(5) %} echo "{{ ds }}" echo "{{ macros.ds_add(ds, 7)}}" {% endfor %} """ )

t3 = BashOperator( task_id="templated", depends_on_past=False, bash_command=templated_command, )

Jinja 有用,但也可能很晦涩。如果不熟悉 Airflow 的宏,你可能一时看不懂这段代码在做什么。

和 DAG 里的其他 operator 一样,这个任务是把 Jinja 模板渲染成 bash 命令来执行,从而用不同日期多次运行 `echo` 命令。

Dagster

Dagster 鼓励你写更明确的代码。不依赖模板或 Jinja,而是直接使用标准 Python。同样的效果可以这样实现:

from datetime import datetime, timedelta

@dg.asset() def templated(): ds = datetime.today().strftime("%Y-%m-%d") ds_add = (datetime.today() + timedelta(days=7)).strftime("%Y-%m-%d") for _ in range(5): formatted_string = f""" echo "{ds}" echo "{ds_add}" """ subprocess.run(formatted_string, shell=True, check=True)

文档

Airflow

Airflow 中可以为任务和 DAG 编写文档,但并不总是直观。任务定义之后,你可以这样关联文档:

t1.doc_md = textwrap.dedent( """\ #### Task Documentation You can document your task using the attributes `doc_md` (markdown), `doc` (plain text), `doc_rst`, `doc_json`, `doc_yaml` which gets rendered in the UI's Task Instance Details page. ![img](https://imgs.xkcd.com/comics/fixing_problems.png) Image Credit: Randall Munroe, [XKCD](https://xkcd.com/license.html) """ )

Dagster

Dagster 通过 Python docstring 支持资产文档。文档直接与函数绑定在一起,维护起来容易得多:

@dg.asset def print_date(): """ You can use the docstring. ![img](https://imgs.xkcd.com/comics/fixing_problems.png) Image Credit: Randall Munroe, [XKCD](https://xkcd.com/license.html) """ result = subprocess.run(["date"], capture_output=True, text=True) print(result)

docstring 里还支持 Markdown,你写的内容会在查看资产时渲染在 Dagster 的资产目录中。

设置依赖

Airflow

除了定义 DAG 和任务,你还需要显式设置所有任务之间的关系。教程中,任务 `t2` 和 `t3` 依赖于任务 `t1`。要构建这个图结构,你需要这样设置依赖:

t1 >> [t2, t3]

Dagster

在生产环境中,你可能有几十甚至上百个资产,显式定义所有关系是不现实的。因此 Dagster 把节点之间的关系直接写在资产定义内部。要构建同样的图,只需在 `sleep` 和 `templated` 的资产装饰器中把 `print_date` 声明为依赖:

@dg.asset def print_date(): ...

@dg.asset( deps=[print_date], ) def sleep(): # asset depends on print_date ...

@dg.asset( deps=[print_date], ) def templated(): # asset depends on print_date ...

让 Dagster 来维护你的依赖图,更易管理、也更具扩展性。在资产目录中可以查看所有资产,并深入查看具体节点的关系。由于 Dagster 不把节点限制在 DAG 层面,你能获得整个 data mesh 的完整全景视图。

启动

Airflow

Airflow 教程的最后一步是对 DAG 做一次测试运行。你可以通过 CLI 执行 Airflow 运行:

airflow tasks test tutorial print_date 2015-06-01

用 CLI 执行运行可以免去启动 Airflow 服务的麻烦。

Airflow 作为一个服务由多个组件构成。要在本地跑起完整的 Airflow,你需要启动数据库、webserver 和 scheduler。没有这些组件,就无法通过 Airflow 的 UI 执行 pipeline。

要在本地启动这些组件,你可以配置并运行 `airflow standalone`,或者分别启动:

airflow db init

airflow users create \ --username admin \ --firstname Peter \ --lastname Parker \ --role Admin \ --email spiderman@superhero.org

airflow webserver --port 8080

airflow scheduler

Dagster

Dagster 让你能尽快进入 UI。假设资产代码已保存在文件中(比如 `tutorial.py`),用 Dagster 自带的 CLI 就能启动 UI:

dagster dev -f tutorial.py

这一条命令会启动一个包含 Dagster 全部功能的临时实例。你可以在资产目录中查看资产,体验 schedules、sensors 等功能。

在 Dagster 中映射依赖也更容易。任何运行 `dagster dev` 的环境都可以成为一个 code location,而 code location 的数量不限。这样你可以为资产量身定制环境,同时仍然通过同一个编排层统一管理。

如果你想通过 API 调用 Dagster,Dagster UI 的所有操作背后都有完整的 GraphQL 层支撑。借助这个 API,你可以在 UI 之外以编程方式完成任何需要的操作。

结论

对比一下 Dagster 和 Airflow 的最终代码:

Airflow

import textwrap from datetime import datetime, timedelta

# The DAG object; we'll need this to instantiate a DAG from airflow.models.dag import DAG

# Operators; we need this to operate! from airflow.operators.bash import BashOperator with DAG( "tutorial", # These args will get passed on to each operator # You can override them on a per-task basis during operator initialization default_args={ "depends_on_past": False, "email": ["airflow@example.com"], "email_on_failure": False, "email_on_retry": False, "retries": 1, "retry_delay": timedelta(minutes=5), # 'queue': 'bash_queue', # 'pool': 'backfill', # 'priority_weight': 10, # 'end_date': datetime(2016, 1, 1), # 'wait_for_downstream': False, # 'sla': timedelta(hours=2), # 'execution_timeout': timedelta(seconds=300), # 'on_failure_callback': some_function, # or list of functions # 'on_success_callback': some_other_function, # or list of functions # 'on_retry_callback': another_function, # or list of functions # 'sla_miss_callback': yet_another_function, # or list of functions # 'on_skipped_callback': another_function, #or list of functions # 'trigger_rule': 'all_success' }, description="A simple tutorial DAG", schedule=timedelta(days=1), start_date=datetime(2021, 1, 1), catchup=False, tags=["example"], ) as dag:

# t1, t2 and t3 are examples of tasks created by instantiating operators t1 = BashOperator( task_id="print_date", bash_command="date", )

t2 = BashOperator( task_id="sleep", depends_on_past=False, bash_command="sleep 5", retries=3, ) t1.doc_md = textwrap.dedent( """\ #### Task Documentation You can document your task using the attributes `doc_md` (markdown), `doc` (plain text), `doc_rst`, `doc_json`, `doc_yaml` which gets rendered in the UI's Task Instance Details page. ![img](https://imgs.xkcd.com/comics/fixing_problems.png) Image Credit: Randall Munroe, [XKCD](https://xkcd.com/license.html) """ )

dag.doc_md = __doc__ # providing that you have a docstring at the beginning of the DAG; OR dag.doc_md = """ This is a documentation placed anywhere """ # otherwise, type it like this templated_command = textwrap.dedent( """ {% for i in range(5) %} echo "{{ ds }}" echo "{{ macros.ds_add(ds, 7)}}" {% endfor %} """ )

t3 = BashOperator( task_id="templated", depends_on_past=False, bash_command=templated_command, )

t1 >> [t2, t3]

Dagster

import subprocess from datetime import datetime, timedelta

import dagster as dg

@dg.asset def print_date(): """ You can use the docstring.

![img](https://imgs.xkcd.com/comics/fixing_problems.png)

Image Credit: Randall Munroe, [XKCD](https://xkcd.com/license.html) """ subprocess.run(["date"])

@dg.asset( deps=[print_date], retry_policy=dg.RetryPolicy(max_retries=3), ) def sleep(): subprocess.run(["sleep", "5"])

@dg.asset( deps=[print_date], ) def templated(): ds = datetime.today().strftime("%Y-%m-%d") ds_add = (datetime.today() + timedelta(days=7)).strftime("%Y-%m-%d")

for _ in range(5): formatted_string = f""" echo "{ds}" echo "{ds_add}" """ subprocess.run(formatted_string, shell=True, check=True)

Airflow UI

Dagster UI

即使你从未用过 Dagster,也会觉得 Dagster 的代码更 Pythonic。做数据开发时,你希望尽量减少摩擦,专注于资产本身——资产才是数据平台价值所在。

本教程只展示了 Dagster 能力的一小部分。随着深入使用,你会进一步体会到它在规模化执行资产、构建可持续演进的应用方面的特性。但现在最重要的是看到上手有多简单。如果你一直想构建某个数据 pipeline,或者一直想重构某个 Airflow 里的 DAG,不妨在 Dagster 里试一试。

如果你在数据领域工作了一段时间,很可能对 Airflow 很熟悉。自十多年前发布以来,Airflow 建立了许多数据 pipeline 构建方面的概念。但随着数据工程的不断演进,Airflow 的某些实践在数据工程中已经显得不那么顺手了。

有反馈或问题?欢迎在 Slack 或 GitHub 上发起讨论。

有兴趣和我们一起工作?查看我们的在招职位。

想要更多这类内容?在 LinkedIn 上关注我们。

评论 (0)