OpenLineage + Marquez 实战:5 步搭起端到端血缘系统,2,604 ⭐ + 2,259 ⭐ 的标准 + 后端组合拳

43次阅读
OpenLineage + Marquez 实战:5 步搭起端到端血缘系统,2,604 ⭐ + 2,259 ⭐ 的标准 + 后端组合拳

飞熊实战向 · 数据更新于 2026-08-17 · 目标:从 Docker 部署到 Airflow/Spark/dbt 全栈集成, 让你 30 分钟拥有完整的端到端血缘系统 。这一篇不讲理论, 只讲操作


0. 实战目标

部署一套完整的血缘收集 + 可视化系统:

Spark / Airflow / dbt / Flink
    ↓ 发射 OpenLineage 事件
Marquez(OpenLineage 官方参考后端)↓ 持久化 + 可视化
http://localhost:3000  ← 在这里看血缘图

预计耗时 :30 分钟
预计资源:Docker + 4GB 内存


1. 环境准备

1.1 必备依赖

# Docker + Docker Compose
docker --version   # ≥ 20.10
docker compose version   # ≥ 2.0

# Git
git --version

# 可选:jq(解析 JSON)apt install -y jq   # Ubuntu/Debian
brew install jq     # macOS

1.2 硬件最低要求

资源 最低 推荐
CPU 2 核 4 核
内存 4 GB 8 GB
磁盘 10 GB 20 GB

1.3 网络要求

  • 能拉 Docker Hub 镜像
  • 能访问 GitHub(克隆仓库)

2. Step 1:部署 Marquez(5 分钟)

Marquez 是 OpenLineage 官方参考后端,自带 Web UI + PostgreSQL + REST/GraphQL API。

2.1 克隆仓库

git clone https://github.com/MarquezProject/marquez.git
cd marquez

# MacOS 用户:先配置 Git 自动转换
git config --global core.autocrlf false

2.2 配置 Git 大文件

Marquez 用了 Git LFS 存 demo 数据:

git lfs install
git lfs pull

2.3 一键启动

# 默认端口:HTTP API 5000 + Web UI 3000
./docker/up.sh

# MacOS 用户:5000 被 AirPlay 占用
./docker/up.sh --api-port 9000

# 想看样例血缘:加 --seed
./docker/up.sh --seed

2.4 验证启动

# 检查容器
docker compose ps
# 应该看到:marquez-web / marquez-api / postgres

# 检查 API 健康
curl http://localhost:5001/healthcheck
# 预期输出:{"status":"UP"}  或类似 JSON

# 打开 Web UI
open http://localhost:3000

2.5 看到 Marquez Web UI

首次打开你会看到:
– 左侧栏:Namespaces / Datasets / Jobs
– 中间:Lineage Graph(血缘图)
– 顶部:搜索框

如果加了 --seed,会有 example_namespace 的样例数据。


3. Step 2:集成 Airflow(10 分钟)

Airflow 是最常见的 ETL 调度器,OpenLineage 提供原生 Hook

3.1 安装 OpenLineage Airflow 库

pip install openlineage-airflow

3.2 配置 OpenLineage 连接

# 环境变量方式(推荐)export OPENLINEAGE_URL=http://localhost:5000
export OPENLINEAGE_API_KEY=   # 本地不设
export OPENLINEAGE_NAMESPACE=production
export OPENLINEAGE_AIRFLOW_SOURCE=my-airflow

或者写入 airflow.cfg

[openlineage]
transport = {
  "type": "http",
  "url": "http://localhost:5000",
  "auth": {"type": "api_key", "api_key": ""}
}
namespace = production

3.3 修改 DAG 自动捕获

# dags/etl_orders.py
from airflow import DAG
from airflow.providers.http.operators.http import SimpleHttpOperator
from airflow.operators.python import PythonOperator
from datetime import datetime

default_args = {
    'owner': 'data-eng',
    'start_date': datetime(2026, 1, 1),
}

with DAG(
    'etl_orders',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False,
) as dag:

    extract = PythonOperator(
        task_id='extract',
        python_callable=lambda: print("extracting orders"),
    )

    transform = PythonOperator(
        task_id='transform',
        python_callable=lambda: print("transforming orders"),
    )

    load = PythonOperator(
        task_id='load',
        python_callable=lambda: print("loading orders"),
    )

    extract >> transform >> load

无需额外代码!OpenLineage Airflow Operator 会自动捕获 DAG 结构、task 执行时间、状态、参数,并发射 OpenLineage 事件。

3.4 触发 DAG + 验证

# Airflow unpaused
airflow dags unpause etl_orders

# 触发一次
airflow dags trigger etl_orders

# 1-2 分钟后看 Marquez
open http://localhost:3000

你应该看到:

production.etl_orders job 出现在左侧
– 一次 Run 记录(COMPLETED 状态)
– 血缘图自动生成


4. Step 3:集成 Spark(10 分钟)

