📖 数据密集型设计

数据仓库架构与设计

深入探讨数据仓库架构与维度建模技术

一、数据仓库概述

数据仓库是用于存储和分析企业历史数据的系统,能够支持决策分析和业务洞察。在数据密集型应用中,数据仓库是构建数据分析平台的核心组件。

二、数据仓库架构

2.1 经典数据仓库架构

graph TD A[源系统] --> B[ODS层] B --> C[DWD层] C --> D[DWS层] D --> E[ADS层] E --> F[数据应用] A --> A1[业务数据库] A --> A2[日志文件] A --> A3[API接口] B --> B1[原始数据] B --> B2[增量数据] C --> C1[明细数据] C --> C2[维度表] D --> D1[汇总数据] D --> D2[主题域数据] E --> E1[指标数据] E --> E2[报表数据]

2.2 数据仓库分层架构

层级 名称 作用 数据特点
ODS 操作数据存储 存储原始数据 原始、增量、保留历史
DWD 明细数据层 清洗、转换原始数据 干净、关联、业务事实
DWS 汇总数据层 主题域汇总 聚合、宽表、主题域
ADS 应用数据层 面向应用的数据 指标、报表、分析

三、维度建模

3.1 星型模型

graph TD A[事实表] --> B[时间维度] A --> C[产品维度] A --> D[用户维度] A --> E[地域维度] A --> F[渠道维度] A --> A1[订单ID] A --> A2[金额] A --> A3[数量] B --> B1[日期] B --> B2[月份] B --> B3[季度] B --> B4[年份] C --> C1[产品ID] C --> C2[产品名称] C --> C3[产品类别] D --> D1[用户ID] D --> D2[用户名] D --> D3[用户等级] E --> E1[地域ID] E --> E2[省份] E --> E3[城市] F --> F1[渠道ID] F --> F2[渠道名称] F --> F3[渠道类型]

3.2 维度建模设计

public class DimensionModelingService
{
    public async Task CreateFactTableAsync(string tableName, List columns)
    {
        var columnDefinitions = string.Join(", ", columns.Select(c => 
            $"{c.Name} {c.Type}"));
        
        var sql = $"CREATE TABLE {tableName} ({columnDefinitions})";
        await _dataWarehouseClient.ExecuteAsync(sql);
    }
    
    public async Task CreateDimensionTableAsync(string tableName, List columns)
    {
        var columnDefinitions = string.Join(", ", columns.Select(c => 
            $"{c.Name} {c.Type}"));
        
        var sql = $"CREATE TABLE {tableName} ({columnDefinitions}, PRIMARY KEY(id))";
        await _dataWarehouseClient.ExecuteAsync(sql);
    }
    
    public async Task CreateStarSchemaAsync(string factTableName, List dimensionTableNames)
    {
        foreach (var dimensionTable in dimensionTableNames)
        {
            var constraintName = $"fk_{factTableName}_{dimensionTable}";
            var sql = $"ALTER TABLE {factTableName} ADD CONSTRAINT {constraintName} FOREIGN KEY ({dimensionTable}_id) REFERENCES {dimensionTable}(id)";
            await _dataWarehouseClient.ExecuteAsync(sql);
        }
    }
}

public class FactColumn
{
    public string Name { get; set; }
    public string Type { get; set; }
    public bool IsMeasure { get; set; }
    public bool IsForeignKey { get; set; }
}

public class DimensionColumn
{
    public string Name { get; set; }
    public string Type { get; set; }
    public bool IsPrimaryKey { get; set; }
}

3.3 缓慢变化维度(SCD)

public class SlowlyChangingDimensionService
{
    public async Task HandleSCDType1Async(string tableName, DimensionRecord record)
    {
        var sql = $"UPDATE {tableName} SET {BuildUpdateClause(record)} WHERE id = @Id";
        await _dataWarehouseClient.ExecuteAsync(sql, record);
    }
    
