Java数据血缘案例

wen java案例 2

Java 数据血缘案例

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

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": "总销售额"}
]

这个案例展示了:

  1. 数据模型设计:定义数据节点和血缘关系
  2. 血缘追踪逻辑:基于图算法的上下游遍历
  3. 业务场景应用:电商订单数据从订单表到报表的完整链路
  4. 自动化记录:通过AOP在ETL过程中自动捕获血缘
  5. 查询与可视化:提供API供前端展示数据血缘图谱

实际生产环境中,可以结合Apache Atlas、DataHub等开源数据血缘工具,或在Hadoop/Spark生态中使用Atlas插件自动收集血缘信息。

抱歉,评论功能暂时关闭!