Spark 是大数据处理的事实标准,OpenLineage 通过 SparkListener 自动捕获血缘

4.1 安装 Spark OpenLineage 集成

# 下载 OpenLineage Spark JAR
wget https://repo1.maven.org/maven2/io/openlineage/openlineage-spark_2.12/1.52.0/openlineage-spark_2.12-1.52.0.jar

或者 Maven 依赖:

<dependency>
    <groupId>io.openlineage</groupId>
    <artifactId>openlineage-spark_2.12</artifactId>
    <version>1.52.0</version>
</dependency>

4.2 配置 spark-submit

spark-submit \
  --packages io.openlineage:openlineage-spark_2.12:1.52.0 \
  --conf spark.openlineage.url=http://localhost:5000 \
  --conf spark.openlineage.namespace=production \
  --conf spark.openlineage.transport.type=http \
  --conf spark.sql.query.lineage.enabled=true \
  your_job.py

4.3 写一个简单的 Spark 血缘测试

# spark_lineage_test.py
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("LineageTest") \
    .config("spark.openlineage.url", "http://localhost:5000") \
    .config("spark.openlineage.namespace", "production") \
    .getOrCreate()

# 读 + 写 → 自动捕获血缘
df = spark.read.parquet("s3://my-bucket/raw/orders/")
result = df.filter(df.amount > 100).select("order_id", "amount")
result.write.mode("overwrite").parquet("s3://my-bucket/processed/orders/")

spark.stop()

4.4 验证

spark-submit --master local[*] spark_lineage_test.py

# 在 Marquez Web UI 看:# - jobs: production.LineageTest
# - datasets: production.{raw_orders, processed_orders}
# - lineage graph: raw_orders → processed_orders

5. Step 4:集成 dbt(5 分钟)

dbt 是 ELT 时代的 SQL 转换工具,OpenLineage 提供 dbt 插件

5.1 安装 dbt OpenLineage 插件

pip install openlineage-dbt

5.2 配置 dbt 项目

dbt_project.yml

name: 'analytics'
version: '1.0.0'

profile: 'analytics'

vars:
  openlineage_emitter: {
    "type": "http",
    "url": "http://localhost:5000",
    "namespace": "production"
  }

models:
  analytics:
    +materialized: table

5.3 配置 profiles.yml

~/.dbt/profiles.yml

analytics:
  target: dev
  outputs:
    dev:
      type: postgres
      host: localhost
      port: 5432
      user: postgres
      password: postgres
      dbname: analytics
      schema: public

5.4 配置 dbt 运行

# 通过环境变量指定 emitter
DBT_EMITTER_TYPE=http \
DBT_EMITTER_URL=http://localhost:5000 \
DBT_EMITTER_NAMESPACE=production \
dbt-ol run

5.5 验证

dbt-ol 会在每次 dbt run 后发射 OpenLineage 事件,包含:
– 每个 model 的输入输出表
– SQL 文本
– 列级血缘(如果 dbt project 配置了)
– Run 时间 / 状态

Marquez Web UI 会自动显示:

analytics.{model_name} 作为 jobs
– 源表 + 目标表作为 datasets
– dbt run 之间的依赖图


6. Step 5:联合验证(5 分钟)

6.1 跑一个完整 ETL 流水线

# 1. Spark 任务(生产 raw_orders)spark-submit --packages io.openlineage:openlineage-spark_2.12:1.52.0 \
  --conf spark.openlineage.url=http://localhost:5000 \
  job_extract.py

# 2. dbt 任务(转换 + 业务建模)DBT_EMITTER_TYPE=http \
DBT_EMITTER_URL=http://localhost:5000 \
DBT_EMITTER_NAMESPACE=production \
dbt-ol run

# 3. Airflow DAG(编排 + 通知)airflow dags trigger etl_orders

6.2 在 Marquez 看完整血缘

打开 http://localhost:3000,你应该看到:

Lineage Graph

[S3 raw/orders]
     ↓ (Spark job_extract)
[S3 processed/orders]
     ↓ (dbt model stg_orders)
[dbt model fct_orders]
     ↓ (Airflow DAG etl_orders)
   [BI / 消费端]

Job 列表

production.job_extract (Spark)

analytics.stg_orders (dbt)

analytics.fct_orders (dbt)

production.etl_orders (Airflow)

Dataset 列表

s3://my-bucket/raw/orders


s3://my-bucket/processed/orders


analytics.public.stg_orders


analytics.public.fct_orders

Run 历史
– 每次执行的时间、状态、参数

6.3 GraphQL 查询

# GraphQL playground
open http://localhost:5000/graphql-playground

# 查询 namespaces
{
  namespaces {
    name
    currentSource {
      type
      name
    }
  }
}

6.4 REST API 查询

# 列出所有 jobs
curl http://localhost:5000/api/v1/jobs | jq

# 列出某个 namespace 的 datasets
curl http://localhost:5000/api/v1/namespaces/production/datasets | jq

