双层数据血缘追踪: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 Foundry | coomia-dip |
|---|---|---|
| Schema 血缘 | Dataset → Transform 链路 | ObjectType → Property 级 |
| 实例血缘 | Transaction-level | Record + 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 级和实例级两个层面,覆盖了从架构治理到数据调查的全部血缘需求。关键设计亮点:
- 双层设计:Schema 级满足影响分析,实例级满足数据溯源
- 统一模型:两层共用 LineageNode/LineageEdge 模型,降低复杂度
- 混合存储:图数据库 + Iceberg,兼顾查询性能和存储经济性
- 安全集成:血缘驱动分类传播,与审计系统深度集成
- 可扩展:通过 Collector 插件机制支持新的数据变换场景
下一篇将探讨 coomia-dip 的历史数据回放能力。