    public async Task HandleSCDType2Async(string tableName, DimensionRecord record)
    {
        var expireSql = $"UPDATE {tableName} SET end_date = @CurrentDate WHERE id = @Id AND end_date IS NULL";
        await _dataWarehouseClient.ExecuteAsync(expireSql, new { record.Id, CurrentDate = DateTime.Now });
        
        record.StartDate = DateTime.Now;
        record.EndDate = null;
        
        var insertSql = $"INSERT INTO {tableName} ({BuildInsertColumns(record)}) VALUES ({BuildInsertValues(record)})";
        await _dataWarehouseClient.ExecuteAsync(insertSql, record);
    }
    
    public async Task HandleSCDType3Async(string tableName, DimensionRecord record)
    {
        var sql = $"UPDATE {tableName} SET previous_value = current_value, current_value = @NewValue WHERE id = @Id";
        await _dataWarehouseClient.ExecuteAsync(sql, record);
    }
}

public class DimensionRecord
{
    public long Id { get; set; }
    public string Name { get; set; }
    public DateTime StartDate { get; set; }
    public DateTime? EndDate { get; set; }
    public int Version { get; set; }
}

四、ETL流程

4.1 ETL架构

graph TD A[抽取 Extract] --> B[转换 Transform] B --> C[加载 Load] A --> A1[全量抽取] A --> A2[增量抽取] A --> A3[CDC抽取] B --> B1[数据清洗] B --> B2[数据转换] B --> B3[数据验证] C --> C1[全量加载] C --> C2[增量加载] C --> C3[UPSERT]

4.2 ETL实现

public class EtlService
{
    public async Task ExecuteEtlAsync(EtlJobConfig config)
    {
        await ExtractAsync(config);
        await TransformAsync(config);
        await LoadAsync(config);
    }
    
    private async Task ExtractAsync(EtlJobConfig config)
    {
        switch (config.ExtractType)
        {
            case ExtractType.Full:
                await ExtractFullAsync(config);
                break;
            case ExtractType.Incremental:
                await ExtractIncrementalAsync(config);
                break;
            case ExtractType.Cdc:
                await ExtractCdcAsync(config);
                break;
        }
    }
    
    private async Task TransformAsync(EtlJobConfig config)
    {
        var data = await _storageService.ReadAsync(config.StagingPath);
        
        data = await _dataCleaner.CleanAsync(data);
        data = await _dataTransformer.TransformAsync(data, config.TransformRules);
        data = await _dataValidator.ValidateAsync(data);
        
        await _storageService.WriteAsync(config.StagingPath, data);
    }
    
    private async Task LoadAsync(EtlJobConfig config)
    {
        var data = await _storageService.ReadAsync(config.StagingPath);
        
        switch (config.LoadType)
        {
            case LoadType.Full:
                await LoadFullAsync(config, data);
                break;
            case LoadType.Incremental:
                await LoadIncrementalAsync(config, data);
                break;
            case LoadType.Upsert:
                await LoadUpsertAsync(config, data);
                break;
        }
    }
}

public class EtlJobConfig
{
    public string JobName { get; set; }
    public ExtractType ExtractType { get; set; }
    public LoadType LoadType { get; set; }
    public string SourcePath { get; set; }
    public string StagingPath { get; set; }
    public string TargetTable { get; set; }
    public List TransformRules { get; set; }
}

public enum ExtractType { Full, Incremental, Cdc }
public enum LoadType { Full, Incremental, Upsert }

4.3 ETL数据清洗

public class DataCleaner
{
    public async Task>> CleanAsync(List> data)
    {
        var cleanedData = new List>();
        
        foreach (var record in data)
        {
            var cleanedRecord = await CleanRecordAsync(record);
            
            if (cleanedRecord != null)
            {
                cleanedData.Add(cleanedRecord);
            }
        }
        
        return cleanedData;
    }
    
    private async Task> CleanRecordAsync(Dictionary record)
    {
        foreach (var key in record.Keys.ToList())
        {
            if (record[key] == null)
            {
                record[key] = GetDefaultValue(key);
            }
            
            record[key] = await NormalizeValueAsync(key, record[key]);
            
            if (!await ValidateValueAsync(key, record[key]))
            {
                return null;
            }
        }
        
        return record;
    }
    