# 获取某 dataset 的血缘
curl http://localhost:5000/api/v1/datasets/namespace/name/lineage | jq

7. 生产考虑

7.1 性能

组件 默认配置 生产建议
PostgreSQL 1 GB 内存 8 GB+ · SSD · 定期 vacuum
Marquez API 512 MB 2 GB+ · JVM 调优
Marquez Web 256 MB 1 GB
Kafka(可选) 5 节点集群

7.2 灾备

# 备份 PostgreSQL
docker exec marquez-postgres pg_dump -U postgres marquez > backup.sql

# 恢复
cat backup.sql | docker exec -i marquez-postgres psql -U postgres marquez

7.3 升级路径

# 拉最新代码
cd marquez && git pull

# 重启
./docker/down.sh && ./docker/up.sh --build

注意:Marquez 与 OpenLineage 兼容表——保持两者最新,避开 0.49.x 之前的版本(已 maintenance)。

7.4 安全

# Marquez 默认无鉴权(开发模式)# 生产部署务必:# 1. 加 Nginx 反向代理 + Basic Auth
# 2. PostgreSQL 启用 TLS
# 3. OpenLineage API 启用 API_KEY

export OPENLINEAGE_API_KEY=<strong-random-key>

7.5 K8s 部署

helm repo add marquez https://marquezproject.github.io/marquez
helm install marquez marquez/marquez \
  --set postgresql.auth.password=secret \
  --namespace marquez --create-namespace

8. 飞熊实战建议

8.1 客户场景选型

场景 组合 原因
小团队 <5 人 Marquez + OpenLineage Airflow Docker 一键起,零运维成本
中团队 5-20 人 Marquez + OpenLineage Airflow + Spark + dbt 全栈集成,统一命名空间
大团队 20+ 人 Marquez(轻量后端)+ DataHub/OpenMetadata(一站式) 后端 + 前端分离
AI 时代数据平台 OpenMetadata + OpenLineage 主流云原生组合

8.2 30 分钟落地清单

  • [] 克隆 Marquez(1 分钟)
  • [] ./docker/up.sh --seed(2 分钟)
  • [] 打开 Web UI 验证(1 分钟)
  • [] 配置 Airflow OPENLINEAGE_URL(3 分钟)
  • [] 触发 DAG + 看 Marquez(5 分钟)
  • [] Spark JAR 配置(5 分钟)
  • [] 跑测试 Spark job(3 分钟)
  • [] dbt-ol 配置 + 跑模型(5 分钟)
  • [] 看完整血缘图(4 分钟)

总计 29 分钟,全栈血缘上线。

8.3 调试技巧

看不到事件?

# 1. 检查 Marquez 是否收到
curl http://localhost:5001/healthcheck

# 2. 检查 OpenLineage 客户端日志
# Airflow: 看 task log 末尾
# Spark: 看 driver log 的 "OpenLineageSparkListener"
# dbt: 看 dbt run output 末尾

# 3. 测试 OpenLineage 端点
curl -X POST http://localhost:5000/api/v1/lineage \
  -H "Content-Type: application/json" \
  -d '{
    "eventType": "COMPLETE",
    "eventTime": "2026-08-17T00:00:00Z",
    "run": {"runId": "test-1"},
    "job": {"namespace": "test", "name": "manual-test"},
    "inputs": [],
    "outputs": [],
    "producer": "manual-test",
    "schemaURL": "https://openlineage.io/spec/2-0-2/OpenLineage.json"
  }'

如果这个 curl 能在 Marquez 看到 test.manual-test job,说明后端 OK,问题在客户端。


9. 总结

9.1 实战核心步骤

步骤 耗时 产出
1. 部署 Marquez 5 分钟 后端 + Web UI + PostgreSQL
2. Airflow 集成 10 分钟 DAG 自动捕获血缘
3. Spark 集成 10 分钟 Spark job 自动捕获血缘
4. dbt 集成 5 分钟 dbt 模型自动捕获血缘
5. 联合验证 5 分钟 完整血缘图可视化

总计 35 分钟,5 步拥有端到端血缘系统。

9.2 关键技术决策

  • OpenLineage 标准 + Marquez 后端 —— LF AI 治理双保险
  • Apache-2.0 + Apache-2.0 —— 双 License 商用友好
  • Java 17 + PostgreSQL 14 —— Marquez 架构稳定
  • Spark/Airflow/dbt/Flink 全部原生 —— 无需自研集成

9.3 下一步

实战完成后,你应该考虑:

  1. 加 OpenMetadata(如果需要 catalog + 质量 + AI 上下文)
  2. 接 Apache Atlas(如果需要 Hadoop 圈 PII 自动 masking)
  3. 加 Prometheus + Grafana(Marquez 自带 /metrics 端口 5001)

📎 参考


本文数据基准 2026-08-17(Marquez 0.50.0 + OpenLineage spec 2-0-2)。如需引用: [《OpenLineage + Marquez 实战》](https://east196.cn/?p=363)

正文完