📖 数据密集型设计

数据血缘追踪与影响分析

深入探讨数据血缘追踪技术、影响分析方法及数据变更管理

一、数据血缘概述

数据血缘是描述数据从产生到消费的完整路径,包括数据的来源、处理过程和去向。数据血缘追踪是数据治理的重要组成部分,能够帮助理解数据的流转过程,支持影响分析和变更管理。

二、数据血缘类型

2.1 数据血缘分类

类型 描述 示例 用途
技术血缘 数据技术层面流转 表到表、字段到字段 数据集成、质量监控
业务血缘 数据业务层面流转 指标计算、报表生成 业务分析、指标溯源
流程血缘 数据处理流程 ETL任务、数据流 任务调度、故障排查
对象血缘 数据对象依赖 视图、存储过程、API 对象管理、变更分析

2.2 数据血缘追踪架构

graph TD A[数据源] --> B[血缘采集] B --> C[血缘存储] C --> C1[技术血缘] C --> C2[业务血缘] C --> C3[流程血缘] C --> C4[对象血缘] C --> D[血缘分析] D --> E[向上追踪] D --> F[向下追踪] D --> G[影响分析] D --> H[变更管理]

三、数据血缘采集

3.1 SQL解析血缘采集

public class SqlLineageExtractor
{
    public LineageResult ExtractLineage(string sql)
    {
        var result = new LineageResult();
        
        var parser = new SqlParser(sql);
        var statements = parser.Parse();
        
        foreach (var statement in statements)
        {
            if (statement is SelectStatement selectStmt)
            {
                result.Sources.AddRange(ExtractSources(selectStmt));
                result.Targets.AddRange(ExtractTargets(selectStmt));
                result.ColumnMappings.AddRange(ExtractColumnMappings(selectStmt));
            }
            
            if (statement is InsertStatement insertStmt)
            {
                result.Targets.Add(insertStmt.TargetTable);
                result.Sources.AddRange(ExtractSources(insertStmt.Source));
                result.ColumnMappings.AddRange(ExtractColumnMappings(insertStmt));
            }
            
            if (statement is CreateViewStatement viewStmt)
            {
                result.Targets.Add(viewStmt.ViewName);
                result.Sources.AddRange(ExtractSources(viewStmt.Query));
            }
        }
        
        return result;
    }
    
    private List ExtractSources(SqlNode node)
    {
        var sources = new List();
        
        if (node is FromClause fromClause)
        {
            sources.AddRange(fromClause.Tables.Select(t => t.Name));
        }
        
        if (node is JoinClause joinClause)
        {
            sources.Add(joinClause.Table.Name);
            sources.AddRange(ExtractSources(joinClause.Condition));
        }
        
        if (node is SubQuery subQuery)
        {
            sources.AddRange(ExtractSources(subQuery.Query));
        }
        
        return sources;
    }
    
    private List ExtractTargets(SqlNode node)
    {
        var targets = new List();
        
        if (node is SelectStatement selectStmt)
        {
            foreach (var column in selectStmt.Columns)
            {
                if (column.Alias != null)
                {
                    targets.Add(column.Alias);
                }
            }
        }
        
        return targets;
    }
    
    private List ExtractColumnMappings(SqlNode node)
    {
        var mappings = new List();
        
        if (node is SelectStatement selectStmt)
        {
            foreach (var column in selectStmt.Columns)
            {
                mappings.Add(new ColumnMapping
                {
                    SourceColumn = column.Expression.ToString(),
                    TargetColumn = column.Alias ?? column.Expression.ToString()
                });
            }
        }
        
        return mappings;
    }
}

public class LineageResult
{
    public List Sources { get; set; } = new();
    public List Targets { get; set; } = new();
    public List ColumnMappings { get; set; } = new();
}

public class ColumnMapping
{
    public string SourceColumn { get; set; }
    public string TargetColumn { get; set; }
}

3.2 ETL任务血缘采集

public class EtlLineageExtractor
{
    public async Task ExtractFromEtlTaskAsync(EtlTask task)
    {
        var result = new LineageResult();
        
        foreach (var source in task.Sources)
        {
            result.Sources.Add(source.TableName);
            
            var sourceMetadata = await _metadataService.GetTableMetadataAsync(source.ConnectionString, source.TableName);
            
            foreach (var column in sourceMetadata.Columns)
            {
                result.ColumnMappings.Add(new ColumnMapping
                {
                    SourceColumn = $"{source.TableName}.{column.ColumnName}",
                    TargetColumn = FindTargetColumn(task, column.ColumnName)
                });
            }
        }
        
        foreach (var target in task.Targets)
        {
            result.Targets.Add(target.TableName);
        }
        
        return result;
    }
    
