site logo

Marico's space

如何将 Apache Airflow 与 OpenLineage 集成以实现端到端可追溯性

AI技术与应用 2026-09-04 11:30:52 3

最近折腾了 Airflow 和 OpenLineage 的集成,把数据管道的血缘关系跑通了。这篇说说具体怎么配,以及踩过的几个坑。

跑完之后,你 Airflow 实例里的每个 DAG Run 都会发出结构化的血缘事件,精确记录每个任务读写了哪些表,打开图就能回答"这个数字是哪个上游任务算出来的",不用再 grep scheduler 日志。

这就是核心价值。不需要手动写文档,不需要维护一个三周后就过时的血缘表格。编排器在运行的同时报告它实际做了什么。

配置本身不难。下面的步骤顺序是精心安排的,先把最容易出问题的三个地方暴露出来,省得你花一下午瞎猜。

环境准备和版本要求

  • Apache Airflow 2.11.0 或更高版本,或任意 Airflow 3.x 版本。这是当前 provider 分发包支持的最低版本。
  • Python 3.9 到 3.12。
  • Docker,用于本地运行血缘后端。
  • Airflow 里的一个 Postgres 连接(postgres_default),如果你想完整复现后面的 SQL 示例的话。

关键就两个包:Airflow OpenLineage provider负责从 Airflow 提取元数据并转换成事件,openlineage-python负责发送。客户端可以独立于 provider 升级,这意味着你需要传输层修复时不用动 Airflow 版本。

第一步:事件模型,动手安装之前先搞明白

OpenLineage 有三个核心对象和一个扩展机制。大多数首次集成最后图是空的,就是跳过这一步导致的。

  • Job(任务):运行的东西。你的 DAG 是一个 Job,每个 Task 也是。
  • Run(运行):一次任务执行,有唯一的 Run ID。
  • Dataset(数据集):被读取或写入的东西。由 namespacename 标识。
  • Facet(切面):附加到上述任意对象的原子元数据块。Schema、SQL 文本、列级血缘、运行状态以及自定义字段都以 Facet 形式传递,OpenLineage 规范列出了标准切面。

事件在状态转换时触发:STARTRUNNINGCOMPLETEFAILABORTOTHER。血缘关系是在下游通过跨 Run 连接数据集重建的,这意味着数据集标识决定了图的质量。后面会详细说,因为这是图有节点无边最常见的原因。

第二步:安装 provider

pip install apache-airflow-providers-openlineage

官方 Airflow Docker 镜像可能已经内置了。先检查一下再往 requirements 文件里加:

airflow providers list | grep openlineage

这时候还没有任何事件发出。Provider 会保持沉默,直到它知道该把事件发到哪里。

第三步:配置传输方式

先从 Console 传输开始。它把事件写到任务日志里,零成本运行,而且能立即告诉你提取功能是否正常工作。

export AIRFLOW__OPENLINEAGE__TRANSPORT='{"type": "console"}'

日志里看到事件之后,换成真实后端。Marquez 是一个不错的选择,任何兼容 OpenLineage 的后端都可以用。它是标准的参考实现,能最快速度跑起血缘 UI:

git clone https://github.com/MarquezProject/marquez
cd marquez
./docker/up.sh

API 监听 5000 端口,管理界面在 5001,Web UI 在 3000。在 macOS 上,5000 端口被系统占用了,需要运行 ./docker/up.sh --api-port 9000,然后相应调整下面的 URL。

现在切换传输方式:

export AIRFLOW__OPENLINEAGE__TRANSPORT='{"type": "http", "url": "http://localhost:5000", "endpoint": "api/v1/lineage"}'
export AIRFLOW__OPENLINEAGE__NAMESPACE='airflow-local'

同样的配置写到 airflow.cfg 里:

[openlineage]
transport = {"type": "http", "url": "http://localhost:5000", "endpoint": "api/v1/lineage"}
namespace = airflow-local
disabled = False

namespace 要主动设置。它在逻辑上隔离不同的生产者,这样测试环境 Airflow 和生产环境 Airflow 不会混到一张图里给你误导。如果不设置,所有数据都会落到 default 里。

除了本地环境,不要把凭证写到 airflow.cfg 里。Provider 支持传入一个通用的 Airflow connection ID,传输配置(包括认证)放在 connection extra 里。

第四步:跑一个真正能产出血缘的 DAG

SQL 操作符是最适合入门的,因为 provider 会解析查询,自动推导出输入、输出和列级关系,不用你写一行代码:

# dags/openlineage_demo.py
from datetime import datetime from airflow import DAG
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator with DAG( dag_id="openlineage_demo", start_date=datetime(2026, 1, 1), schedule="@daily", catchup=False,
) as dag: build_daily_orders = SQLExecuteQueryOperator( task_id="build_daily_orders", conn_id="postgres_default", sql=""" CREATE TABLE IF NOT EXISTS analytics.daily_orders AS SELECT o.order_date, c.region, COUNT(*) AS order_count, SUM(o.amount) AS revenue FROM raw.orders o JOIN raw.customers c ON c.customer_id = o.customer_id GROUP BY o.order_date, c.region; """, )

