返回博客

双层数据血缘追踪:Schema 级与实例级的统一设计

coomia-dip 实现了双层数据血缘追踪系统——Schema 级血缘记录对象类型和属性之间的定义依赖关系,实例级血缘记录具体数据记录的来源和变换链路。两层血缘通过统一的 LineageGraph 模型管理,支持正向影响分析("这个字段变了会影响什么")和反向溯源("这个值是怎么来的")。本文从双层设计理念、图模型实现、血缘采集、查询引擎到可视化集成,完整阐述这一企业级数据血缘能力。

Coomia发布于 2025年9月23日10 分钟阅读
分享本文Twitter / X

系列:S6 平台工程 · 第 9 篇 | 难度:高级 | 阅读时间:18 分钟

双层数据血缘追踪:Schema 级与实例级的统一设计

#TL;DR

coomia-dip 实现了双层数据血缘追踪系统——Schema 级血缘记录对象类型和属性之间的定义依赖关系,实例级血缘记录具体数据记录的来源和变换链路。两层血缘通过统一的 LineageGraph 模型管理,支持正向影响分析("这个字段变了会影响什么")和反向溯源("这个值是怎么来的")。本文从双层设计理念、图模型实现、血缘采集、查询引擎到可视化集成,完整阐述这一企业级数据血缘能力。

#1. 数据血缘的核心价值

#1.1 为什么需要双层血缘

单层血缘无法同时满足架构治理和数据调查两个场景:

  • Schema 级血缘:回答"Employee.salary 字段来自哪些源系统的哪些字段?"——支持架构变更的影响分析
  • 实例级血缘:回答"员工张三的薪资数字 28500 是怎么算出来的?"——支持数据质量调查和合规审计
维度Schema 级血缘实例级血缘
粒度对象类型 / 属性数据记录 / 字段值
变化频率低(Schema 变更时)高(每次数据写入)
存储量小(千级节点)大(百万级节点)
主要用户架构师、治理团队数据分析师、审计团队
典型查询影响分析数据溯源

#1.2 对标 Palantir Foundry

能力Palantir Foundrycoomia-dip
Schema 血缘Dataset → Transform 链路ObjectType → Property 级
实例血缘Transaction-levelRecord + Field 级
可视化Monocle内置 DAG 可视化
查询 API内部gRPC + SDK
存储内部Iceberg + Neo4j

#2. 双层血缘架构

#2.1 系统架构

Code
┌────────────────────────────────────────────────────┐
│                Lineage Query API                    │
│           (gRPC + REST, SDK 封装)                    │
└───────────────────────┬────────────────────────────┘
                        │
┌───────────────────────▼────────────────────────────┐
│              Lineage Graph Engine                   │
│  ┌─────────────────┐  ┌─────────────────────────┐  │
│  │ Schema Lineage  │  │ Instance Lineage        │  │
│  │ (Graph DB)      │  │ (Iceberg + Graph Index) │  │
│  └────────┬────────┘  └────────────┬────────────┘  │
│           │                        │                │
│  ┌────────▼────────────────────────▼────────────┐  │
│  │         Unified Lineage Model                │  │
│  │    (LineageNode, LineageEdge, LineageGraph)   │  │
│  └──────────────────────────────────────────────┘  │
└────────────────────────────────────────────────────┘
                        ▲
┌───────────────────────┼────────────────────────────┐
│           Lineage Collectors                        │
│  ┌──────────┐  ┌──────────┐  ┌─────────────────┐  │
│  │ Schema   │  │ Pipeline │  │ Action          │  │
│  │ Collector│  │ Collector│  │ Collector       │  │
│  └──────────┘  └──────────┘  └─────────────────┘  │
└────────────────────────────────────────────────────┘

#2.2 统一血缘模型

Python
class LineageNode(BaseModel):
    """血缘节点"""

    node_id: str = Field(description="节点唯一标识")
    node_type: LineageNodeType = Field(description="节点类型")
    layer: LineageLayer = Field(description="血缘层级")

    # 标识信息
    name: str
    qualified_name: str  # 全限定名(如 ObjectType.Property)

    # 元数据
    metadata: dict = Field(default_factory=dict)
    classification: ClassificationLevel | None = None

    # 时间
    created_at: datetime
    updated_at: datetime | None = None


class LineageNodeType(str, Enum):
    # Schema 级
    OBJECT_TYPE = "object_type"
    PROPERTY = "property"
    LINK_TYPE = "link_type"
    ACTION_TYPE = "action_type"
    EXTERNAL_SOURCE = "external_source"

    # 实例级
    RECORD = "record"
    FIELD_VALUE = "field_value"
    TRANSFORM_EXECUTION = "transform_execution"
    PIPELINE_RUN = "pipeline_run"