    private object GetDefaultValue(string key)
    {
        return key switch
        {
            "amount" => 0,
            "count" => 0,
            "date" => DateTime.MinValue,
            _ => string.Empty
        };
    }
}

五、OLAP分析

5.1 OLAP多维分析

graph TD A[多维立方体] --> B[时间维度] A --> C[产品维度] A --> D[地域维度] B --> B1[日] B --> B2[周] B --> B3[月] B --> B4[年] C --> C1[类别] C --> C2[品牌] C --> C3[价格区间] D --> D1[国家] D --> D2[省份] D --> D3[城市] A --> E[指标] E --> E1[销售额] E --> E2[订单数] E --> E3[毛利]

5.2 OLAP查询优化

public class OlapQueryOptimizer
{
    public string OptimizeQuery(string query)
    {
        query = PushDownFilters(query);
        query = RemoveUnusedDimensions(query);
        query = UseAggregateTables(query);
        
        return query;
    }
    
    private string PushDownFilters(string query)
    {
        return query.Replace("WHERE", "PREWHERE");
    }
    
    private string RemoveUnusedDimensions(string query)
    {
        return query;
    }
    
    private string UseAggregateTables(string query)
    {
        if (query.Contains("GROUP BY month"))
        {
            return query.Replace("fact_sales", "fact_sales_monthly");
        }
        
        return query;
    }
    
    public async Task> ExecuteQueryAsync(string query)
    {
        var optimizedQuery = OptimizeQuery(query);
        
        return await _olapClient.QueryAsync(optimizedQuery);
    }
}

public class AggregationResult
{
    public Dictionary Dimensions { get; set; } = new Dictionary();
    public Dictionary Measures { get; set; } = new Dictionary();
}

5.3 OLAP指标计算

public class OlapMetricCalculator
{
    public async Task CalculateMetricAsync(MetricDefinition metric)
    {
        var data = await _olapClient.QueryAsync(metric.Query);
        
        return new MetricResult
        {
            Name = metric.Name,
            Value = CalculateValue(data, metric.AggregationType),
            Dimensions = metric.Dimensions,
            TimeRange = metric.TimeRange
        };
    }
    
    private decimal CalculateValue(List data, AggregationType aggregationType)
    {
        return aggregationType switch
        {
            AggregationType.Sum => data.Sum(r => r.Measures["value"]),
            AggregationType.Avg => data.Average(r => r.Measures["value"]),
            AggregationType.Max => data.Max(r => r.Measures["value"]),
            AggregationType.Min => data.Min(r => r.Measures["value"]),
            _ => 0
        };
    }
}

public class MetricDefinition
{
    public string Name { get; set; }
    public string Query { get; set; }
    public AggregationType AggregationType { get; set; }
    public List Dimensions { get; set; } = new List();
    public TimeRange TimeRange { get; set; }
}

public enum AggregationType { Sum, Avg, Max, Min, Count }

六、数据仓库监控

6.1 ETL监控

public class EtlMonitor
{
    public async Task GetMetricsAsync(string jobName)
    {
        var jobStatus = await _etlRepository.GetJobStatusAsync(jobName);
        var executionHistory = await _etlRepository.GetExecutionHistoryAsync(jobName, 10);
        
        return new EtlMetrics
        {
            JobName = jobName,
            Status = jobStatus.Status,
            LastExecutionTime = executionHistory.FirstOrDefault()?.Duration ?? TimeSpan.Zero,
            AverageExecutionTime = TimeSpan.FromMilliseconds(executionHistory.Average(e => e.Duration.TotalMilliseconds)),
            SuccessRate = executionHistory.Count(e => e.Success) / (double)executionHistory.Count * 100,
            RowCount = jobStatus.RowCount
        };
    }
    
