OpenMetadata 深度实战:从 JSON Schema 实体模型、Ingestion Framework 到血缘图与数据质量的元数据平台工程全解

当一家公司的数据栈同时跑着 Snowflake、Trino、Airflow、dbt 和三套自研服务时,真正稀缺的从来不是"能不能查到数据",而是"敢不敢信这份数据"。OpenMetadata 想解决的就是这个信任问题——把所有系统的元数据收敛到一个带血缘、带质量信号、带所有权的图里。这不是一个 UI 工程,而是一套完整的元数据基础设施。

一、元数据平台的本质难题

数据目录(Data Catalog)听起来像是"给表加个描述",但落地时你会发现三个绕不开的坑:

  1. Schema 熵增:每个系统的元数据形状都不同。Snowflake 的 COLUMNS、Trino 的 information_schema、Airflow 的 DAG 序列化对象、dbt 的 manifest.json——它们连"一张表"的定义都不一致。
  2. 新鲜度与规模的矛盾:全量扫描十万张表意味着数小时的连接风暴;增量又要求你准确识别"谁变了"。
  3. 元数据本身的质量:抽上来一堆没有 owner、没有描述、名字叫 tmp_v2_final_final 的表,目录只会变成垃圾场。

OpenMetadata 的解法是把这三点拆成三个正交的子系统:类型系统(JSON Schema 实体模型)、管道框架(Ingestion Framework)、服务层(血缘图 + 质量信号)。下面逐个拆。

二、类型系统:JSON Schema 就是契约

OpenMetadata 最有争议也最聪明的设计决策是:所有实体都用 JSON Schema 定义,代码从 Schema 生成。

openmetadata-spec/src/main/resources/json/schema/
├── entity/
│   ├── data/table.json          # Table 实体
│   ├── data/databaseSchema.json
│   ├── services/databaseService.json
│   └── ...
├── type/                        # 可复用类型定义
│   ├── entityReference.json
│   ├── tagLabel.json
│   └── changeEvent.json
└── api/                         # API 请求/响应契约

一张表的 Schema 片段(简化):

{
  "$id": "https://open-metadata.org/schema/entity/data/table.json",
  "title": "Table",
  "type": "object",
  "javaType": "org.openmetadata.schema.entity.data.Table",
  "properties": {
    "id":            { "$ref": "../../type/basic.json#/definitions/uuid" },
    "name":          { "$ref": "../../type/basic.json#/definitions/entityName" },
    "fullyQualifiedName": { "$ref": "../../type/basic.json#/definitions/fullyQualifiedName" },
    "columns": {
      "type": "array",
      "items": { "$ref": "../../type/entity/data/column.json" }
    },
    "tableConstraints": {
      "type": "array",
      "items": { "$ref": "../../type/entity/data/tableConstraint.json" }
    },
    "tableType": { "enum": ["Regular","External","View","SecureView","Iceberg"] },
    "profile":   { "$ref": "table.json#/definitions/tableProfile" },
    "extension": { "type": "object" }
  },
  "required": ["id", "name", "columns"]
}

几个关键点值得注意:

  • fullyQualifiedName(FQN)是全局主键。所有跨实体引用都通过 service.db.schema.table.column 这种路径式的 FQN,而不是 UUID。这让血缘合并天然幂等——两个系统上报同一个 FQN,自动对齐。
  • extension 字段是逃生舱。任何 Schema 没覆盖的业务元数据(比如"这张表归属于哪个成本中心")可以塞进 extension,不需要改 spec、不需要重新生成代码。
  • javaType 注解驱动代码生成,同时生成 Python Pydantic 模型。这意味着 Python 侧的 ingestion 代码与 Java 侧的服务端代码共享同一份契约,接口漂移在编译期就被捕获。

这套设计让"加一个新实体类型"的成本从"改 5 个服务 + 3 个 SDK"降到"加一个 JSON 文件"。代价是 Schema 演进必须严格遵守向后兼容(不能删字段、不能改类型),这也是项目里 Deprecate 注解满天飞的原因。

三、Ingestion Framework:元数据抽取的编译器