class LineageLayer(str, Enum):
    SCHEMA = "schema"
    INSTANCE = "instance"


class LineageEdge(BaseModel):
    """血缘边"""

    edge_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
    source_id: str
    target_id: str
    edge_type: LineageEdgeType
    layer: LineageLayer

    # 变换信息
    transform_type: str | None = None    # map, filter, aggregate, join, derive
    transform_expression: str | None = None
    confidence: float = Field(default=1.0, ge=0.0, le=1.0)

    # 时间
    created_at: datetime
    valid_from: datetime | None = None
    valid_until: datetime | None = None


class LineageEdgeType(str, Enum):
    DERIVES_FROM = "derives_from"
    TRANSFORMS_TO = "transforms_to"
    COPIES_FROM = "copies_from"
    AGGREGATES_FROM = "aggregates_from"
    JOINS_WITH = "joins_with"
    REFERENCES = "references"

#3. Schema 级血缘

#3.1 Schema 血缘采集器

Python
class SchemaLineageCollector:
    """Schema 级血缘采集器"""

    async def on_object_type_created(self, event: SchemaChangeEvent) -> None:
        """对象类型创建时采集血缘"""
        obj_node = LineageNode(
            node_id=f"schema:{event.schema_id}",
            node_type=LineageNodeType.OBJECT_TYPE,
            layer=LineageLayer.SCHEMA,
            name=event.schema_name,
            qualified_name=event.schema_name,
            created_at=event.timestamp,
        )
        await self._graph.add_node(obj_node)

        # 为每个属性创建节点
        for prop in event.properties:
            prop_node = LineageNode(
                node_id=f"schema:{event.schema_id}.{prop.name}",
                node_type=LineageNodeType.PROPERTY,
                layer=LineageLayer.SCHEMA,
                name=prop.name,
                qualified_name=f"{event.schema_name}.{prop.name}",
                classification=prop.classification.level if prop.classification else None,
                created_at=event.timestamp,
            )
            await self._graph.add_node(prop_node)

    async def on_derived_property_defined(
        self,
        object_type: str,
        property_name: str,
        source_expression: str,
        source_properties: list[tuple[str, str]],
    ) -> None:
        """派生属性定义时建立血缘关系"""
        target_id = f"schema:{object_type}.{property_name}"

        for src_type, src_prop in source_properties:
            source_id = f"schema:{src_type}.{src_prop}"
            edge = LineageEdge(
                source_id=source_id,
                target_id=target_id,
                edge_type=LineageEdgeType.DERIVES_FROM,
                layer=LineageLayer.SCHEMA,
                transform_type="derive",
                transform_expression=source_expression,
                created_at=datetime.utcnow(),
            )
            await self._graph.add_edge(edge)

#3.2 Schema 级查询

Python
class SchemaLineageQuery:
    """Schema 级血缘查询"""

    async def get_upstream(
        self,
        qualified_name: str,
        max_depth: int = 10,
    ) -> LineageGraph:
        """获取上游血缘(数据从哪来)"""
        return await self._graph.traverse(
            start_node=f"schema:{qualified_name}",
            direction="upstream",
            edge_types=[LineageEdgeType.DERIVES_FROM, LineageEdgeType.COPIES_FROM],
            max_depth=max_depth,
            layer=LineageLayer.SCHEMA,
        )

    async def get_downstream(
        self,
        qualified_name: str,
        max_depth: int = 10,
    ) -> LineageGraph:
        """获取下游血缘(数据影响什么)"""
        return await self._graph.traverse(
            start_node=f"schema:{qualified_name}",
            direction="downstream",
            edge_types=[LineageEdgeType.TRANSFORMS_TO, LineageEdgeType.DERIVES_FROM],
            max_depth=max_depth,
            layer=LineageLayer.SCHEMA,
        )

    async def impact_analysis(
        self,
        qualified_name: str,
    ) -> ImpactReport:
        """变更影响分析"""
        downstream = await self.get_downstream(qualified_name, max_depth=20)
        affected_objects = set()
        affected_properties = set()
        affected_actions = set()

        for node in downstream.nodes:
            if node.node_type == LineageNodeType.OBJECT_TYPE:
                affected_objects.add(node.qualified_name)
            elif node.node_type == LineageNodeType.PROPERTY:
                affected_properties.add(node.qualified_name)
            elif node.node_type == LineageNodeType.ACTION_TYPE:
                affected_actions.add(node.qualified_name)

        return ImpactReport(
            source=qualified_name,
            affected_object_types=list(affected_objects),
            affected_properties=list(affected_properties),
            affected_actions=list(affected_actions),
            total_affected=len(downstream.nodes),
        )