    public async Task MonitorAsync()
    {
        var jobs = await _etlRepository.GetAllJobsAsync();
        
        foreach (var job in jobs)
        {
            var metrics = await GetMetricsAsync(job.Name);
            
            if (metrics.Status == EtlStatus.Failed)
            {
                await _alertService.SendAlert("ETL任务失败", 
                    $"任务: {job.Name}, 错误: {job.LastError}");
            }
            
            if (metrics.LastExecutionTime > TimeSpan.FromHours(1))
            {
                await _alertService.SendAlert("ETL任务执行超时", 
                    $"任务: {job.Name}, 耗时: {metrics.LastExecutionTime}");
            }
        }
    }
}

public class EtlMetrics
{
    public string JobName { get; set; }
    public EtlStatus Status { get; set; }
    public TimeSpan LastExecutionTime { get; set; }
    public TimeSpan AverageExecutionTime { get; set; }
    public double SuccessRate { get; set; }
    public long RowCount { get; set; }
}

public enum EtlStatus { Running, Success, Failed, Pending }

6.2 数据质量监控

public class DataQualityMonitor
{
    public async Task CheckDataQualityAsync(string tableName)
    {
        var checks = new List>
        {
            CheckCompletenessAsync(tableName),
            CheckAccuracyAsync(tableName),
            CheckConsistencyAsync(tableName),
            CheckTimelinessAsync(tableName)
        };
        
        var results = await Task.WhenAll(checks);
        
        return new DataQualityReport
        {
            TableName = tableName,
            Checks = results,
            OverallScore = results.Average(r => r.Score)
        };
    }
    
    private async Task CheckCompletenessAsync(string tableName)
    {
        var totalRows = await _dataWarehouseClient.ExecuteScalarAsync($"SELECT COUNT(*) FROM {tableName}");
        var nullRows = await _dataWarehouseClient.ExecuteScalarAsync($"SELECT COUNT(*) FROM {tableName} WHERE id IS NULL");
        
        var completeness = (totalRows - nullRows) / (double)totalRows * 100;
        
        return new QualityCheckResult { Name = "完整性", Score = completeness, Passed = completeness >= 99 };
    }
    
    private async Task CheckAccuracyAsync(string tableName)
    {
        var validRows = await _dataWarehouseClient.ExecuteScalarAsync($"SELECT COUNT(*) FROM {tableName} WHERE amount >= 0");
        var totalRows = await _dataWarehouseClient.ExecuteScalarAsync($"SELECT COUNT(*) FROM {tableName}");
        
        var accuracy = validRows / (double)totalRows * 100;
        
        return new QualityCheckResult { Name = "准确性", Score = accuracy, Passed = accuracy >= 99 };
    }
}

public class DataQualityReport
{
    public string TableName { get; set; }
    public List Checks { get; set; }
    public double OverallScore { get; set; }
}

public class QualityCheckResult
{
    public string Name { get; set; }
    public double Score { get; set; }
    public bool Passed { get; set; }
}

七、数据仓库最佳实践

7.1 分层设计原则

  • ODS层:保留原始数据,支持数据追溯
  • DWD层:清洗转换,建立事实表和维度表
  • DWS层:主题域汇总,减少重复计算
  • ADS层:面向应用,直接支持报表和分析

7.2 维度建模最佳实践

public class DimensionModelingBestPractices
{
    public string DesignFactTable(string businessProcess)
    {
        return businessProcess switch
        {
            "sales" => "fact_sales",
            "orders" => "fact_orders",
            "page_views" => "fact_page_views",
            _ => throw new ArgumentException("Unknown business process")
        };
    }
    
    public List DesignDimensionTables(string businessProcess)
    {
        return businessProcess switch
        {
            "sales" => new List { "dim_time", "dim_product", "dim_customer", "dim_region" },
            "orders" => new List { "dim_time", "dim_order", "dim_customer", "dim_channel" },
            _ => new List()
        };
    }
}

7.3 ETL最佳实践

  • 使用CDC减少对源系统的影响
  • 并行处理提高ETL效率
  • 数据验证确保数据质量
  • 日志记录便于问题排查
  • 增量处理减少数据量

八、总结

数据仓库是数据密集型应用中构建数据分析平台的核心组件。通过分层架构设计、维度建模、ETL流程和OLAP分析,能够支持决策分析和业务洞察。建立完善的监控体系,确保数据质量和ETL任务的可靠性。