触发这个 DAG。生成的事件应该列出 raw.ordersraw.customers 作为输入,analytics.daily_orders 作为输出,并且带有一个列级 Facet,把 revenue 映射回 o.amount

第五步:验证是否生效

三个检查,按顺序来。

看图。打开 http://localhost:3000,找到你配置的 namespace,确认两个源表连接到了输出表。

检查哪些任务在上报。这是最容易漏掉的诊断。DagRunSTART 事件带有一个 AirflowJobFacet,列出了 DAG 里的每个任务,每个任务都有一个 emits_ol_events 布尔值。在图缺了一半之前,这个信息就能告诉你哪些操作符会保持静默,而不是让你瞎猜。

理解静默的原因。EmptyOperator 默认不发出任何事件,因为 Airflow 对它的调度方式和对真正任务的不同。如果需要它在图里出现,加上 on_executeon_success 回调,或者一个 task outlet。当任务级细节不重要时,DagRunCOMPLETE 事件携带一个 AirflowStateRunFacet,包含 Run 里每个任务的状态。

第六步:覆盖那些不上报的操作符

自动提取覆盖了 SQL 操作符和很多 provider 操作符。你自己的操作符在告诉它接触到什么之前什么都不会报。

对于你拥有的操作符,直接实现 OpenLineage 方法:

from airflow.models import BaseOperator class S3ToWarehouseOperator(BaseOperator): def __init__(self, *, source_bucket, source_key, target_table, **kwargs): super().__init__(**kwargs) self.source_bucket = source_bucket self.source_key = source_key self.target_table = target_table def execute(self, context): ... # your copy logic def get_openlineage_facets_on_complete(self, task_instance): # 本地导入:在顶层导入 Airflow 可能产生循环依赖 # 导致提取静默失败 from airflow.providers.common.compat.openlineage.facet import Dataset from airflow.providers.openlineage.extractors import OperatorLineage return OperatorLineage( inputs=[ Dataset(namespace=f"s3://{self.source_bucket}", name=self.source_key) ], outputs=[ Dataset(namespace="postgres://warehouse:5432", name=self.target_table) ], )

几个值得牢记的规则:

  • 必须实现 get_openlineage_facets_on_start()get_openlineage_facets_on_complete(ti) 至少之一。如果缺少 on_complete,provider 会回退到 on_start。还有一个 get_openlineage_facets_on_failure(ti),默认复用 on_complete 的逻辑。
  • 优先使用 on_complete,当真实的数据集名称只在 execute 期间才能确定时。在开始时报告通配符路径但从不纠正,会产生一个看起来很自信但实际错误的图。
  • 在方法内部导入 OpenLineage 对象,不要在模块级别。监听器是在 worker 启动时实例化的,顶层导入 Airflow 可能产生循环依赖,无声地杀死提取过程,还不会给出明显的错误。

对于无法修改的第三方操作符,写一个自定义 extractor 并注册它:

export AIRFLOW__OPENLINEAGE__EXTRACTORS='plugins.extractors.MyCustomExtractor'

第七步:附加你自己的上下文

从 provider 版本 1.10.0 开始,你可以注入任意的 run facets 而不用碰操作符代码。写一个接受 task instance 并返回 facet 字典的函数,然后注册导入路径,用分号分隔:

export AIRFLOW__OPENLINEAGE__CUSTOM_RUN_FACETS='plugins.ol_facets.ownership_facet'

这是把团队、成本中心或变更工单 ID 附加到每个事件上的方式,把"谁负责这个坏掉的管道"从 Slack 讨论串变成一个过滤器。

生产环境里会坏的地方

四个配置项加一个习惯,涵盖了大部分坑。

  • include_full_task_info:很诱人,但很贵。开启后所有可序列化的任务参数都会进入事件。根据你传给任务的内容,单个事件可能达到几兆字节。
  • execution_timeout:限制提取可以运行多久,防止慢血缘调用变成管道事故。
  • dag_state_change_process_pool_size:scheduler 用于异步处理 DAG 状态变化的进程数。在繁忙的实例上值得调优。
  • emission_policy:当前控制发送内容的做法。老的 selective_enabledisable_source_code 标志已在它的 favor 中弃用。
  • 命名规范:血缘通过数据集标识连接。如果一个任务写 analytics.daily_orders,另一个读 ANALYTICS.DAILY_ORDERS,你会得到两个节点和一条边都没有。在规模化之前修复这个约定(database.schema.table、环境作用域的 namespace),因为事后重写标识意味着要重新处理历史数据。

血缘不是答案的场景

你现在拥有的是运营血缘:什么运行了、它读了什么、写了什么、是否失败。足够追溯一个事故的根源,也足够在 schema 变更前做影响分析。

但不足以回答谁拥有一个数据集、它是否经过认证、"活跃客户"在业务层面是什么意思、PII 流向了哪里。这些属于元数据平台,OpenLineage 事件是它的输入而非替代品。如果你朝这个方向走,这篇关于数据治理和端到端血缘与元数据目录的文章涵盖了采集层和治理层如何分工。

但顺序很重要。先发出,再编目。靠手工喂养的目录老化的速度和你要替换的文档一样快。