
飞熊实战向 · 数据更新于 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 下一步
实战完成后,你应该考虑:
- 加 OpenMetadata(如果需要 catalog + 质量 + AI 上下文)
- 接 Apache Atlas(如果需要 Hadoop 圈 PII 自动 masking)
- 加 Prometheus + Grafana(Marquez 自带
/metrics端口 5001)
📎 参考
- Marquez 部署: github.com/MarquezProject/marquez
- OpenLineage 集成: openlineage.io/docs/integrations
- 已有 WP 调研: 《OpenLineage 调研》 · 《Marquez 调研》 · 《数据治理横评》
本文数据基准 2026-08-17(Marquez 0.50.0 + OpenLineage spec 2-0-2)。如需引用: [《OpenLineage + Marquez 实战》](https://east196.cn/?p=363)