OpenLineage 与列级数据血缘深度实战:从 SQL AST 作用域解析到 Marquez 图存储的工程全解

数据血缘(Data Lineage)在数据平台里长期处于一个尴尬的位置:人人都承认它重要,真正落地到列级粒度的团队却凤毛麟角。多数公司最终做出来的东西,是一张"表 A → 表 B → 报表 C"的拓扑图,看起来热闹,一旦有人问"这个风控指标里的 user_score 字段,改动上游哪个字段会炸",图就答不上来了。

这篇文章聊的是把血缘做到列级这件事的完整工程链路:OpenLineage 的事件模型长什么样、SQL 文本如何被解析成列级依赖边、跨引擎(Spark / Flink / dbt)如何统一上报、以及 Marquez 这类后端如何把事件流固化成可查询的血缘图。

一、为什么表级血缘不够用

表级血缘的信息损失是数量级级别的。一个典型数仓的 DWD 宽表可能有 200 列、上游 15 张表,表级图只会告诉你"这 15 张表都影响它",而实际上某个指标列可能只依赖其中 2 张表的 3 个字段。

粒度差异直接决定了四类场景能不能做:

  • 影响分析(Impact Analysis):改 dim_user.kyc_level 的类型,表级图会告警下游 40 张表,列级图只会告诉你真正受影响的 3 张表 7 个字段。前者等于没告警。
  • 根因定位:报表数字异常时,列级血缘可以把排查范围从"上游全链路"收敛到"该列的实际计算路径"。
  • 合规与删除权:GDPR 的被遗忘权要求你定位到"哪些下游资产包含了这个人的手机号"。做不到列级,就只能全表扫描式地人工审计。
  • 成本归因:把计算成本按"最终被消费的列"反推分摊,而不是粗暴地按表均摊。

换句话说,表级血缘是"资产地图",列级血缘才是"可执行的依赖关系"。

二、OpenLineage 的核心模型

OpenLineage 是 LF AI 旗下的开放规范,本质是一套 JSON Schema + 传输协议(HTTP / Kafka),不绑定任何实现。理解它就四个概念:

  • Job:产生数据的过程,用一个稳定的 name + namespace 标识,可版本化。
  • Dataset:数据的逻辑视图,name + namespace + 可选的 physical name。
  • Run:Job 的一次具体执行,带 runId(UUID)、eventTime、eventType(START / RUNNING / COMPLETE / ABORT / FAIL)。
  • Facet:挂在 Job / Dataset / Run / RunEvent 上的任意扩展信息。规范内置的 schema、columnLineage、dataQualityMetrics、symlinks 都是 Facet。

一个带列级血缘的 COMPLETE 事件骨架:

{
  "eventTime": "2026-09-30T03:12:44.512Z",
  "producer": "https://github.com/OpenLineage/OpenLineage/tree/main/integration/spark",
  "schemaURL": "https://openlineage.io/spec/2-0-2/OpenLineage.json#/$defs/RunEvent",
  "eventType": "COMPLETE",
  "job": {
    "namespace": "spark://analytics",
    "name": "dwd_order_wide",
    "facets": { "jobType": { "processingType": "BATCH", "integration": "SPARK" } }
  },
  "run": { "runId": "b3f1c2a4-8d31-4f0e-9a77-2c1e5b8d6f00" },
  "inputs": [{
    "namespace": "hive://warehouse",
    "name": "ods.order_detail",
    "facets": {
      "schema": {
        "fields": [
          { "name": "order_id", "type": "BIGINT" },
          { "name": "sku_price", "type": "DECIMAL(18,2)" },
          { "name": "discount",  "type": "DECIMAL(18,2)" }
        ]
      }
    }
  }],
  "outputs": [{
    "namespace": "hive://warehouse",
    "name": "dwd.order_wide",
    "facets": {
      "schema": { "fields": [{ "name": "pay_amount", "type": "DECIMAL(18,2)" }] },
      "columnLineage": {
        "fields": {
          "pay_amount": {
            "inputFields": [
              { "namespace": "hive://warehouse", "name": "ods.order_detail", "field": "sku_price" },
              { "namespace": "hive://warehouse", "name": "ods.order_detail", "field": "discount" }
            ],
            "transformationType": "TRANSFORMATION",
            "transformationDescription": "sku_price * (1 - discount)"
          }
        }
      }
    }
  }]
}

注意几个工程细节:runId 必须幂等可重放(后端靠它去重);eventTime 必须是 UTC 且带毫秒;columnLevelLineage 挂在 output 的 facet 上,因为语义是"这个输出列由哪些输入列算出"。

三、列级血缘的真正难点:SQL 作用域解析