#4. 实例级血缘

#4.1 实例血缘采集器

Python
class InstanceLineageCollector:
    """实例级血缘采集器"""

    async def on_record_created(
        self,
        object_type: str,
        record_id: str,
        source_records: list[SourceRecord] | None = None,
        transform_context: TransformContext | None = None,
    ) -> None:
        """记录创建时采集实例血缘"""
        # 创建记录节点
        record_node = LineageNode(
            node_id=f"instance:{object_type}:{record_id}",
            node_type=LineageNodeType.RECORD,
            layer=LineageLayer.INSTANCE,
            name=record_id,
            qualified_name=f"{object_type}:{record_id}",
            created_at=datetime.utcnow(),
        )
        await self._graph.add_node(record_node)

        # 建立来源关系
        if source_records:
            for source in source_records:
                source_id = f"instance:{source.object_type}:{source.record_id}"
                edge = LineageEdge(
                    source_id=source_id,
                    target_id=record_node.node_id,
                    edge_type=LineageEdgeType.DERIVES_FROM,
                    layer=LineageLayer.INSTANCE,
                    transform_type=transform_context.transform_type if transform_context else None,
                    transform_expression=transform_context.expression if transform_context else None,
                    created_at=datetime.utcnow(),
                )
                await self._graph.add_edge(edge)

    async def on_field_computed(
        self,
        object_type: str,
        record_id: str,
        field_name: str,
        computed_value: Any,
        source_fields: list[SourceField],
        computation: str,
    ) -> None:
        """字段计算时采集字段级血缘"""
        target_id = f"instance:{object_type}:{record_id}:{field_name}"

        field_node = LineageNode(
            node_id=target_id,
            node_type=LineageNodeType.FIELD_VALUE,
            layer=LineageLayer.INSTANCE,
            name=f"{field_name}={computed_value}",
            qualified_name=f"{object_type}:{record_id}.{field_name}",
            metadata={"value": str(computed_value), "computation": computation},
            created_at=datetime.utcnow(),
        )
        await self._graph.add_node(field_node)

        for source in source_fields:
            source_id = f"instance:{source.object_type}:{source.record_id}:{source.field_name}"
            edge = LineageEdge(
                source_id=source_id,
                target_id=target_id,
                edge_type=LineageEdgeType.DERIVES_FROM,
                layer=LineageLayer.INSTANCE,
                transform_type="compute",
                transform_expression=computation,
                created_at=datetime.utcnow(),
            )
            await self._graph.add_edge(edge)

#4.2 实例级查询

Python
class InstanceLineageQuery:
    """实例级血缘查询"""

    async def trace_value_origin(
        self,
        object_type: str,
        record_id: str,
        field_name: str,
    ) -> ValueOriginTrace:
        """追溯字段值的来源"""
        start_node = f"instance:{object_type}:{record_id}:{field_name}"
        graph = await self._graph.traverse(
            start_node=start_node,
            direction="upstream",
            max_depth=50,
            layer=LineageLayer.INSTANCE,
        )

        # 构建值溯源链
        steps = []
        for edge in graph.edges_in_order():
            steps.append(TraceStep(
                source=edge.source_id,
                target=edge.target_id,
                transform=edge.transform_expression,
                timestamp=edge.created_at,
            ))

        return ValueOriginTrace(
            field=f"{object_type}.{field_name}",
            record_id=record_id,
            steps=steps,
            leaf_sources=[n for n in graph.leaf_nodes()],
        )

    async def get_record_provenance(
        self,
        object_type: str,
        record_id: str,
    ) -> RecordProvenance:
        """获取记录的完整来源"""
        node_id = f"instance:{object_type}:{record_id}"
        graph = await self._graph.traverse(
            start_node=node_id,
            direction="upstream",
            max_depth=20,
            layer=LineageLayer.INSTANCE,
        )

        return RecordProvenance(
            record_id=record_id,
            object_type=object_type,
            source_count=len(graph.leaf_nodes()),
            transform_count=len(graph.edges),
            lineage_graph=graph,
        )

#5. 血缘存储

#5.1 混合存储策略