    private string FindTargetColumn(EtlTask task, string sourceColumnName)
    {
        var mapping = task.ColumnMappings.FirstOrDefault(m => m.SourceColumn == sourceColumnName);
        
        return mapping?.TargetColumn ?? sourceColumnName;
    }
    
    public async Task> ExtractFromAllEtlTasksAsync()
    {
        var results = new List();
        
        var tasks = await _etlRepository.GetAllTasksAsync();
        
        foreach (var task in tasks)
        {
            var result = await ExtractFromEtlTaskAsync(task);
            results.Add(result);
        }
        
        return results;
    }
}

public class EtlTask
{
    public string TaskName { get; set; }
    public List Sources { get; set; } = new();
    public List Targets { get; set; } = new();
    public List ColumnMappings { get; set; } = new();
}

public class EtlSource
{
    public string ConnectionString { get; set; }
    public string TableName { get; set; }
}

public class EtlTarget
{
    public string ConnectionString { get; set; }
    public string TableName { get; set; }
}

四、数据血缘存储与查询

4.1 血缘图存储

public class LineageGraphRepository
{
    private readonly GraphDatabase _graphDb;
    
    public async Task SaveLineageAsync(LineageResult result)
    {
        using var transaction = await _graphDb.BeginTransactionAsync();
        
        foreach (var source in result.Sources)
        {
            await CreateOrGetNodeAsync(transaction, source, NodeType.Source);
        }
        
        foreach (var target in result.Targets)
        {
            await CreateOrGetNodeAsync(transaction, target, NodeType.Target);
        }
        
        foreach (var mapping in result.ColumnMappings)
        {
            var sourceNode = await CreateOrGetNodeAsync(transaction, mapping.SourceColumn, NodeType.Column);
            var targetNode = await CreateOrGetNodeAsync(transaction, mapping.TargetColumn, NodeType.Column);
            
            await CreateEdgeAsync(transaction, sourceNode, targetNode, EdgeType.DataFlow);
        }
        
        await transaction.CommitAsync();
    }
    
    private async Task CreateOrGetNodeAsync(ITransaction transaction, string name, NodeType type)
    {
        var existingNode = await transaction.GetNodeByNameAsync(name);
        
        if (existingNode != null)
        {
            return existingNode;
        }
        
        return await transaction.CreateNodeAsync(name, type);
    }
    
    private async Task CreateEdgeAsync(ITransaction transaction, Node source, Node target, EdgeType type)
    {
        return await transaction.CreateEdgeAsync(source, target, type);
    }
    
    public async Task> FindUpstreamNodesAsync(string nodeName)
    {
        var node = await _graphDb.GetNodeByNameAsync(nodeName);
        
        if (node == null)
        {
            return new List();
        }
        
        var upstreamEdges = await _graphDb.GetIncomingEdgesAsync(node);
        
        return upstreamEdges.Select(e => e.SourceNode.Name).ToList();
    }
    
    public async Task> FindDownstreamNodesAsync(string nodeName)
    {
        var node = await _graphDb.GetNodeByNameAsync(nodeName);
        
        if (node == null)
        {
            return new List();
        }
        
        var downstreamEdges = await _graphDb.GetOutgoingEdgesAsync(node);
        
        return downstreamEdges.Select(e => e.TargetNode.Name).ToList();
    }
}

public enum NodeType { Source, Target, Column, Table, Database, Task }
public enum EdgeType { DataFlow, Dependency, Transformation, Aggregation }

4.2 血缘查询接口

public class LineageQueryService
{
    public async Task GetFullLineagePathAsync(string targetNodeName)
    {
        var path = new LineagePath();
        
        var upstreamNodes = await FindAllUpstreamNodesAsync(targetNodeName);
        var downstreamNodes = await FindAllDownstreamNodesAsync(targetNodeName);
        
        path.Upstream = upstreamNodes;
        path.Downstream = downstreamNodes;
        path.Target = targetNodeName;
        
        return path;
    }
    