生成 columnLineage 只有两条路:解析 SQL 文本,或从执行引擎的物理计划里读(Spark 的 QueryExecutionListener 能拿到 optimized logical plan,Flink 有 JobListener + calcite relNode)。生产上两条路都会走,但 SQL 文本解析是兜底能力——因为视图定义、dbt model、历史脚本你都得能离线解析。

难点不在语法树,而在作用域(Scope)与列解析(Column Resolution)。

3.1 作用域与别名消解

一段 SQL 里,SELECT a + b FROM t1 JOIN t2 USING(id) 里的 a 到底来自哪张表?要正确回答,必须构建一个作用域链:

  1. 每个 SELECT 块是一个 scope,内部有 sources(FROM 子句产出的表/子查询)和 columns(输出列)。
  2. USING / NATURAL JOIN 会引入合并列(coalesced column),它的血缘是参与合并的多个列的并集。
  3. 别名(AS)会遮蔽(shadow)真实列名,输出列的血缘必须回溯到别名覆盖前的表达式。
  4. 关联子查询会引用外部作用域的列,这时候要沿 scope 链向上查找。

下面用 Python + sqlglot 写一个最小可用的列级血缘提取器(生产版在此基础上加方言适配和容错):

import sqlglot
from sqlglot import exp

def build_scopes(sql: str, dialect="hive"):
    """返回 {scope_id: {'sources': {alias: (db, table)}, 'outputs': {col: [ (db, table, col) ]}}}"""
    tree = sqlglot.parse_one(sql, read=dialect)
    result = {}

    for i, select in enumerate(tree.find_all(exp.Select)):
        # 1) FROM 子句 -> 物理表映射(含别名)
        sources = {}
        for tbl in select.find_all(exp.Table):
            alias = tbl.alias_or_name
            sources[alias] = (tbl.db or "default", tbl.name)

        # 2) 输出列 -> 依赖输入列
        outputs = {}
        for proj in select.expressions:
            out_name = proj.alias_or_name
            deps = []
            for col in proj.find_all(exp.Column):
                qualifier = col.table or _resolve_qualifier(col, sources, select)
                src = sources.get(qualifier)
                if src is None:
                    # 关联子查询:回溯外层作用域
                    deps.append(("__outer__", "__outer__", col.name))
                else:
                    deps.append((src[0], src[1], col.name))
            outputs[out_name] = deps
        result[f"scope_{i}"] = {"sources": sources, "outputs": outputs}
    return result

def _resolve_qualifier(col, sources, select):
    """只有一个 source 时,未限定的列直接归属它"""
    if len(sources) == 1:
        return next(iter(sources))
    return None

这个 60 行版本能覆盖 70% 场景。剩下 30% 才是工程量的大头:

  • **SELECT * 展开**:必须查 metastore 拿到源表 schema 才能展开成具体列,否则血缘出现 * 节点,下游全部断链。这要求解析器和 catalog 强耦合。
  • CTE 与视图:WITH x AS (...) 是一个具名 scope,引用处要做视图展开;嵌套视图要递归展开并做环检测(递归视图或同名临时表很容易炸栈)。
  • UNION / UNION ALL:按位置对齐,输出列第 i 位的血缘是各分支第 i 位的并集。
  • JOIN 合并列:USING(k) 的 k 血缘同时指向左右两表。
  • 聚合与窗口函数:SUM(x) 的血缘是 x,但语义上带"聚合"标记(TRANSFORMATION_AGGREGATION),下游做影响分析时要区分"值传递"和"值变换"。

3.2 降级策略比算法更重要

真实的数仓里有大量解析器搞不定的 SQL:UDF 黑盒、动态 SQL、存储过程、INSERT OVERWRITE 混着 SELECT *。工程上必须有明确的降级路径:

列级解析成功        -> 输出 columnLineage facet
仅解析出输入输出表   -> 输出 dataset 级 input/output,schema facet 留空
SQL 获取失败        -> 从引擎侧(Spark listener)拿物理计划兜底
全部失败            -> 上报一个 ABORT/FAIL 事件,绝不静默丢弃

静默丢弃是血缘系统死亡的第一原因。一条链路漏了,用户看到的就是一张错误的图,信任一旦崩塌,这系统就没人用了。我的做法是给每个 run 强制落一条审计记录,解析覆盖率做成dashboard 指标,低于阈值直接告警。

四、跨引擎统一上报

