📖 数据密集型设计

元数据管理与数据目录

深入探讨元数据管理、数据目录建设及数据血缘追踪技术

一、元数据概述

元数据是描述数据的数据,在数据密集型应用中,元数据管理是数据治理的核心。通过元数据管理,可以实现数据资产的统一管理、数据血缘的追踪和数据目录的建设。

二、元数据分类

2.1 元数据类型

类型 描述 示例 用途
技术元数据 数据技术属性 表结构、字段类型、索引 数据集成、查询优化
业务元数据 数据业务含义 业务术语、数据定义、指标 数据理解、业务分析
流程元数据 数据处理流程 ETL任务、数据流 数据血缘、影响分析
操作元数据 数据使用情况 访问日志、查询统计 数据治理、权限管理

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 元数据采集

public class MetadataCollector
{
    public async Task CollectDatabaseMetadataAsync(string connectionString, string databaseName)
    {
        var metadata = new Metadata
        {
            DatabaseName = databaseName,
            CollectionTime = DateTime.UtcNow,
            Tables = new List()
        };
        
        using var connection = new SqlConnection(connectionString);
        await connection.OpenAsync();
        
        var tables = await GetTablesAsync(connection, databaseName);
        
        foreach (var table in tables)
        {
            var tableMetadata = await CollectTableMetadataAsync(connection, databaseName, table);
            metadata.Tables.Add(tableMetadata);
        }
        
        return metadata;
    }
    
    private async Task> GetTablesAsync(SqlConnection connection, string databaseName)
    {
        var query = $@"SELECT TABLE_NAME FROM INFORMATION_SCHEMA.TABLES 
                      WHERE TABLE_TYPE = 'BASE TABLE' AND TABLE_CATALOG = '{databaseName}'";
        
        var tables = new List();
        
        using var command = new SqlCommand(query, connection);
        using var reader = await command.ExecuteReaderAsync();
        
        while (await reader.ReadAsync())
        {
            tables.Add(reader.GetString(0));
        }
        
        return tables;
    }
    
    private async Task CollectTableMetadataAsync(SqlConnection connection, string databaseName, string tableName)
    {
        var tableMetadata = new TableMetadata
        {
            TableName = tableName,
            Columns = new List()
        };
        
        var query = $@"SELECT COLUMN_NAME, DATA_TYPE, IS_NULLABLE, COLUMN_DEFAULT 
                      FROM INFORMATION_SCHEMA.COLUMNS 
                      WHERE TABLE_NAME = '{tableName}' AND TABLE_CATALOG = '{databaseName}'";
        
        using var command = new SqlCommand(query, connection);
        using var reader = await command.ExecuteReaderAsync();
        
        while (await reader.ReadAsync())
        {
            tableMetadata.Columns.Add(new ColumnMetadata
            {
                ColumnName = reader.GetString(0),
                DataType = reader.GetString(1),
                IsNullable = reader.GetString(2) == "YES",
                DefaultValue = reader.IsDBNull(3) ? null : reader.GetString(3)
            });
        }
        
        return tableMetadata;
    }
}

public class Metadata
{
    public string DatabaseName { get; set; }
    public DateTime CollectionTime { get; set; }
    public List Tables { get; set; }
}

public class TableMetadata
{
    public string TableName { get; set; }
    public List Columns { get; set; }
}

public class ColumnMetadata
{
    public string ColumnName { get; set; }
    public string DataType { get; set; }
    public bool IsNullable { get; set; }
    public string DefaultValue { get; set; }
}

3.2 元数据存储与管理

public class MetadataRepository
{
    private readonly IMongoCollection _metadataCollection;
    
    public async Task SaveMetadataAsync(Metadata metadata)
    {
        await _metadataCollection.ReplaceOneAsync(
            m => m.DatabaseName == metadata.DatabaseName,
            metadata,
            new ReplaceOptions { IsUpsert = true });
    }
    
    public async Task GetMetadataAsync(string databaseName)
    {
        return await _metadataCollection.Find(m => m.DatabaseName == databaseName).FirstOrDefaultAsync();
    }
    
