一、数据仓库概述
数据仓库是用于存储和分析企业历史数据的系统,能够支持决策分析和业务洞察。在数据密集型应用中,数据仓库是构建数据分析平台的核心组件。
二、数据仓库架构
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任务的可靠性。