OpenLineage 的价值在于它把"N 个引擎 × M 个血缘后端"的 M×N 问题变成了 N+M。

  • Spark:官方 openlineage-spark 用 SparkListenerSQLExecutionStart/End 钩子,从 QueryExecution 里读逻辑计划和 output 表,直接生成 columnLineage。这是最省事的一条路,因为它绕过了 SQL 文本解析——物理计划里已经没有别名歧义了。
  • Flink:流式作业没有"一次执行"的概念。做法是长周期作业发 START,然后按 checkpoint 周期发 RUNNING 事件,把 Kafka topic 作为 input/output dataset,用 facets.symlinks 把 topic 名映射到逻辑 dataset。
  • dbt:dbt 的 manifest.json 里已经有完整的 depends_on 和列信息,写一个脚本在 dbt run 后把 manifest 转成 OpenLineage 事件是最划算的投入——几乎零解析器工作量,就能拿到全公司最核心的那批 model 的列级血缘。
  • Airflow:把 lineage backend 配成 OpenLineage,每个 task 实例天然就是一个 run,runId 用 dag_id + task_id + run_id 哈希生成,保证重试幂等。

统一的关键点是 namespace 规范。建议提前定死:kafka://cluster、hive://warehouse、s3://bucket/prefix、postgres://host:port/db。namespace 命名不一致会导致同一份数据在血缘图里分裂成两个节点,这是后期最难修的数据质量问题。

五、Marquez 后端:事件流如何变成图

Marquez 是 OpenLineage 的参考实现,存储层是 PostgreSQL。它的模型设计有几个值得学的点:

Dataset 版本化。Dataset 不是静态节点,而是由 (namespace, name) 标识的逻辑实体 + 一系列 dataset_versions(UUID)。每次 run 产生了新 schema,就开一个新版本。查询血缘时可以指定"当前版本"或"某历史时间点版本"——这对"当时那个报表是怎么算出来的"这类审计问题至关重要。

Job 版本化同理。Job 的 location(源码路径)+ schema 的哈希变化会触发新版本,于是"这个指标上周和这周算得不一样"可以从血缘图直接定位到代码变更。

图遍历。血缘查询本质是有向图遍历,Marquez 用递归 CTE 做:

WITH RECURSIVE lineage AS (
  SELECT d.uuid, d.name, 1 AS depth
  FROM datasets d WHERE d.name = 'dwd.order_wide'
  UNION ALL
  SELECT i.uuid, i.name, l.depth + 1
  FROM lineage l
  JOIN dataset_versions_io io ON io.output_dataset_version_uuid = l.uuid
  JOIN datasets i ON i.uuid = io.input_dataset_uuid
  WHERE l.depth < 10
)
SELECT DISTINCT name, min(depth) FROM lineage GROUP BY name;

depth < 10 这个硬截断是必须的——数仓里交叉依赖形成的环会导致无限递归。生产上我还会加一层:查询超时 + 结果集上限,超出的部分返回"已截断"标记,前端明确告知用户图不完整,而不是假装完整。

乱序与幂等。事件经 Kafka 到达时可能乱序(COMPLETE 早于 START 到达)。Marquez 的做法是按 runId 做 upsert,用 eventTime 做行级版本控制,后到的同 runId 事件只在其 eventTime 更新时才覆盖。

六、生产落地的几个硬观点

  1. 先做读取侧,别先做全链路。从 BI 报表和核心指标的入口反向追溯,覆盖"最常被问到的 20 个字段"带来的价值,比全量铺开高得多。
  2. 血缘质量要量化。定义"解析覆盖率 = 有 columnLineage 的 output / 总 output",低于 85% 就别对外宣传,因为错误的血缘比没有血缘更危险。
  3. PII 标签要随血缘传播。把 PII: true 做成 column facet,列级血缘能自动算出"哪些下游资产含 PII",这是合规审计最直接的收益,也是最容易争取到预算的场景。
  4. 控制 dataset 基数。每个 run 都开新 dataset version 会让图爆炸。生产上要按 schema 哈希去重,schema 未变则复用旧版本,否则半年后你的血缘表会有几千万行版本记录。
  5. 把它接进 CI。dbt / SQL 脚本的 PR 阶段就跑一次血缘解析,把"影响到的下游资产"作为 PR comment 输出,血缘就从"事后查询工具"变成了"事前防御工具"——这才是它真正产生价值的位置。

小结

列级血缘的技术门槛不在图数据库、不在前端可视化,而在 SQL 作用域解析的覆盖率 和 namespace 治理 这两件脏活上。OpenLineage 提供的价值是把事件模型标准化,让你不必重复造轮子;Marquez 提供的是一套可参考的版本化图存储设计。剩下 80% 的工作量,是你公司那些祖传 SQL 和没写注释的 UDF——对此没有捷径,只能靠降级策略 + 覆盖率指标,一寸一寸地把可信度啃上去。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部