Python
class HybridLineageStore:
    """混合血缘存储 - Schema 用图数据库,Instance 用 Iceberg"""

    def __init__(
        self,
        graph_store: GraphStore,       # Neo4j / JanusGraph
        table_store: IcebergStore,     # Iceberg tables
    ):
        self._graph = graph_store      # Schema 级 + 热实例级
        self._table = table_store      # 冷实例级

    async def add_node(self, node: LineageNode) -> None:
        if node.layer == LineageLayer.SCHEMA:
            await self._graph.upsert_node(node)
        else:
            # 实例级:双写(图索引 + Iceberg 持久化)
            await self._graph.upsert_node(node)
            await self._table.append_node(node)

    async def add_edge(self, edge: LineageEdge) -> None:
        if edge.layer == LineageLayer.SCHEMA:
            await self._graph.upsert_edge(edge)
        else:
            await self._graph.upsert_edge(edge)
            await self._table.append_edge(edge)

    async def archive_old_instances(self, before: datetime) -> int:
        """归档旧的实例级血缘(从图索引移除,保留 Iceberg)"""
        return await self._graph.delete_nodes(
            layer=LineageLayer.INSTANCE,
            created_before=before,
        )

#6. 血缘与安全集成

#6.1 分类传播

血缘关系驱动数据分类的自动传播:

Python
class LineageDrivenClassification:
    """基于血缘的分类传播"""

    async def propagate_classification(
        self,
        source_node_id: str,
        new_classification: ClassificationLevel,
    ) -> list[str]:
        """当源节点分类变更时,传播到下游"""
        downstream = await self._lineage_query.get_downstream(
            source_node_id, max_depth=20,
        )

        affected_nodes = []
        for node in downstream.nodes:
            if node.classification is None or node.classification < new_classification:
                node.classification = new_classification
                await self._graph.update_node(node)
                affected_nodes.append(node.node_id)

        return affected_nodes

#6.2 审计集成

血缘查询操作本身也被审计记录:

Python
@audit_tracked(AuditEventType.DATA_ACCESS, action="lineage_query")
async def trace_value_origin(self, object_type, record_id, field_name):
    ...

#7. 测试策略

Python
class TestDualLayerLineage:
    async def test_schema_lineage_creation(self):
        collector = SchemaLineageCollector(graph)
        await collector.on_derived_property_defined(
            object_type="Employee",
            property_name="total_compensation",
            source_expression="base_salary + bonus",
            source_properties=[
                ("Employee", "base_salary"),
                ("Employee", "bonus"),
            ],
        )

        upstream = await query.get_upstream("Employee.total_compensation")
        assert len(upstream.nodes) == 3  # target + 2 sources
        assert len(upstream.edges) == 2

    async def test_instance_lineage_trace(self):
        collector = InstanceLineageCollector(graph)
        await collector.on_field_computed(
            object_type="Employee",
            record_id="emp-001",
            field_name="total_compensation",
            computed_value=35000,
            source_fields=[
                SourceField("Employee", "emp-001", "base_salary"),
                SourceField("Employee", "emp-001", "bonus"),
            ],
            computation="base_salary + bonus",
        )

        trace = await query.trace_value_origin("Employee", "emp-001", "total_compensation")
        assert len(trace.steps) == 2
        assert len(trace.leaf_sources) == 2

    async def test_impact_analysis(self):
        report = await query.impact_analysis("SourceSystem.raw_salary")
        assert "Employee.base_salary" in report.affected_properties
        assert "Employee.total_compensation" in report.affected_properties

#8. 生产最佳实践

#8.1 存储策略

血缘层级热存储温存储冷存储
Schema 级永久(图数据库)--
实例级30 天(图索引)365 天(Iceberg)5 年(归档)

#8.2 性能优化

  • Schema 级查询:< 50ms(图索引)
  • 实例级溯源:< 200ms(热数据)
  • 批量血缘采集:异步管道,不阻塞写入路径
  • 定期归档旧的实例级数据,保持图索引轻量

#8.3 数据质量

  • 血缘完整性检查:定期验证所有派生属性都有血缘记录
  • 孤立节点检测:定期清理无边连接的血缘节点
  • 血缘断裂告警:当数据变换未产生血缘记录时触发告警

#9. 总结

coomia-dip 的双层数据血缘系统通过 Schema 级和实例级两个层面,覆盖了从架构治理到数据调查的全部血缘需求。关键设计亮点:

  1. 双层设计:Schema 级满足影响分析,实例级满足数据溯源
  2. 统一模型:两层共用 LineageNode/LineageEdge 模型,降低复杂度
  3. 混合存储:图数据库 + Iceberg,兼顾查询性能和存储经济性
  4. 安全集成:血缘驱动分类传播,与审计系统深度集成
  5. 可扩展:通过 Collector 插件机制支持新的数据变换场景

下一篇将探讨 coomia-dip 的历史数据回放能力。