元数据到底是什么

元数据就是"关于数据的数据"。听起来像绕口令,但道理很简单:你的 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%?”

没有数据血缘时,排查过程是这样的:

  1. 找到日报的 SQL 脚本(可能在某个 Git 仓库里)
  2. 查看 SQL 里引用了哪些表(从几十行 SQL 中手动解析)
  3. 追踪这些表的上游 ETL 任务(去 Airflow 里一个个找)
  4. 检查每个任务的执行状态(翻日志)
  5. 定位到某个任务失败了,或者某个上游数据源出问题了

这一步排查下来,快的半小时,慢的要半天。有了数据血缘,你只需要点几下鼠标就可以看到完整的数据链路。

血缘的基本单元

一条血缘记录包含三个要素:

原始字段 → 处理过程 → 目标字段
  ↓           ↓          ↓
 来源       转换       去向

最简单的血缘示例如下:

{
  "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: 元数据三层分类与数据血缘核心概念,从仓库设计到工具选型与影响分析实践。