Java 数据血缘案例
数据血缘(Data Lineage)追踪数据从源头到目标的完整生命周期,以下通过一个电商数据分析的实际案例,展示如何在 Java 中实现数据血缘追踪。

案例场景:电商订单分析系统
假设有一个电商平台,需要追踪订单数据从产生到最终报表的完整流转过程。
数据血缘模型定义
// 数据节点(数据元素)
@Data
@Builder
public class DataNode {
private String nodeId; // 节点ID
private String name; // 节点名称
private String type; // 类型:TABLE/COLUMN/FIELD
private String database; // 所属数据库
private String schema; // 所属Schema
private String table; // 所属表
private String column; // 所属列
}
// 数据关系(血缘边)
@Data
@Builder
public class DataLineage {
private String lineageId; // 血缘ID
private String sourceNodeId; // 源节点
private String targetNodeId; // 目标节点
private String transformation; // 转换逻辑
private String processType; // 处理类型:SELECT/JOIN/AGGREGATE/ETL
private LocalDateTime timestamp;
}
血缘追踪器实现
@Service
public class DataLineageTracker {
private final Graph<String, DefaultEdge> lineageGraph = new DirectedMultigraph<>(DefaultEdge.class);
private final Map<String, DataNode> nodeMap = new ConcurrentHashMap<>();
/**
* 注册数据节点
*/
public void registerNode(DataNode node) {
nodeMap.put(node.getNodeId(), node);
lineageGraph.addVertex(node.getNodeId());
}
/**
* 记录数据流转关系
*/
public void recordLineage(DataLineage lineage) {
lineageGraph.addEdge(lineage.getSourceNodeId(),
lineage.getTargetNodeId(),
new DefaultEdge());
}
/**
* 获取数据来源(溯源)
*/
public List<DataNode> getDataSources(String targetNodeId) {
List<DataNode> sources = new ArrayList<>();
Set<String> visited = new HashSet<>();
traverseUpstream(targetNodeId, visited, sources);
return sources;
}
private void traverseUpstream(String nodeId, Set<String> visited,
List<DataNode> sources) {
if (visited.contains(nodeId)) return;
visited.add(nodeId);
Set<DefaultEdge> incomingEdges = lineageGraph.incomingEdgesOf(nodeId);
if (incomingEdges.isEmpty()) {
// 这是源头节点
sources.add(nodeMap.get(nodeId));
} else {
for (DefaultEdge edge : incomingEdges) {
String source = lineageGraph.getEdgeSource(edge);
traverseUpstream(source, visited, sources);
}
}
}
/**
* 获取数据影响范围(下游)
*/
public List<DataNode> getDataImpact(String sourceNodeId) {
List<DataNode> impacted = new ArrayList<>();
Set<String> visited = new HashSet<>();
traverseDownstream(sourceNodeId, visited, impacted);
return impacted;
}
private void traverseDownstream(String nodeId, Set<String> visited,
List<DataNode> impacted) {
if (visited.contains(nodeId)) return;
visited.add(nodeId);
Set<DefaultEdge> outgoingEdges = lineageGraph.outgoingEdgesOf(nodeId);
if (!outgoingEdges.isEmpty()) {
for (DefaultEdge edge : outgoingEdges) {
String target = lineageGraph.getEdgeTarget(edge);
impacted.add(nodeMap.get(target));
traverseDownstream(target, visited, impacted);
}
}
}
}
实际业务场景模拟
@Component
public class OrderDataLineageDemo {
@Autowired
private DataLineageTracker tracker;
@PostConstruct
public void simulateOrderLineage() {
// 1. 注册源数据节点
DataNode orderTable = DataNode.builder()
.nodeId("order_db.orders")
.name("订单表")
.type("TABLE")
.database("order_db")
.table("orders")
.build();
DataNode orderIdColumn = DataNode.builder()
.nodeId("order_db.orders.order_id")
.name("订单ID")
.type("COLUMN")
.database("order_db")
.table("orders")
.column("order_id")
.build();
DataNode amountColumn = DataNode.builder()
.nodeId("order_db.orders.amount")
.name("订单金额")
.type("COLUMN")
.database("order_db")
.table("orders")
.column("amount")
.build();
// 2. 注册中间计算节点
DataNode dailySummary = DataNode.builder()
.nodeId("dw.daily_order_summary")
.name("每日订单汇总")
.type("TABLE")
.database("dw")
.table("daily_order_summary")
.build();
DataNode totalAmount = DataNode.builder()
.nodeId("dw.daily_order_summary.total_amount")
.name("总金额")
.type("COLUMN")
.database("dw")
.table("daily_order_summary")
.column("total_amount")
.build();
DataNode orderCount = DataNode.builder()
.nodeId("dw.daily_order_summary.order_count")
.name("订单数量")
.type("COLUMN")
.database("dw")
.table("daily_order_summary")
.column("order_count")
.build();
// 3. 注册报表节点
DataNode report = DataNode.builder()
.nodeId("report.daily_sales")
.name("日报表")
.type("TABLE")
.database("report")
.table("daily_sales")
.build();
DataNode reportTotalSales = DataNode.builder()
.nodeId("report.daily_sales.total_sales")
.name("总销售额")
.type("COLUMN")
.database("report")
.table("daily_sales")
.column("total_sales")
.build();
// 注册所有节点
Arrays.asList(orderTable, orderIdColumn, amountColumn,
dailySummary, totalAmount, orderCount,
report, reportTotalSales)
.forEach(tracker::registerNode);
// 4. 记录血缘关系
// 订单列 -> 汇总表
tracker.recordLineage(DataLineage.builder()
.sourceNodeId("order_db.orders.amount")
.targetNodeId("dw.daily_order_summary.total_amount")
.transformation("SUM(amount)")
.processType("AGGREGATE")
.build());
tracker.recordLineage(DataLineage.builder()
.sourceNodeId("order_db.orders.order_id")
.targetNodeId("dw.daily_order_summary.order_count")
.transformation("COUNT(order_id)")
.processType("AGGREGATE")
.build());
// 汇总表 -> 报表
tracker.recordLineage(DataLineage.builder()
.sourceNodeId("dw.daily_order_summary.total_amount")
.targetNodeId("report.daily_sales.total_sales")
.transformation("total_amount * exchange_rate")
.processType("ETL")
.build());
}
}
血缘查询与可视化
@RestController
@RequestMapping("/api/lineage")
public class LineageQueryController {
@Autowired
private DataLineageTracker tracker;
/**
* 查询数据来源
*/
@GetMapping("/sources/{targetNodeId}")
public ResponseEntity<List<DataNode>> getSources(
@PathVariable String targetNodeId) {
List<DataNode> sources = tracker.getDataSources(targetNodeId);
return ResponseEntity.ok(sources);
}
/**
* 查询数据影响范围
*/
@GetMapping("/impact/{sourceNodeId}")
public ResponseEntity<List<DataNode>> getImpact(
@PathVariable String sourceNodeId) {
List<DataNode> impacted = tracker.getDataImpact(sourceNodeId);
return ResponseEntity.ok(impacted);
}
/**
* 生成血缘图(D3.js兼容格式)
*/
@GetMapping("/graph")
public ResponseEntity<Map<String, Object>> getLineageGraph() {
// 转换为前端可用的格式
Map<String, Object> graph = new HashMap<>();
List<Map<String, String>> nodes = new ArrayList<>();
List<Map<String, String>> edges = new ArrayList<>();
// 构建节点和边
// ... 转换逻辑
graph.put("nodes", nodes);
graph.put("edges", edges);
return ResponseEntity.ok(graph);
}
}
数据血缘的存储实现
@Entity
@Table(name = "data_lineage_records")
public class LineageRecord {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
@Column(name = "source_table")
private String sourceTable;
@Column(name = "source_column")
private String sourceColumn;
@Column(name = "target_table")
private String targetTable;
@Column(name = "target_column")
private String targetColumn;
@Column(name = "transformation_logic")
private String transformationLogic;
@Column(name = "etl_job_id")
private String etlJobId;
@Column(name = "created_at")
private LocalDateTime createdAt;
}
@Repository
public interface LineageRecordRepository extends JpaRepository<LineageRecord, Long> {
List<LineageRecord> findByTargetTable(String targetTable);
List<LineageRecord> findBySourceTable(String sourceTable);
@Query("SELECT DISTINCT l.sourceTable FROM LineageRecord l WHERE l.targetTable = :targetTable")
List<String> findSourceTablesByTarget(@Param("targetTable") String targetTable);
}
在ETL中自动记录血缘
@Component
public class ETLDataLineageAspect {
@Autowired
private LineageRecordRepository lineageRepo;
@Around("@annotation(etlMapping)")
public Object recordLineage(ProceedingJoinPoint joinPoint,
ETLLineageMapping etlMapping) throws Throwable {
// 执行前记录源数据信息
String sourceTable = etlMapping.sourceTable();
String sourceQuery = etlMapping.sourceQuery();
// 执行ETL
Object result = joinPoint.proceed();
// 执行后记录目标数据信息
String targetTable = etlMapping.targetTable();
// 解析SQL,提取字段级血缘关系
List<FieldLineage> fieldLineages = parseSQLFieldLineage(
sourceQuery, sourceTable, targetTable);
// 保存血缘记录
fieldLineages.forEach(lineage -> {
LineageRecord record = new LineageRecord();
record.setSourceTable(sourceTable);
record.setSourceColumn(lineage.getSourceColumn());
record.setTargetTable(targetTable);
record.setTargetColumn(lineage.getTargetColumn());
record.setTransformationLogic(lineage.getTransformation());
record.setCreatedAt(LocalDateTime.now());
lineageRepo.save(record);
});
return result;
}
}
使用示例
// 1. 查询报表数据来源
GET /api/lineage/sources/report.daily_sales.total_sales
Response:
[
{"nodeId": "dw.daily_order_summary.total_amount", "name": "汇总总金额"},
{"nodeId": "order_db.orders.amount", "name": "订单金额"}
]
// 2. 查询某字段影响范围
GET /api/lineage/impact/order_db.orders.amount
Response:
[
{"nodeId": "dw.daily_order_summary.total_amount", "name": "汇总总金额"},
{"nodeId": "report.daily_sales.total_sales", "name": "总销售额"}
]
这个案例展示了:
- 数据模型设计:定义数据节点和血缘关系
- 血缘追踪逻辑:基于图算法的上下游遍历
- 业务场景应用:电商订单数据从订单表到报表的完整链路
- 自动化记录:通过AOP在ETL过程中自动捕获血缘
- 查询与可视化:提供API供前端展示数据血缘图谱
实际生产环境中,可以结合Apache Atlas、DataHub等开源数据血缘工具,或在Hadoop/Spark生态中使用Atlas插件自动收集血缘信息。