OpenMetadata 的 ingestion 不是一堆脚本,而是一个带中间表示(IR)的管道编译器。核心抽象是:

Source → [Processor] → Sink
  • Source:从外部系统读出元数据,产出 Entity 的迭代器。
  • Processor:纯函数式的变换(过滤、打标签、改写 FQN)。
  • Sink:写入 OpenMetadata 服务端(主要是 metadata sink,还有 file、fivetran 等)。

一个生产级的 Trino ingestion 配置:

source:
  type: trino
  serviceName: prod-trino
  serviceConnection:
    config:
      type: Trino
      hostPort: trino-coordinator:8080
      catalog: hive
      connectionOptions:
        protocol: https
      connectionArguments:
        verify: /etc/ssl/certs/ca.pem
  sourceConfig:
    config:
      type: DatabaseMetadata
      # 关键:增量与规模控制
      markDeletedTables: true
      includeTables: true
      includeViews: true
      schemaFilterPattern:
        excludes: ["^tmp_.*", "^staging_.*"]
      tableFilterPattern:
        includes:
          - { pattern: ".*" }
      # 采样与 profiling 分离,避免拖垮 ingestion
      tableProfilerConfig:
        fullyQualifiedName: ".*"
        profileSample: 75.0
        profileQuery: "SELECT * FROM {table} TABLESAMPLE SYSTEM (75)"
processor:
  type: auto-classification   # 自动识别 PII
  config:
    confidenceThreshold: 0.75
sink:
  type: metadata-rest
  config: {}
workflowConfig:
  openMetadataServerConfig:
    hostPort: http://openmetadata-server:8585/api
    authProvider: openmetadata
    securityConfig:
      jwtToken: ${OM_JWT_TOKEN}
    # 校验顺序很重要
    enableVersionValidation: false

3.1 自定义 Processor:把治理规则变成代码

Processor 是治理落地的最佳位置。比如"所有含 email 的字段自动打 PII 标签":

from metadata.ingestion.api.processor import Processor, ProcessorStatus
from metadata.ingestion.ometa.ometa_api import OpenMetadata
from metadata.generated.schema.entity.data.table import Column, Table
from metadata.generated.schema.type.tagLabel import TagLabel, LabelType, State, TagSource

class PIIAutoTagProcessor(Processor):
    def __init__(self, config, metadata: OpenMetadata):
        super().__init__()
        self.metadata = metadata
        self.rules = {"email": "PII.Sensitive", "phone": "PII.Sensitive"}

    def process(self, record: Table):
        changed = False
        for column in record.columns or []:
            col_name = column.name.root.lower()
            for pattern, tag_fqn in self.rules.items():
                if pattern in col_name:
                    column.tags = (column.tags or []) + [
                        TagLabel(
                            tagFQN=tag_fqn,
                            labelType=LabelType.AUTOMATED,
                            state=State.SUGGESTED,   # 建议态,等待人工确认
                            source=TagSource.CLASSIFICATION,
                        )
                    ]
                    changed = True
        return record if changed else None   # 返回 None 则不落库

    def close(self):
        self.status = ProcessorStatus()

注意 state=State.SUGGESTED。这是 OpenMetadata 治理模型里非常实用的一个设计:标签有 CONFIRMED / SUGGESTED 两种状态,自动分类产出的永远是 SUGGESTED,需要数据 owner 确认后才升级为 CONFIRMED。这避免了"自动打标污染权威元数据"的经典问题。

3.2 增量抽取的正确姿势

全量 ingestion 在万表规模下会打爆元数据数据库连接。三个必做的优化:

  1. markDeletedTables: true:靠 FQN 比对识别删除,而不是每次 TRUNCATE 重建。
  2. 按 databaseSchema 分片并行:把整个 catalog 拆成多个 workflow,用 Airflow 的 Dynamic DAG 或 K8s Job 并发跑,每个 workflow 独立失败重试。
  3. Profiling 与 Metadata 分离:profiling 是全表扫描,耗时长且容易超时。把 tableProfilerConfig 拆到独立的 workflow,用低频(每周)+ 采样率控制。

