6 minutes
元数据管理与数据血缘
元数据到底是什么
元数据就是"关于数据的数据"。听起来像绕口令,但道理很简单:你的 ETL 管道每天处理几百张表、几千个字段、几十个任务。当有人问"这张表的数据从哪来的"或者"这个字段的转换规则是什么"时,你需要一个系统来回答这些问题。这个系统就是元数据管理。
元数据分为三个层次:
技术元数据
技术元数据描述数据处理的技术细节:
- 数据库连接信息
- 表的 schema(字段名、类型、约束)
- 分区和文件路径
- ETL 任务的调度配置
- 数据量和大小
# 典型的表级技术元数据
table_metadata = {
"table_name": "fact_orders",
"database": "dw_prod",
"schema": "public",
"columns": [
{"name": "order_id", "type": "BIGINT", "nullable": False, "pk": True},
{"name": "user_id", "type": "BIGINT", "nullable": False, "fk": "dim_users"},
{"name": "amount", "type": "DECIMAL(12,2)", "nullable": False},
{"name": "status", "type": "VARCHAR(20)", "nullable": False},
{"name": "created_at", "type": "TIMESTAMP", "nullable": False},
],
"partitions": ["dt"],
"format": "parquet",
"location": "s3://dw-prod/orders/",
"row_count": 50000000,
"last_refreshed": "2026-04-19 03:15:00",
}
业务元数据
业务元数据回答"这是什么意思"的问题:
- 字段的业务定义(“什么是活跃用户”)
- 数据的所有者和责任人
- 数据的使用限制和安全等级
- 计算口径和业务规则
例如,“订单金额"这个字段在技术元数据里是一个 DECIMAL(12,2),但在业务元数据里需要写明:这是用户下单时实际支付的金额(不含运费),单位是人民币元,计算口径是 sum(商品单价 × 数量 - 优惠金额)。
操作元数据
操作元数据记录数据处理过程本身的信息:
- ETL 任务的执行历史
- 每次运行的耗时和资源消耗
- 数据质量检查结果
- 错误日志和故障记录
# 典型的任务执行元数据
run_metadata = {
"task_id": "daily_order_etl",
"dag_id": "order_pipeline",
"execution_date": "2026-04-19",
"start_time": "2026-04-19 03:00:00",
"end_time": "2026-04-19 03:25:30",
"status": "success",
"records_processed": 125000,
"records_inserted": 120000,
"records_updated": 5000,
"records_rejected": 150,
"error_rate_pct": 0.12,
"peak_memory_mb": 2048,
"peak_cpu_pct": 85,
"retry_attempts": 0,
}
元数据仓库设计
有了元数据的概念,接下来要解决的是"存在哪里”。元数据仓库就是所有元数据的集中存储。
核心表设计
一个最简单的元数据仓库只需要几张表:
-- 数据源表
CREATE TABLE metadata_data_sources (
source_id SERIAL PRIMARY KEY,
source_name VARCHAR(255) NOT NULL,
source_type VARCHAR(50) NOT NULL, -- MySQL, PostgreSQL, S3, API
connection_info TEXT,
owner VARCHAR(100),
created_at TIMESTAMP DEFAULT NOW(),
updated_at TIMESTAMP DEFAULT NOW()
);
-- 表/数据集元数据
CREATE TABLE metadata_tables (
table_id SERIAL PRIMARY KEY,
source_id INT REFERENCES metadata_data_sources(source_id),
table_name VARCHAR(255) NOT NULL,
table_type VARCHAR(50) DEFAULT 'table', -- table, view, external, temp
description TEXT,
row_count BIGINT,
size_bytes BIGINT,
last_analyzed TIMESTAMP,
created_at TIMESTAMP DEFAULT NOW()
);
-- 字段元数据
CREATE TABLE metadata_columns (
column_id SERIAL PRIMARY KEY,
table_id INT REFERENCES metadata_tables(table_id),
column_name VARCHAR(255) NOT NULL,
data_type VARCHAR(100),
is_nullable BOOLEAN DEFAULT TRUE,
is_primary_key BOOLEAN DEFAULT FALSE,
default_value TEXT,
business_description TEXT,
created_at TIMESTAMP DEFAULT NOW()
);
-- ETL 任务元数据
CREATE TABLE metadata_etl_tasks (
task_id SERIAL PRIMARY KEY,
task_name VARCHAR(255) NOT NULL,
source_table_id INT REFERENCES metadata_tables(table_id),
target_table_id INT REFERENCES metadata_tables(table_id),
transformation_logic TEXT,
schedule_expression VARCHAR(100),
owner VARCHAR(100),
created_at TIMESTAMP DEFAULT NOW()
);
-- 任务执行记录
CREATE TABLE metadata_execution_logs (
log_id SERIAL PRIMARY KEY,
task_id INT REFERENCES metadata_etl_tasks(task_id),
execution_date DATE NOT NULL,
start_time TIMESTAMP,
end_time TIMESTAMP,
status VARCHAR(20), -- running, success, failed, retrying
rows_processed INT,
error_message TEXT,
duration_seconds INT
);
元数据采集
元数据不应该手动录入,而应该自动采集。采集方式有几种:
主动扫描型:定期连接数据库,读取 information_schema,获取表和字段信息。
def scan_database_metadata(conn, source_name):
"""扫描数据库获取元数据"""
metadata = {
"source": source_name,
"tables": [],
}
cursor = conn.cursor()
# 获取所有表
cursor.execute("""
SELECT table_name, table_type, pg_size_pretty(pg_total_relation_size(quote_ident(table_name)))
FROM information_schema.tables
WHERE table_schema = 'public'
ORDER BY table_name
""")
for table_name, table_type, size in cursor.fetchall():
table_info = {
"name": table_name,
"type": table_type,
"size": size,
"columns": [],
}
# 获取字段信息
cursor.execute(f"""
SELECT column_name, data_type, is_nullable, column_default,
character_maximum_length
FROM information_schema.columns
WHERE table_schema = 'public' AND table_name = '{table_name}'
ORDER BY ordinal_position
""")
for col_name, col_type, nullable, default, char_len in cursor.fetchall():
table_info["columns"].append({
"name": col_name,
"type": col_type,
"nullable": nullable == "YES",
"default": default,
"max_length": char_len,
})
metadata["tables"].append(table_info)
return metadata
事件捕获型:在 ETL 任务执行时自动上报元数据,比如 Airflow 的日志本身就包含了执行元数据。
解析推断型:从 SQL 脚本、Python 代码中解析出数据流信息。这是最难的,因为代码表达方式千差万别。
数据血缘
数据血缘(Data Lineage)是元数据管理的核心应用。它回答一个简单的问题:某份数据从哪里来,经过了什么处理,最终用在了哪里。
为什么需要数据血缘
假设业务方问了一个问题:“昨天的日报里,GMV 为什么比前一天少了 20%?”
没有数据血缘时,排查过程是这样的:
- 找到日报的 SQL 脚本(可能在某个 Git 仓库里)
- 查看 SQL 里引用了哪些表(从几十行 SQL 中手动解析)
- 追踪这些表的上游 ETL 任务(去 Airflow 里一个个找)
- 检查每个任务的执行状态(翻日志)
- 定位到某个任务失败了,或者某个上游数据源出问题了
这一步排查下来,快的半小时,慢的要半天。有了数据血缘,你只需要点几下鼠标就可以看到完整的数据链路。
血缘的基本单元
一条血缘记录包含三个要素:
原始字段 → 处理过程 → 目标字段
↓ ↓ ↓
来源 转换 去向
最简单的血缘示例如下:
{
"lineage": [
{
"source": {
"database": "oltp_prod",
"table": "orders",
"column": "amount"
},
"transformation": "SUM(amount) - SUM(discount)",
"target": {
"database": "dw_prod",
"table": "fact_daily_sales",
"column": "net_revenue"
}
}
]
}
构建血缘的三种方式
手动维护:由数据工程师在开发时记录。比如在 dbt 中,模型之间的引用关系就是手动声明的。
# dbt 模型中的血缘声明
version: 2
models:
- name: fact_daily_sales
description: "每日销售事实表"
columns:
- name: net_revenue
description: "净收入 (销售额 - 折扣)"
depends_on:
- staging_orders
- dim_products
手动维护的好处是准确,但问题也明显——容易遗漏,而且开发人员不一定会主动记录。
基于解析自动构建:通过解析 SQL 和 Python 代码,自动推断出数据流向。
import re
import sqlparse
def extract_table_lineage(sql_query):
"""从 SQL 查询中提取数据血缘"""
parsed = sqlparse.parse(sql_query)[0]
from_tables = []
insert_table = None
for token in parsed.tokens:
# 提取 INSERT 目标表
if token.ttype is None and isinstance(token, sqlparse.sql.Identifier):
if any(p.value == "INSERT" for p in parsed.tokens):
insert_table = str(token)
# 提取 FROM 中的源表
if token.ttype is None and isinstance(token, sqlparse.sql.Where):
continue
# 使用正则查找表名
from_pattern = r'FROM\s+(\w+)'
join_pattern = r'JOIN\s+(\w+)'
from_matches = re.findall(from_pattern, sql_query, re.IGNORECASE)
join_matches = re.findall(join_pattern, sql_query, re.IGNORECASE)
# 查找 INSERT 或 CREATE TABLE 的目标表
insert_pattern = r'INSERT\s+INTO\s+(\w+)'
create_pattern = r'CREATE\s+TABLE\s+(\w+)'
insert_match = re.search(insert_pattern, sql_query, re.IGNORECASE)
create_match = re.search(create_pattern, sql_query, re.IGNORECASE)
target = (insert_match or create_match).group(1) if (insert_match or create_match) else None
return {
"target": target,
"sources": from_matches + join_matches,
"query": sql_query[:100] + "..."
}
# 使用示例
sql = """
INSERT INTO dw.fact_daily_sales
SELECT o.order_id, o.amount, p.category, u.region
FROM staging.orders o
JOIN staging.products p ON o.product_id = p.product_id
JOIN staging.users u ON o.user_id = u.user_id
WHERE o.created_at >= CURRENT_DATE - 1
"""
lineage = extract_table_lineage(sql)
print(f"目标表: {lineage['target']}")
print(f"源表: {lineage['sources']}")
解析方式的难点在于 SQL 的复杂性(CTE、子查询、动态 SQL)和 Python 的动态特性。
运行时捕获:在 ETL 框架的执行过程中自动记录。这是最可靠的方案。
Airflow 中可以通过自定义回调来记录运行时血缘:
from airflow.models import BaseOperator
from airflow.utils.decorators import apply_defaults
class LineageAwareOperator(BaseOperator):
"""带血缘感知的自定义算子"""
@apply_defaults
def __init__(self, source_tables=None, target_tables=None, *args, **kwargs):
super().__init__(*args, **kwargs)
self.source_tables = source_tables or []
self.target_tables = target_tables or []
def execute(self, context):
# 记录血缘信息
lineage_record = {
"dag_id": context["dag"].dag_id,
"task_id": self.task_id,
"execution_date": str(context["execution_date"]),
"sources": self.source_tables,
"targets": self.target_tables,
"start_time": datetime.now().isoformat(),
}
try:
# 实际执行逻辑
result = self._execute_logic(context)
lineage_record["status"] = "success"
lineage_record["end_time"] = datetime.now().isoformat()
lineage_record["records_affected"] = len(result) if result else 0
return result
except Exception as e:
lineage_record["status"] = "failed"
lineage_record["error"] = str(e)
raise
finally:
# 写入血缘记录
self._write_lineage(lineage_record)
OpenLineage 标准
OpenLineage 是一个开源标准,定义了数据血缘的通用格式。它的目标是让不同工具之间能够交换血缘信息。
一个 OpenLineage 格式的运行事件示例:
{
"eventType": "COMPLETE",
"eventTime": "2026-04-19T03:15:00Z",
"run": {
"runId": "d478e8a2-9f1c-4a5e-b2d4-8c3a6f9b0e21"
},
"job": {
"namespace": "airflow",
"name": "order_pipeline.daily_order_etl"
},
"inputs": [
{
"namespace": "postgresql://oltp:5432",
"name": "ecommerce.orders",
"facets": {
"schema": {
"fields": [
{"name": "order_id", "type": "int8"},
{"name": "amount", "type": "numeric"},
{"name": "status", "type": "varchar"}
]
}
}
}
],
"outputs": [
{
"namespace": "s3://dw-prod",
"name": "fact_orders",
"facets": {
"schema": {
"fields": [
{"name": "order_id", "type": "bigint"},
{"name": "net_amount", "type": "decimal"},
{"name": "order_date", "type": "date"}
]
}
}
}
]
}
OpenLineage 的优点是标准化,支持它的工具(Airflow、dbt、Spark、Flink)可以互相共享血缘数据。Marquez 是 OpenLineage 的一个参考实现,提供了血缘数据的存储和可视化。
元数据工具对比
| 工具 | 定位 | 部署 | 血缘支持 | 数据目录 | 学习成本 |
|---|---|---|---|---|---|
| Apache Atlas | 企业级数据治理 | 重(Hadoop 生态) | 深度支持 | 是 | 高 |
| DataHub | 数据发现平台 | 中(Kafka + Elasticsearch) | 强 | 是 | 中 |
| Amundsen | 数据发现 | 中 | 基础 | 是 | 中 |
| Marquez | 血缘追踪 | 轻 | 专注血缘 | 否 | 低 |
| dbt docs | 转换层文档 | 极轻 | dbt 模型内 | 部分 | 低 |
选型建议:
- 如果团队刚起步,先做好 dbt docs 或简单的元数据表收集,不要过早引入重量级工具
- 如果已经有 Hadoop/Spark 生态,Atlas 是自然的选择
- 如果需要企业级数据发现和目录,DataHub 是目前社区最活跃的选择
- 如果只需要血缘追踪,Marquez 轻量且专注
数据目录的概念
数据目录(Data Catalog)是元数据管理的上层应用。它解决一个实际痛点:业务分析师不知道公司有哪些数据可用,也不知道这些数据什么意思。
一个好的数据目录应该提供:
- 数据搜索:根据表名、字段名、标签搜索数据集
- 数据预览:查看数据样例,不需要写 SQL
- 字段说明:每个字段的业务含义
- 使用统计:哪些表被频繁查询,谁是数据的热点
- 所有者标注:有问题时知道找谁
class SimpleDataCatalog:
"""极简数据目录实现"""
def __init__(self):
self.entries = {}
def register_table(self, table_name, description, owner, tags):
"""注册表到目录"""
self.entries[table_name] = {
"description": description,
"owner": owner,
"tags": tags,
"columns": [],
"popularity": 0,
}
def search(self, keyword):
"""搜索数据目录"""
results = []
keyword = keyword.lower()
for name, meta in self.entries.items():
score = 0
if keyword in name.lower():
score += 10
if keyword in meta["description"].lower():
score += 5
for tag in meta["tags"]:
if keyword in tag.lower():
score += 3
if score > 0:
results.append((name, meta, score))
results.sort(key=lambda x: -x[2])
return [(name, meta) for name, meta, _ in results]
影响分析与根因溯源
元数据管理最实在的两个应用场景就是影响分析和根因溯源。
影响分析(Impact Analysis)
当你计划修改一个数据源的 schema 时,需要知道下游有哪些系统会受到影响。这就是影响分析。
要修改: oltp.orders.amount (字段类型 INT → DECIMAL)
影响链:
└─ staging.orders (直接依赖)
└─ dw.fact_orders (计算: amount * 1.13)
└─ dw.daily_sales_summary (聚合计算)
└─ bi报表: 每日销售看板
└─ 管理层日报邮件
一个简单的血缘查询来实现影响分析:
def impact_analysis(column_identifier, lineage_graph, max_depth=5):
"""影响分析:如果修改某字段,哪些下游会受影响"""
impacted = []
queue = [(column_identifier, 0)]
visited = set()
while queue:
current, depth = queue.pop(0)
if current in visited or depth > max_depth:
continue
visited.add(current)
for downstream in lineage_graph.get_downstream(current):
impacted.append({
"target": downstream,
"hops": depth + 1,
"path": f"{current} → {downstream}"
})
queue.append((downstream, depth + 1))
return impacted
根因溯源(Root Cause Tracing)
当日报的数字不对时,从最终报表反向追溯到数据源,找到问题所在。
报表: daily_sales (GMV = 1.2M)
└─ dw.daily_sales_summary (最后更新: 04:30)
└─ dw.fact_orders (最后更新: 04:15)
└─ staging.orders (任务: daily_order_etl, 状态: FAILED)
└─ oltp.orders (源表正常)
└─ 发现问题: staging.orders 的 ETL 任务因为磁盘满而失败
小结
元数据管理是数据工程从"能用"走向"好用"的关键一步。技术元数据帮你管理数据资产,业务元数据让数据能被理解,操作元数据让数据处理过程可追溯。数据血缘作为元数据管理的核心应用,解决了数据从哪里来、经过什么处理、用到哪里去的问题。OpenLineage 标准化了血缘数据的交换格式。在实际建设中,建议先从简单的元数据收集开始,逐步引入专用的元数据管理工具,最终建立起覆盖全链路的数据治理体系。
Summary: 元数据三层分类与数据血缘核心概念,从仓库设计到工具选型与影响分析实践。