    private async Task> FindAllUpstreamNodesAsync(string nodeName, int depth = 5)
    {
        if (depth <= 0)
        {
            return new List();
        }
        
        var nodes = new List();
        
        var directUpstream = await _graphRepository.FindUpstreamNodesAsync(nodeName);
        
        foreach (var upstream in directUpstream)
        {
            nodes.Add(new LineageNode
            {
                Name = upstream,
                Depth = 5 - depth + 1,
                Children = await FindAllUpstreamNodesAsync(upstream, depth - 1)
            });
        }
        
        return nodes;
    }
    
    private async Task> FindAllDownstreamNodesAsync(string nodeName, int depth = 5)
    {
        if (depth <= 0)
        {
            return new List();
        }
        
        var nodes = new List();
        
        var directDownstream = await _graphRepository.FindDownstreamNodesAsync(nodeName);
        
        foreach (var downstream in directDownstream)
        {
            nodes.Add(new LineageNode
            {
                Name = downstream,
                Depth = 5 - depth + 1,
                Children = await FindAllDownstreamNodesAsync(downstream, depth - 1)
            });
        }
        
        return nodes;
    }
    
    public async Task> FindAllAffectedObjectsAsync(string sourceNodeName)
    {
        var affected = new HashSet();
        
        var downstreamNodes = await FindAllDownstreamNodesAsync(sourceNodeName);
        
        CollectAffectedObjects(downstreamNodes, affected);
        
        return affected.ToList();
    }
    
    private void CollectAffectedObjects(List nodes, HashSet affected)
    {
        foreach (var node in nodes)
        {
            affected.Add(node.Name);
            CollectAffectedObjects(node.Children, affected);
        }
    }
}

public class LineagePath
{
    public string Target { get; set; }
    public List Upstream { get; set; } = new();
    public List Downstream { get; set; } = new();
}

public class LineageNode
{
    public string Name { get; set; }
    public int Depth { get; set; }
    public List Children { get; set; } = new();
}

五、影响分析

5.1 影响分析流程

sequenceDiagram participant User as 用户 participant Service as 影响分析服务 participant Repository as 血缘存储 participant Alert as 告警系统 User->>Service: 发起变更请求(表名, 变更类型) Service->>Repository: 查询下游依赖 Repository-->>Service: 返回下游对象列表 Service->>Service: 分析影响范围 Service->>Service: 评估影响级别 Service-->>User: 返回影响分析报告 Service->>Alert: 发送影响通知

5.2 影响分析实现

public class ImpactAnalysisService
{
    public async Task AnalyzeImpactAsync(string objectName, ChangeType changeType)
    {
        var report = new ImpactAnalysisReport
        {
            ObjectName = objectName,
            ChangeType = changeType,
            AnalysisTime = DateTime.UtcNow
        };
        
        var downstreamObjects = await _lineageQueryService.FindAllAffectedObjectsAsync(objectName);
        
        report.AffectedObjects = downstreamObjects;
        report.AffectedObjectCount = downstreamObjects.Count;
        
        foreach (var obj in downstreamObjects)
        {
            var objectType = await DetermineObjectTypeAsync(obj);
            report.AffectedByType[objectType] = report.AffectedByType.GetValueOrDefault(objectType, 0) + 1;
        }
        
        report.ImpactLevel = CalculateImpactLevel(report);
        report.RiskAssessment = AssessRisk(report);
        
        return report;
    }
    
    private async Task DetermineObjectTypeAsync(string objectName)
    {
        if (objectName.Contains(".sql")) return ObjectType.Query;
        if (objectName.Contains(".etl")) return ObjectType.EtlTask;
        if (objectName.Contains(".view")) return ObjectType.View;
        if (objectName.Contains(".report")) return ObjectType.Report;
        if (objectName.Contains(".api")) return ObjectType.Api;
        
        return ObjectType.Table;
    }
    
    private ImpactLevel CalculateImpactLevel(ImpactAnalysisReport report)
    {
        if (report.AffectedObjectCount > 100) return ImpactLevel.Critical;
        if (report.AffectedObjectCount > 30) return ImpactLevel.High;
        if (report.AffectedObjectCount > 10) return ImpactLevel.Medium;
        if (report.AffectedObjectCount > 0) return ImpactLevel.Low;
        
        return ImpactLevel.None;
    }
    