    public async Task> GetAllMetadataAsync()
    {
        return await _metadataCollection.Find(_ => true).ToListAsync();
    }
    
    public async Task DeleteMetadataAsync(string databaseName)
    {
        await _metadataCollection.DeleteOneAsync(m => m.DatabaseName == databaseName);
    }
    
    public async Task> GetAllDatabaseNamesAsync()
    {
        var metadataList = await _metadataCollection.Find(_ => true).ToListAsync();
        
        return metadataList.Select(m => m.DatabaseName).ToList();
    }
    
    public async Task> GetTablesAsync(string databaseName)
    {
        var metadata = await GetMetadataAsync(databaseName);
        
        return metadata?.Tables ?? new List();
    }
    
    public async Task GetTableAsync(string databaseName, string tableName)
    {
        var metadata = await GetMetadataAsync(databaseName);
        
        return metadata?.Tables.FirstOrDefault(t => t.TableName == tableName);
    }
}

四、数据目录

4.1 数据目录架构

graph TD A[数据目录] --> B[搜索功能] A --> C[数据资产浏览] A --> D[数据详情查看] A --> E[数据使用统计] B --> B1[关键词搜索] B --> B2[高级筛选] B --> B3[全文检索] C --> C1[数据库列表] C --> C2[表列表] C --> C3[字段列表] D --> D1[表结构] D --> D2[字段定义] D --> D3[数据血缘] D --> D4[使用统计]

4.2 数据目录实现

public class DataCatalogService
{
    public async Task SearchAsync(string query, SearchOptions options)
    {
        var results = new SearchResults();
        
        var databases = await SearchDatabasesAsync(query);
        var tables = await SearchTablesAsync(query);
        var columns = await SearchColumnsAsync(query);
        
        results.Databases = databases.Take(options.Limit).ToList();
        results.Tables = tables.Take(options.Limit).ToList();
        results.Columns = columns.Take(options.Limit).ToList();
        
        results.TotalResults = databases.Count + tables.Count + columns.Count;
        
        return results;
    }
    
    private async Task> SearchDatabasesAsync(string query)
    {
        var metadataList = await _metadataRepository.GetAllMetadataAsync();
        
        return metadataList
            .Where(m => m.DatabaseName.Contains(query, StringComparison.OrdinalIgnoreCase))
            .Select(m => new DatabaseResult
            {
                Name = m.DatabaseName,
                TableCount = m.Tables.Count,
                LastUpdated = m.CollectionTime
            })
            .ToList();
    }
    
    private async Task> SearchTablesAsync(string query)
    {
        var results = new List();
        
        var metadataList = await _metadataRepository.GetAllMetadataAsync();
        
        foreach (var metadata in metadataList)
        {
            var matchingTables = metadata.Tables
                .Where(t => t.TableName.Contains(query, StringComparison.OrdinalIgnoreCase))
                .Select(t => new TableResult
                {
                    DatabaseName = metadata.DatabaseName,
                    TableName = t.TableName,
                    ColumnCount = t.Columns.Count,
                    LastUpdated = metadata.CollectionTime
                });
            
            results.AddRange(matchingTables);
        }
        
        return results;
    }
    
    private async Task> SearchColumnsAsync(string query)
    {
        var results = new List();
        
        var metadataList = await _metadataRepository.GetAllMetadataAsync();
        
        foreach (var metadata in metadataList)
        {
            foreach (var table in metadata.Tables)
            {
                var matchingColumns = table.Columns
                    .Where(c => c.ColumnName.Contains(query, StringComparison.OrdinalIgnoreCase))
                    .Select(c => new ColumnResult
                    {
                        DatabaseName = metadata.DatabaseName,
                        TableName = table.TableName,
                        ColumnName = c.ColumnName,
                        DataType = c.DataType
                    });
                
                results.AddRange(matchingColumns);
            }
        }
        
        return results;
    }
    