一个真实踩坑:profileSample 单位是百分比但底层走的是 TABLESAMPLE SYSTEM,它按数据块抽样而非按行。在只有几个大 Parquet 文件的分区表上,SYSTEM (75) 可能返回 100% 的数据。正确做法是配合 profileQuery 显式指定带 WHERE 的分区过滤。

四、血缘:从 FQN 到图

OpenMetadata 的血缘是边表(EntityLineage)+ 图遍历模型,不是存一个大 JSON。

血缘上报有三条路径:

  1. Push 模式(推荐):外部系统通过 OpenLineage 协议主动推。Airflow 侧只需装 provider:
pip install openlineage-airflow openmetadata-ingestion
# airflow.cfg 或环境变量
AIRFLOW__OPENLINEAGE__NAMESPACE = "prod-airflow"
AIRFLOW__OPENLINEAGE__TRANSPORT = json.dumps({
    "type": "http",
    "url": "http://openmetadata-server:8585/api/v1/openLineage",
    "auth": {"type": "api_key", "apiKey": "${OM_JWT_TOKEN}"}
})
AIRFLOW__OPENLINEAGE__DISABLED_FOR_OPERATORS = "airflow.operators.empty.EmptyOperator"

OpenMetadata 内置了 /api/v1/openLineage 端点,直接把 OpenLineage 的 RunEvent 翻译成内部的 AddLineageRequest。

  1. Pull 模式:从 dbt manifest.json、Looker、Tableau 等带静态依赖图的系统解析。
  1. SQL 解析模式:对无法提供血缘的查询型系统,用内置解析器(基于 sqlglot)从查询历史中抽取 table -> table 依赖。这条路径噪音最大,必须配合置信度阈值。

4.1 血缘合并的幂等性

核心在于 LineageDetails 里的 source 字段:

{
  "edge": {
    "fromEntity": { "id": "...", "type": "table", "fullyQualifiedName": "svc.db.raw.orders" },
    "toEntity":   { "id": "...", "type": "table", "fullyQualifiedName": "svc.db.dwd.orders_dtl" }
  },
  "lineageDetails": {
    "sqlQuery": "INSERT INTO dwd.orders_dtl SELECT * FROM raw.orders",
    "source": "DBT",
    "pipeline": { "id": "...", "type": "pipeline", "fullyQualifiedName": "svc.orders_etl" },
    "columnsLineage": [
      { "fromColumns": ["raw.orders.order_id"], "toColumns": ["dwd.orders_dtl.order_id"] }
    ]
  }
}

同一条边被多个 source 上报时,OpenMetadata 会合并 lineageDetails 而不是覆盖,并记录每个来源的贡献。这就是为什么 FQN 主键如此重要——它让去中心化的血缘上报天然收敛。

4.2 查询图

服务端提供向下/向上的 N 层遍历:

curl -s "http://om:8585/api/v1/lineage/table/name/svc.db.dwd.orders_dtl?upstreamDepth=3&downstreamDepth=3" \
  -H "Authorization: Bearer $OM_JWT" | jq '.nodes | length'

注意深度不要超过 5。图遍历走的是递归 SQL(WITH RECURSIVE),在宽扇出节点(比如一张被 500 张下游表引用的维度表)上,upstreamDepth=10 足以把 PostgreSQL 连接池打满。生产环境应该给这个接口加缓存和超时熔断。

五、数据质量:把断言变成一等公民

OpenMetadata 的数据质量模型由三层组成:

  • Test Case:单个断言("这一列的 null 值比例 < 5%")
  • Test Suite:测试集合,绑定到 Table 或 Column
  • Test Definition:可复用的测试模板,包含 SQL 表达式的占位符

定义一个自定义测试(通过 API):

POST /api/v1/dataQuality/testDefinitions
{
  "name": "columnValueMinToBeBetween",
  "entityType": "COLUMN",
  "testPlatforms": ["OpenMetadata"],
  "parameterDefinition": [
    { "name": "minValue", "dataType": "NUMBER", "required": true },
    { "name": "maxValue", "dataType": "NUMBER", "required": true }
  ]
}