    private RiskAssessment AssessRisk(ImpactAnalysisReport report)
    {
        var risk = new RiskAssessment();
        
        if (report.AffectedByType.ContainsKey(ObjectType.Report) && report.AffectedByType[ObjectType.Report] > 5)
        {
            risk.BusinessRisk = RiskLevel.High;
            risk.Recommendations.Add("业务报表受影响,请通知业务部门");
        }
        
        if (report.AffectedByType.ContainsKey(ObjectType.Api) && report.AffectedByType[ObjectType.Api] > 10)
        {
            risk.SystemRisk = RiskLevel.High;
            risk.Recommendations.Add("API接口受影响,请通知开发团队");
        }
        
        if (report.ImpactLevel == ImpactLevel.Critical)
        {
            risk.OverallRisk = RiskLevel.High;
            risk.Recommendations.Add("建议暂停变更,进行详细评估");
        }
        
        return risk;
    }
    
    public async Task> GetAffectedOwnersAsync(ImpactAnalysisReport report)
    {
        var owners = new HashSet();
        
        foreach (var obj in report.AffectedObjects)
        {
            var owner = await _metadataService.GetObjectOwnerAsync(obj);
            if (!string.IsNullOrEmpty(owner))
            {
                owners.Add(owner);
            }
        }
        
        return owners.ToList();
    }
}

public class ImpactAnalysisReport
{
    public string ObjectName { get; set; }
    public ChangeType ChangeType { get; set; }
    public DateTime AnalysisTime { get; set; }
    public List AffectedObjects { get; set; } = new();
    public int AffectedObjectCount { get; set; }
    public Dictionary AffectedByType { get; set; } = new();
    public ImpactLevel ImpactLevel { get; set; }
    public RiskAssessment RiskAssessment { get; set; }
}

public enum ChangeType { AddColumn, DropColumn, AlterColumn, RenameColumn, DropTable, TruncateTable }
public enum ObjectType { Table, View, Query, EtlTask, Report, Api, Dashboard }
public enum ImpactLevel { None, Low, Medium, High, Critical }
public enum RiskLevel { Low, Medium, High, Critical }

5.3 变更管理流程

public class ChangeManagementService
{
    public async Task CreateChangeRequestAsync(ChangeRequestDto dto)
    {
        var request = new ChangeRequest
        {
            Id = Guid.NewGuid().ToString(),
            ObjectName = dto.ObjectName,
            ChangeType = dto.ChangeType,
            Description = dto.Description,
            Requester = dto.Requester,
            Status = ChangeStatus.PendingAnalysis,
            CreatedAt = DateTime.UtcNow
        };
        
        await _changeRepository.SaveChangeRequestAsync(request);
        
        return request;
    }
    
    public async Task AnalyzeChangeAsync(string requestId)
    {
        var request = await _changeRepository.GetChangeRequestAsync(requestId);
        
        if (request == null)
        {
            throw new NotFoundException("变更请求不存在");
        }
        
        var impactReport = await _impactAnalysisService.AnalyzeImpactAsync(request.ObjectName, request.ChangeType);
        
        request.ImpactReport = impactReport;
        request.Status = ChangeStatus.Analyzed;
        request.AnalyzedAt = DateTime.UtcNow;
        
        await _changeRepository.SaveChangeRequestAsync(request);
        
        return request;
    }
    
    public async Task ApproveChangeAsync(string requestId, string approver)
    {
        var request = await _changeRepository.GetChangeRequestAsync(requestId);
        
        if (request == null)
        {
            throw new NotFoundException("变更请求不存在");
        }
        
        if (request.ImpactReport.ImpactLevel == ImpactLevel.Critical)
        {
            throw new ValidationException("严重影响的变更需要高级审批");
        }
        
        request.Status = ChangeStatus.Approved;
        request.Approver = approver;
        request.ApprovedAt = DateTime.UtcNow;
        
        await _changeRepository.SaveChangeRequestAsync(request);
        
        return request;
    }
    
    public async Task ExecuteChangeAsync(string requestId)
    {
        var request = await _changeRepository.GetChangeRequestAsync(requestId);
        
        if (request == null)
        {
            throw new NotFoundException("变更请求不存在");
        }
        
        if (request.Status != ChangeStatus.Approved)
        {
            throw new ValidationException("变更请求未批准");
        }
        
        try
        {
            await _dataService.ExecuteChangeAsync(request.ObjectName, request.ChangeType);
            
            request.Status = ChangeStatus.Completed;
            request.ExecutedAt = DateTime.UtcNow;
            
            await _lineageService.RefreshLineageAsync(request.ObjectName);
        }
        catch (Exception ex)
        {
            request.Status = ChangeStatus.Failed;
            request.ErrorMessage = ex.Message;
            request.FailedAt = DateTime.UtcNow;
        }
        
        await _changeRepository.SaveChangeRequestAsync(request);
        
        return request;
    }
}