    public async Task GetAssetDetailAsync(string databaseName, string tableName = null, string columnName = null)
    {
        var metadata = await _metadataRepository.GetMetadataAsync(databaseName);
        
        if (metadata == null)
        {
            return null;
        }
        
        if (!string.IsNullOrEmpty(columnName) && !string.IsNullOrEmpty(tableName))
        {
            var table = metadata.Tables.FirstOrDefault(t => t.TableName == tableName);
            var column = table?.Columns.FirstOrDefault(c => c.ColumnName == columnName);
            
            return new DataAssetDetail
            {
                Type = AssetType.Column,
                DatabaseName = databaseName,
                TableName = tableName,
                ColumnName = columnName,
                Column = column
            };
        }
        
        if (!string.IsNullOrEmpty(tableName))
        {
            var table = metadata.Tables.FirstOrDefault(t => t.TableName == tableName);
            
            return new DataAssetDetail
            {
                Type = AssetType.Table,
                DatabaseName = databaseName,
                TableName = tableName,
                Table = table
            };
        }
        
        return new DataAssetDetail
        {
            Type = AssetType.Database,
            DatabaseName = databaseName,
            Database = metadata
        };
    }
}

public enum AssetType { Database, Table, Column }

五、数据血缘追踪

5.1 数据血缘概念

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

5.2 数据血缘实现

public class DataLineageService
{
    public async Task BuildLineageGraphAsync(string targetTableName)
    {
        var graph = new LineageGraph();
        
        var sources = await FindSourceTablesAsync(targetTableName);
        
        foreach (var source in sources)
        {
            graph.Nodes.Add(new LineageNode
            {
                TableName = source.TableName,
                Type = LineageNodeType.Source
            });
            
            graph.Edges.Add(new LineageEdge
            {
                SourceTableName = source.TableName,
                TargetTableName = targetTableName,
                RelationshipType = source.RelationshipType
            });
            
            var upstreamSources = await FindSourceTablesAsync(source.TableName);
            
            foreach (var upstream in upstreamSources)
            {
                graph.Nodes.Add(new LineageNode
                {
                    TableName = upstream.TableName,
                    Type = LineageNodeType.Source
                });
                
                graph.Edges.Add(new LineageEdge
                {
                    SourceTableName = upstream.TableName,
                    TargetTableName = source.TableName,
                    RelationshipType = upstream.RelationshipType
                });
            }
        }
        
        graph.Nodes.Add(new LineageNode
        {
            TableName = targetTableName,
            Type = LineageNodeType.Target
        });
        
        return graph;
    }
    
    private async Task> FindSourceTablesAsync(string targetTableName)
    {
        var lineageRecords = await _lineageRepository.GetByTargetTableAsync(targetTableName);
        
        return lineageRecords.Select(r => new LineageSource
        {
            TableName = r.SourceTableName,
            RelationshipType = r.RelationshipType
        }).ToList();
    }
    
    public async Task> FindDownstreamTablesAsync(string sourceTableName)
    {
        var lineageRecords = await _lineageRepository.GetBySourceTableAsync(sourceTableName);
        
        return lineageRecords.Select(r => r.TargetTableName).Distinct().ToList();
    }
    
    public async Task> FindUpstreamTablesAsync(string targetTableName)
    {
        var lineageRecords = await _lineageRepository.GetByTargetTableAsync(targetTableName);
        
        return lineageRecords.Select(r => r.SourceTableName).Distinct().ToList();
    }
    
    public async Task AnalyzeImpactAsync(string tableName, string columnName = null)
    {
        var downstreamTables = await FindDownstreamTablesAsync(tableName);
        
        return new ImpactAnalysis
        {
            SourceTable = tableName,
            SourceColumn = columnName,
            AffectedTables = downstreamTables,
            ImpactLevel = downstreamTables.Count > 10 ? ImpactLevel.High : 
                          downstreamTables.Count > 3 ? ImpactLevel.Medium : ImpactLevel.Low
        };
    }
}

public class LineageGraph
{
    public List Nodes { get; set; } = new();
    public List Edges { get; set; } = new();
}

public class LineageNode
{
    public string TableName { get; set; }
    public LineageNodeType Type { get; set; }
}

public class LineageEdge
{
    public string SourceTableName { get; set; }
    public string TargetTableName { get; set; }
    public RelationshipType RelationshipType { get; set; }
}

