源码精读:LineageService — 实体级+字段级血缘
coomia-dip 的血缘追踪由两个互补的 LineageService 组成:Pipeline & Orchestration Layer(数据工程视角)提供 Dataset-centric 的 Pipeline 血缘,Metadata & Governance Layer(治理视角)提供 Entity-centric 的业务血缘。两者共同实现了从数据源到决策结果的端到端血缘追踪。本文深入剖析两套 Proto 合约的节点/边类型设计差异、双视角血缘图的合并策略、字段级血缘的 FieldMapping 六种转换类型、AnalyzeImpact 的影响分析算法、以及血缘摘要(LineageSummary)的可视化数据模型。
源码精读:LineageService — 实体级+字段级血缘
“系列:S9 源码精读 · 第 18 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
coomia-dip 的血缘追踪由两个互补的 LineageService 组成:Pipeline & Orchestration Layer(数据工程视角)提供 Dataset-centric 的 Pipeline 血缘,Metadata & Governance Layer(治理视角)提供 Entity-centric 的业务血缘。两者共同实现了从数据源到决策结果的端到端血缘追踪。本文深入剖析两套 Proto 合约的节点/边类型设计差异、双视角血缘图的合并策略、字段级血缘的 FieldMapping 六种转换类型、AnalyzeImpact 的影响分析算法、以及血缘摘要(LineageSummary)的可视化数据模型。
#目录
- 双视角架构:Pipeline & Orchestration Layer vs Metadata & Governance Layer
- Pipeline & Orchestration Layer Proto:Dataset-centric 血缘
- Metadata & Governance Layer Proto:Entity-centric 业务血缘
- 节点类型设计:13 种 NodeType 的分类逻辑
- 边类型设计:12 种 EdgeType 的语义
- 字段级血缘:FieldMapping 与 TransformType
- GetLineage:深度控制与时间点查询
- AnalyzeImpact:变更影响分析
- TraceOrigin:数据溯源
- 血缘摘要与统计
- Key Takeaways
#1. 双视角架构:Pipeline & Orchestration Layer vs Metadata & Governance Layer
coomia-dip 的血缘追踪采用"双视角"设计,两个 LineageService 服务于不同的使用场景:
proto/plane_f/lineage.proto # 数据工程视角(291 行)
proto/plane_g/lineage_service.proto # 业务治理视角(429 行)
| 维度 | Pipeline & Orchestration Layer (数据工程) | Metadata & Governance Layer (业务治理) |
|---|---|---|
| 关注点 | Dataset、Pipeline、Step | Entity、Reasoning、Decision |
| 粒度 | 数据集级 + 字段级 | 实体级 + 字段级 |
| 触发方式 | Pipeline 执行自动记录 | 跨 Layer 手动/自动记录 |
| 典型用户 | 数据工程师 | 业务分析师、合规官 |
#2. Pipeline & Orchestration Layer Proto:Dataset-centric 血缘
Pipeline & Orchestration Layer 的 LineageService 围绕 Pipeline 执行构建血缘图:
service LineageService {
rpc RecordLineage(RecordLineageRequest) returns (RecordLineageResponse);
rpc BatchRecordLineage(BatchRecordLineageRequest) returns (BatchRecordResponse);
rpc GetUpstream(GetLineageRequest) returns (LineageGraphResponse);
rpc GetDownstream(GetLineageRequest) returns (LineageGraphResponse);
rpc GetFullLineage(GetFullLineageRequest) returns (LineageGraphResponse);
rpc GetFieldImpact(FieldImpactRequest) returns (FieldImpactResponse);
rpc GetFieldOrigin(FieldOriginRequest) returns (FieldOriginResponse);
rpc GetLineageSummary(GetSummaryRequest) returns (SummaryResponse);
rpc GetLineageStats(GetLineageStatsRequest) returns (LineageStatsResponse);
rpc DeleteLineage(DeleteLineageRequest) returns (DeleteResponse);
}
RecordLineage 是 Pipeline 执行的核心产出——每个 Step 完成后记录其输入/输出和字段映射:
message RecordLineageRequest {
ExecutionContext context = 1;
repeated DatasetRef inputs = 2;
repeated DatasetRef outputs = 3;
repeated FieldMapping field_mappings = 4;
map<string, string> metadata = 5;
}
ExecutionContext 携带了完整的 Pipeline 执行上下文:
message ExecutionContext {
string pipeline_id = 1;
string execution_id = 2;
string step_id = 3;
string step_name = 4;
string user_id = 5;
google.protobuf.Timestamp timestamp = 6;
google.protobuf.Duration duration = 7;
}
duration 字段使得血缘图不仅展示"数据从哪来",还能展示"处理花了多久"——为性能瓶颈分析提供了数据基础。
#3. Metadata & Governance Layer Proto:Entity-centric 业务血缘
Metadata & Governance Layer 的 LineageService 关注业务语义层面的血缘关系:
service LineageService {
rpc GetLineage(GetLineageRequest) returns (GetLineageResponse);
rpc TraceOrigin(TraceOriginRequest) returns (TraceOriginResponse);
rpc GetDownstream(GetDownstreamRequest) returns (GetDownstreamResponse);
rpc GetUpstream(GetUpstreamRequest) returns (GetUpstreamResponse);
rpc AnalyzeImpact(AnalyzeImpactRequest) returns (ImpactAnalysisResponse);
rpc GetFieldImpact(GetFieldImpactRequest) returns (FieldImpactResponse);
rpc GetFieldOrigin(GetFieldOriginRequest) returns (FieldOriginResponse);
rpc RecordLineage(RecordLineageRequest) returns (RecordLineageResponse);
rpc BatchRecordLineage(BatchRecordLineageRequest) returns (BatchRecordLineageResponse);
rpc GetLineageSummary(GetLineageSummaryRequest) returns (LineageSummaryResponse);
rpc DeleteLineage(DeleteLineageRequest) returns (DeleteLineageResponse);
}
相比 Pipeline & Orchestration Layer,Metadata & Governance Layer 增加了 TraceOrigin(溯源到原始数据源)和 AnalyzeImpact(变更影响分析)两个高级查询。
#4. 节点类型设计:13 种 NodeType 的分类逻辑
Metadata & Governance Layer 定义了 13 种节点类型,分为三个逻辑组:
enum NodeType {
// Entity-centric types (业务血缘)
NODE_TYPE_ENTITY = 1; // 本体实例
NODE_TYPE_REASONING = 2; // 推理过程
NODE_TYPE_DECISION = 3; // 决策结果
NODE_TYPE_PIPELINE = 4; // 数据管道
NODE_TYPE_SOURCE_SYSTEM = 5; // 外部源系统
NODE_TYPE_SCHEMA = 6; // Schema 定义
// Dataset-centric types (数据工程血缘)
NODE_TYPE_DATASET = 10; // 数据集
NODE_TYPE_PIPELINE_STEP = 11; // Pipeline 步骤
NODE_TYPE_FIELD = 12; // 字段
NODE_TYPE_EXTERNAL = 13; // 外部数据源
}
编号设计值得关注:Entity-centric 类型使用 1-6,Dataset-centric 类型从 10 开始,留出了 7-9 的扩展空间。这是一个常见的 Proto 枚举预留策略。
Pipeline & Orchestration Layer 的 NodeType 更简洁,只有 5 种:
enum NodeType {
NODE_TYPE_DATASET = 1;
NODE_TYPE_PIPELINE = 2;
NODE_TYPE_STEP = 3;
NODE_TYPE_FIELD = 4;
NODE_TYPE_EXTERNAL = 5;
}
#5. 边类型设计:12 种 EdgeType 的语义
Metadata & Governance Layer 定义了 12 种边类型,分为三个语义组:
enum EdgeType {
// Entity-centric edges (业务血缘关系)
EDGE_TYPE_DERIVED_FROM = 1; // 派生自
EDGE_TYPE_USED_BY = 2; // 被使用
EDGE_TYPE_TRANSFORMED_BY = 3; // 被转换
EDGE_TYPE_INFERRED_FROM = 4; // 推导自
EDGE_TYPE_SUPPORTED_BY = 5; // 支撑了
EDGE_TYPE_TRIGGERED_BY = 6; // 触发了
// Dataset-centric edges (数据工程血缘)
EDGE_TYPE_READ = 10; // 读取
EDGE_TYPE_WRITE = 11; // 写入
EDGE_TYPE_MAPS_TO = 12; // 字段映射
EDGE_TYPE_CONTAINS = 13; // 包含关系
// Cross-layer edges (跨层关联)
EDGE_TYPE_PRODUCES = 20; // Pipeline 产出 Entity
EDGE_TYPE_CONSUMES = 21; // Pipeline 消费 Entity
}
跨层边(20-21)是双视角血缘图的桥梁——它们连接了数据工程层面的 Pipeline 与业务层面的 Entity,实现了端到端的血缘追踪。
#6. 字段级血缘:FieldMapping 与 TransformType
字段级血缘是两个 LineageService 共有的高级特性:
message FieldMapping {
repeated FieldRef source_fields = 1; // 多对一映射
FieldRef target_field = 2;
string transform_expression = 3;
TransformType transform_type = 4;
map<string, string> properties = 5;
}
多对一映射:source_fields 是 repeated,支持 JOIN 或聚合等多输入场景。例如 revenue = price * quantity 涉及两个源字段。
六种转换类型:
enum TransformType {
TRANSFORM_TYPE_DIRECT = 1; // 直接复制
TRANSFORM_TYPE_AGGREGATE = 2; // 聚合
TRANSFORM_TYPE_JOIN = 3; // Join
TRANSFORM_TYPE_FILTER = 4; // 过滤
TRANSFORM_TYPE_UDF = 5; // 用户自定义函数
TRANSFORM_TYPE_EXPRESSION = 6; // 表达式计算
}
transform_expression 字段存储实际的转换表达式(如 SQL 片段),使得事后审计时可以理解字段是"如何"从源头到达目标的。
#7. GetLineage:深度控制与时间点查询
Metadata & Governance Layer 的 GetLineage 支持丰富的查询控制:
message GetLineageRequest {
string entity_id = 2;
LineageDirection direction = 3; // UPSTREAM / DOWNSTREAM / BOTH
int32 max_depth = 4; // 默认 3
google.protobuf.Timestamp as_of_time = 5; // 时间点查询
repeated NodeType node_type_filter = 6;
}
max_depth 控制图遍历的深度,默认 3 层已能覆盖大多数业务场景。设置过大可能导致返回海量数据。
as_of_time 支持历史时间点查询——"这个实体在上周三的血缘图是什么样的?"这对事后审计至关重要。
node_type_filter 允许只查看特定类型的节点,例如只关注 SOURCE_SYSTEM 类型来快速定位数据源头。
响应返回完整的图结构:
message GetLineageResponse {
repeated LineageNode nodes = 1;
repeated LineageEdge edges = 2;
string root_node_id = 3;
int32 total_nodes = 4;
int32 total_edges = 5;
}
#8. AnalyzeImpact:变更影响分析
AnalyzeImpact 是 Metadata & Governance Layer 独有的高级功能,用于评估变更的影响范围:
message AnalyzeImpactRequest {
string target_id = 2;
string change_type = 3; // SCHEMA_CHANGE, DATA_UPDATE, DELETE
google.protobuf.Struct change_details = 4;
}
message ImpactAnalysisResponse {
string analysis_id = 1;
int32 total_affected_count = 2;
repeated ImpactGroup groups = 3;
ImpactSeverity overall_severity = 4;
string risk_level = 5;
message ImpactGroup {
NodeType type = 1;
int32 count = 2;
repeated string sample_node_ids = 3;
ImpactSeverity severity = 4;
}
}
四级影响严重度:
enum ImpactSeverity {
IMPACT_SEVERITY_LOW = 1; // 仅元数据变更
IMPACT_SEVERITY_MEDIUM = 2; // 中间推理步骤变化
IMPACT_SEVERITY_HIGH = 3; // 推理置信度大幅下降
IMPACT_SEVERITY_CRITICAL = 4; // 直接导致决策变化
}
ImpactGroup 按节点类型分组统计影响,sample_node_ids 提供了少量示例供人工确认。这种"统计 + 采样"的模式避免了返回完整受影响实体列表的性能问题。
#9. TraceOrigin:数据溯源
TraceOrigin 递归查找所有 SOURCE_SYSTEM 类型的上游节点:
message TraceOriginRequest {
string entity_id = 2;
int32 max_depth = 3;
}
message TraceOriginResponse {
repeated LineageNode source_systems = 1;
repeated LineageEdge paths = 2;
int32 total_sources = 3;
}
响应直接返回所有源系统节点和到达路径,无需客户端遍历完整血缘图再过滤。这是一个面向用例的 API 设计——"这个决策的数据最终来自哪些系统?"是最常见的合规审计问题。
#10. 血缘摘要与统计
Metadata & Governance Layer 的血缘摘要提供了快速概览:
message LineageSummaryResponse {
string entity_id = 1;
int32 upstream_count = 2;
int32 downstream_count = 3;
int32 field_count = 4;
int32 pipeline_count = 5;
google.protobuf.Timestamp last_updated = 6;
map<string, int32> node_type_counts = 7;
}
node_type_counts 提供了按节点类型的分布统计,可以直接用于前端的血缘概览仪表板。
Pipeline & Orchestration Layer 的统计更偏向运维指标:
message LineageStatsResponse {
int64 total_nodes = 1;
int64 total_edges = 2;
int64 total_field_mappings = 3;
map<string, int64> nodes_by_type = 4;
map<string, int64> edges_by_type = 5;
int64 events_recorded = 6;
}
events_recorded 是一个关键运维指标——如果此值停止增长,说明 Pipeline 血缘记录可能中断了。
#11. Key Takeaways
- 双视角设计:Pipeline & Orchestration Layer(数据工程)和 Metadata & Governance Layer(业务治理)各自关注不同层面,通过跨层边连接。
- 13 种节点类型 + 12 种边类型构建了丰富的血缘语义,覆盖了从数据源到决策的完整链路。
- 字段级血缘的
FieldMapping支持多对一映射和六种转换类型,提供了细粒度的数据追溯能力。 - AnalyzeImpact 的"统计 + 采样"模式在影响分析的完整性和响应性能之间取得了平衡。
- TraceOrigin 是面向用例的 API 设计——直接回答"数据来自哪里"这个最常见的合规问题。
- 编号预留策略:NodeType 和 EdgeType 的编号设计留出了扩展空间。
#下一篇
S9-19:OntoPlatform SDK — 28 个子模块的 Facade,我们将深入 Python SDK 的 OntoPlatform 类,了解如何通过一个入口点访问整个平台的能力。
Tags: #coomia-dip #source-code-reading #lineage-service #entity-lineage #field-lineage #grpc #Layer-g