public enum ChangeStatus { PendingAnalysis, Analyzed, Approved, Rejected, Executing, Completed, Failed }

六、数据血缘可视化

6.1 血缘图渲染

public class LineageVisualizationService
{
    public string GenerateLineageGraphSvg(LineagePath path)
    {
        var svgBuilder = new StringBuilder();
        
        svgBuilder.AppendLine("");
        
        svgBuilder.AppendLine("");
        svgBuilder.AppendLine("");
        svgBuilder.AppendLine("");
        svgBuilder.AppendLine("");
        svgBuilder.AppendLine("");
        
        var yOffset = 400;
        var xOffset = 600;
        var nodeWidth = 150;
        var nodeHeight = 40;
        var levelSpacing = 180;
        
        DrawNode(svgBuilder, path.Target, xOffset, yOffset, "#ff6b6b", nodeWidth, nodeHeight);
        
        DrawUpstreamNodes(svgBuilder, path.Upstream, xOffset, yOffset - levelSpacing, nodeWidth, nodeHeight);
        
        DrawDownstreamNodes(svgBuilder, path.Downstream, xOffset, yOffset + levelSpacing, nodeWidth, nodeHeight);
        
        svgBuilder.AppendLine("");
        
        return svgBuilder.ToString();
    }
    
    private void DrawNode(StringBuilder svg, string name, int x, int y, string color, int width, int height)
    {
        var halfWidth = width / 2;
        var halfHeight = height / 2;
        
        svg.AppendLine($"");
        svg.AppendLine($"{name}");
    }
    
    private void DrawUpstreamNodes(StringBuilder svg, List nodes, int centerX, int y, int width, int height)
    {
        if (nodes.Count == 0) return;
        
        var spacing = 200;
        var startX = centerX - ((nodes.Count - 1) * spacing) / 2;
        
        for (int i = 0; i < nodes.Count; i++)
        {
            var x = startX + i * spacing;
            
            DrawNode(svg, nodes[i].Name, x, y, "#4ecdc4", width, height);
            
            svg.AppendLine($"");
            
            if (nodes[i].Children.Count > 0)
            {
                DrawUpstreamNodes(svg, nodes[i].Children, x, y - 180, width, height);
            }
        }
    }
    
    private void DrawDownstreamNodes(StringBuilder svg, List nodes, int centerX, int y, int width, int height)
    {
        if (nodes.Count == 0) return;
        
        var spacing = 200;
        var startX = centerX - ((nodes.Count - 1) * spacing) / 2;
        
        for (int i = 0; i < nodes.Count; i++)
        {
            var x = startX + i * spacing;
            
            DrawNode(svg, nodes[i].Name, x, y, "#45b7d1", width, height);
            
            svg.AppendLine($"");
            
            if (nodes[i].Children.Count > 0)
            {
                DrawDownstreamNodes(svg, nodes[i].Children, x, y + 180, width, height);
            }
        }
    }
}

七、数据血缘最佳实践

7.1 血缘管理策略

策略 描述 实现方式
自动采集 自动解析SQL和ETL任务 SQL解析器、ETL监控
实时更新 变更后及时更新血缘 触发器、CDC
图形存储 使用图数据库存储 Neo4j、JanusGraph
影响分析 变更前评估影响 影响分析服务
可视化 直观展示血缘关系 血缘图、流程图

7.2 影响分析最佳实践

  • 变更前进行影响分析
  • 评估影响级别和风险
  • 通知受影响的所有者
  • 建立变更审批流程
  • 变更后验证血缘完整性

八、总结

数据血缘追踪与影响分析是数据治理的重要组成部分。通过SQL解析和ETL任务分析,能够自动采集数据血缘。使用图数据库存储血缘关系,支持高效的查询和分析。影响分析能够帮助评估数据变更的影响范围和风险,支持变更管理流程。建设完善的数据血缘追踪系统,能够提升数据治理水平,保障数据变更的安全性。