public enum LineageNodeType { Source, Intermediate, Target }
public enum RelationshipType { Direct, Indirect, Transformed }
public enum ImpactLevel { Low, Medium, High }

5.3 数据血缘可视化

graph LR A[订单源表] -->|ETL| B[订单DWD层] B -->|聚合| C[订单DWS层] A -->|关联| D[用户DWD层] D -->|聚合| E[用户DWS层] C -->|关联| F[销售ADS层] E -->|关联| F F -->|查询| G[BI报表] F -->|查询| H[数据接口] style A fill:#90EE90 style B fill:#87CEEB style C fill:#87CEEB style D fill:#87CEEB style E fill:#87CEEB style F fill:#DDA0DD style G fill:#FFD700 style H fill:#FFD700

六、数据字典

6.1 数据字典设计

public class DataDictionaryService
{
    public async Task GetEntryAsync(string databaseName, string tableName, string columnName)
    {
        return await _dictionaryRepository.GetEntryAsync(databaseName, tableName, columnName);
    }
    
    public async Task> GetEntriesAsync(string databaseName, string tableName = null)
    {
        return await _dictionaryRepository.GetEntriesAsync(databaseName, tableName);
    }
    
    public async Task SaveEntryAsync(DataDictionaryEntry entry)
    {
        await _dictionaryRepository.SaveEntryAsync(entry);
    }
    
    public async Task DeleteEntryAsync(string databaseName, string tableName, string columnName)
    {
        await _dictionaryRepository.DeleteEntryAsync(databaseName, tableName, columnName);
    }
    
    public async Task> SearchEntriesAsync(string keyword)
    {
        return await _dictionaryRepository.SearchEntriesAsync(keyword);
    }
}

public class DataDictionaryEntry
{
    public string DatabaseName { get; set; }
    public string TableName { get; set; }
    public string ColumnName { get; set; }
    public string BusinessDefinition { get; set; }
    public string TechnicalDefinition { get; set; }
    public string Owner { get; set; }
    public string DataClassification { get; set; }
    public DateTime CreatedAt { get; set; }
    public DateTime UpdatedAt { get; set; }
}

七、元数据管理最佳实践

7.1 元数据管理策略

策略 描述 实现方式
自动采集 定时自动采集元数据 调度任务、CDC
手动补充 人工补充业务元数据 数据字典管理界面
版本管理 跟踪元数据变更 版本控制、变更记录
权限管理 控制元数据访问 RBAC、行级安全
数据血缘 追踪数据流转 血缘分析、影响分析

7.2 数据目录建设建议

  • 定义数据资产分类标准
  • 实现元数据自动采集
  • 提供便捷的搜索功能
  • 建设数据血缘追踪系统
  • 建立数据字典管理机制

7.3 数据血缘最佳实践

public class LineageBestPractices
{
    public LineageGraph OptimizeLineageGraph(LineageGraph graph)
    {
        var uniqueNodes = graph.Nodes
            .GroupBy(n => n.TableName)
            .Select(g => g.First())
            .ToList();
        
        var uniqueEdges = graph.Edges
            .GroupBy(e => new { e.SourceTableName, e.TargetTableName })
            .Select(g => g.First())
            .ToList();
        
        return new LineageGraph
        {
            Nodes = uniqueNodes,
            Edges = uniqueEdges
        };
    }
    
    public async Task> FindCriticalTablesAsync()
    {
        var allTables = await _metadataRepository.GetAllTableNamesAsync();
        
        return allTables
            .Select(tableName => new
            {
                TableName = tableName,
                DownstreamCount = _lineageService.FindDownstreamTablesAsync(tableName).Result.Count
            })
            .OrderByDescending(t => t.DownstreamCount)
            .Take(10)
            .Select(t => t.TableName)
            .ToList();
    }
}

八、总结

元数据管理与数据目录是数据治理的核心。通过元数据采集、存储和管理,能够实现数据资产的统一管理。数据目录提供了数据资产的搜索和浏览功能,数据血缘追踪能够帮助理解数据的流转过程。建设完善的元数据管理体系,能够提升数据治理水平,支持数据驱动的决策。