然后绑定:

POST /api/v1/dataQuality/testCases
{
  "name": "orders_amount_min_range",
  "testDefinition": "svc.columnValueMinToBeBetween",
  "entityLink": "<#E::table::svc.db.dwd.orders_dtl::columns::amount>",
  "parameterValues": [
    { "name": "minValue", "value": 0 },
    { "name": "maxValue", "value": 1000000 }
  ],
  "computePassedFailedRowCount": true
}

实际执行时,OpenMetadata 会根据 entityLink 解析出目标表,自动生成并执行:

SELECT MIN(amount) FROM svc.db.dwd.orders_dtl;

结果写入 testCaseResult,并驱动 UI 上的红绿灯、以及可选的告警 Webhook。

工程建议:测试执行一定要走独立的 workflow 和独立的数据库凭据(只读 + 资源隔离)。质量测试和 profiling 一样是全表扫描型负载,混进主 ingestion 会拖垮整个管道。同时给每个 test case 设 computePassedFailedRowCount: true,失败行数比单纯的成功/失败信号有用得多。

六、存储与检索:为什么是"关系库 + Elasticsearch"双写

OpenMetadata 的持久化架构是分层的:

MySQL / PostgreSQL  ← 权威存储(实体、关系、血缘边)
        ↓ (ChangeEvent,CDC 式事件)
Elasticsearch       ← 检索索引(全文、分面、自动补全)
        ↓
    UI / Search API

所有实体变更都会发出 ChangeEvent,由 EventHandler 消费后异步刷新 ES 索引。这个设计带来一个运维事实:ES 索引损坏不会丢数据,可以靠 reindex 从关系库全量重建:

./bootstrap/bootstrap_storage.sh reindex   # 触发全量重建

生产上三个必须关注的点:

  1. PostgreSQL 的 jsonb 性能:OpenMetadata 大量使用 jsonb 存实体。务必给 jsonb 字段加 GIN 索引,且 PostgreSQL 版本不低于 12。
  2. ES 索引分片数:table_search_index 在十万级表量下默认 1 shard 会变成写瓶颈,建议按 5-10 shard 规划。
  3. 事件积压:EventHandler 用的是内存队列(可配 Kafka)。默认内存队列在服务重启时会丢事件,导致 ES 与关系库短暂不一致——生产环境建议直接切 Kafka。

七、落地建议与踩坑清单

  • 先做 ingestion 覆盖率,再做治理。没人会给一个只有 30% 表覆盖率的目录补描述。用 ometa CLI 或自建监控盯住"已纳管表数 / 总表数"。
  • Owner 是治理的原子单位。没有 owner 的表,标签、描述、质量告警全都没有落点。建议从 Airflow DAG owner 或 dbt meta.owner 自动回填。
  • Ingestion 失败要分级。连接超时是 P2,Schema 反序列化失败是 P1(说明上游改了结构),FQN 冲突是 P0(说明有重名实体,会污染血缘图)。
  • 不要指望自动血缘覆盖一切。存储过程、动态 SQL、跨系统的文件交换(比如 Spark 写 S3 后 Snowflake 读)都是静态解析的盲区,必须靠显式 API 上报补齐。
  • 版本校验关掉。生产环境多客户端版本混用时,enableVersionValidation: true 会让老版本 ingestion 直接失败。

八、总结

OpenMetadata 的工程价值不在于"又一个目录 UI",而在于它把元数据基础设施里最脏的三件事——跨系统类型统一、可扩展抽取管道、血缘图收敛——用一套自洽的抽象做掉了。JSON Schema 驱动的类型系统解决了熵增,Processor 抽象把治理规则变成可测试的代码,FQN 主键让去中心化的血缘上报天然幂等。

如果你的数据栈已经超过三个引擎、两个编排系统,那么元数据平台不再是"锦上添花",而是控制复杂度的必要成本。而选型的判断标准很简单:看它的扩展点是靠改代码还是靠加配置。OpenMetadata 属于后者。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部