返回博客

源码精读:LineageService — 实体级+字段级血缘

coomia-dip 的血缘追踪由两个互补的 LineageService 组成:Pipeline & Orchestration Layer(数据工程视角)提供 Dataset-centric 的 Pipeline 血缘,Metadata & Governance Layer(治理视角)提供 Entity-centric 的业务血缘。两者共同实现了从数据源到决策结果的端到端血缘追踪。本文深入剖析两套 Proto 合约的节点/边类型设计差异、双视角血缘图的合并策略、字段级血缘的 FieldMapping 六种转换类型、AnalyzeImpact 的影响分析算法、以及血缘摘要(LineageSummary)的可视化数据模型。

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

源码精读:LineageService — 实体级+字段级血缘

系列:S9 源码精读 · 第 18 篇 | 难度:高级 | 阅读时间:25 分钟

#TL;DR

coomia-dip 的血缘追踪由两个互补的 LineageService 组成:Pipeline & Orchestration Layer(数据工程视角)提供 Dataset-centric 的 Pipeline 血缘,Metadata & Governance Layer(治理视角)提供 Entity-centric 的业务血缘。两者共同实现了从数据源到决策结果的端到端血缘追踪。本文深入剖析两套 Proto 合约的节点/边类型设计差异、双视角血缘图的合并策略、字段级血缘的 FieldMapping 六种转换类型、AnalyzeImpact 的影响分析算法、以及血缘摘要(LineageSummary)的可视化数据模型。

#目录

  1. 双视角架构:Pipeline & Orchestration Layer vs Metadata & Governance Layer
  2. Pipeline & Orchestration Layer Proto:Dataset-centric 血缘
  3. Metadata & Governance Layer Proto:Entity-centric 业务血缘
  4. 节点类型设计:13 种 NodeType 的分类逻辑
  5. 边类型设计:12 种 EdgeType 的语义
  6. 字段级血缘:FieldMapping 与 TransformType
  7. GetLineage:深度控制与时间点查询
  8. AnalyzeImpact:变更影响分析
  9. TraceOrigin:数据溯源
  10. 血缘摘要与统计
  11. Key Takeaways

#1. 双视角架构:Pipeline & Orchestration Layer vs Metadata & Governance Layer

coomia-dip 的血缘追踪采用"双视角"设计,两个 LineageService 服务于不同的使用场景:

Code
proto/plane_f/lineage.proto          # 数据工程视角(291 行)
proto/plane_g/lineage_service.proto  # 业务治理视角(429 行)
维度Pipeline & Orchestration Layer (数据工程)Metadata & Governance Layer (业务治理)
关注点Dataset、Pipeline、StepEntity、Reasoning、Decision
粒度数据集级 + 字段级实体级 + 字段级
触发方式Pipeline 执行自动记录跨 Layer 手动/自动记录
典型用户数据工程师业务分析师、合规官

#2. Pipeline & Orchestration Layer Proto:Dataset-centric 血缘

Pipeline & Orchestration Layer 的 LineageService 围绕 Pipeline 执行构建血缘图:

PROTOBUF
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 完成后记录其输入/输出和字段映射:

PROTOBUF
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 执行上下文:

PROTOBUF
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 关注业务语义层面的血缘关系:

PROTOBUF
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 种节点类型,分为三个逻辑组:

PROTOBUF
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 种:

PROTOBUF
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 种边类型,分为三个语义组:

PROTOBUF
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 共有的高级特性:

PROTOBUF
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_fieldsrepeated,支持 JOIN 或聚合等多输入场景。例如 revenue = price * quantity 涉及两个源字段。

六种转换类型:

PROTOBUF
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 支持丰富的查询控制:

PROTOBUF
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 类型来快速定位数据源头。

响应返回完整的图结构:

PROTOBUF
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 独有的高级功能,用于评估变更的影响范围:

PROTOBUF
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;
  }
}

四级影响严重度

PROTOBUF
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 类型的上游节点:

PROTOBUF
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 的血缘摘要提供了快速概览:

PROTOBUF
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 的统计更偏向运维指标:

PROTOBUF
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

  1. 双视角设计:Pipeline & Orchestration Layer(数据工程)和 Metadata & Governance Layer(业务治理)各自关注不同层面,通过跨层边连接。
  2. 13 种节点类型 + 12 种边类型构建了丰富的血缘语义,覆盖了从数据源到决策的完整链路。
  3. 字段级血缘FieldMapping 支持多对一映射和六种转换类型,提供了细粒度的数据追溯能力。
  4. AnalyzeImpact 的"统计 + 采样"模式在影响分析的完整性和响应性能之间取得了平衡。
  5. TraceOrigin 是面向用例的 API 设计——直接回答"数据来自哪里"这个最常见的合规问题。
  6. 编号预留策略: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