一、数据仓库概述
数据仓库是为数据分析和决策支持而设计的数据库系统,通过整合多个数据源的数据,提供统一的视图和分析能力。OLAP(联机分析处理)是数据仓库的核心技术,支持多维数据分析。
二、数据仓库架构
2.1 数据仓库分层架构
graph TD
A[源系统] --> B[ODS层]
B --> C[DWD层]
C --> D[DWS层]
D --> E[ADS层]
B --> B1[原始数据]
B1 --> B2[数据清洗]
B2 --> C
C --> C1[明细数据]
C1 --> C2[数据整合]
C2 --> D
D --> D1[汇总数据]
D1 --> D2[指标计算]
D2 --> E
E --> E1[分析报表]
E --> E2[数据接口]
E --> E3[数据服务]
F[数据集市] --> F1[销售集市]
F1 --> F2[财务集市]
F2 --> F3[运营集市]
2.2 数据仓库层级对比
| 层级 | 描述 | 数据特点 | 主要用途 |
|---|---|---|---|
| ODS | 操作数据存储 | 原始、增量 | 数据接入 |
| DWD | 数据仓库明细层 | 清洗、整合 | 数据基础 |
| DWS | 数据仓库汇总层 | 汇总、指标 | 分析基础 |
| ADS | 应用数据服务层 | 报表、接口 | 业务应用 |
三、数据建模
3.1 星型模型
erDiagram
FACT_SALES ||--o{ DIM_DATE : "date_key"
FACT_SALES ||--o{ DIM_PRODUCT : "product_key"
FACT_SALES ||--o{ DIM_CUSTOMER : "customer_key"
FACT_SALES ||--o{ DIM_REGION : "region_key"
FACT_SALES {
int sale_id PK
int date_key FK
int product_key FK
int customer_key FK
int region_key FK
decimal amount
int quantity
datetime sale_time
}
DIM_DATE {
int date_key PK
date full_date
int year
int quarter
int month
int day
int day_of_week
string holiday_flag
}
DIM_PRODUCT {
int product_key PK
string product_name
string category
string brand
decimal price
}
DIM_CUSTOMER {
int customer_key PK
string customer_name
string email
string phone
int age
string gender
}
DIM_REGION {
int region_key PK
string region_name
string country
string province
string city
}
3.2 雪花模型
erDiagram
FACT_SALES ||--o{ DIM_DATE : "date_key"
FACT_SALES ||--o{ DIM_PRODUCT : "product_key"
FACT_SALES ||--o{ DIM_CUSTOMER : "customer_key"
FACT_SALES ||--o{ DIM_REGION : "region_key"
DIM_REGION ||--o{ DIM_COUNTRY : "country_key"
DIM_COUNTRY ||--o{ DIM_CONTINENT : "continent_key"
DIM_PRODUCT ||--o{ DIM_CATEGORY : "category_key"
DIM_CUSTOMER ||--o{ DIM_CUSTOMER_TYPE : "type_key"
FACT_SALES {
int sale_id PK
int date_key FK
int product_key FK
int customer_key FK
int region_key FK
decimal amount
}
DIM_COUNTRY {
int country_key PK
string country_name
int continent_key FK
}
DIM_CONTINENT {
int continent_key PK
string continent_name
}
DIM_CATEGORY {
int category_key PK
string category_name
string parent_category
}
3.3 星座模型
erDiagram
FACT_SALES ||--o{ DIM_DATE : "date_key"
FACT_SALES ||--o{ DIM_PRODUCT : "product_key"
FACT_SALES ||--o{ DIM_CUSTOMER : "customer_key"
FACT_INVENTORY ||--o{ DIM_DATE : "date_key"
FACT_INVENTORY ||--o{ DIM_PRODUCT : "product_key"
FACT_INVENTORY ||--o{ DIM_WAREHOUSE : "warehouse_key"
FACT_PROMOTION ||--o{ DIM_DATE : "date_key"
FACT_PROMOTION ||--o{ DIM_PRODUCT : "product_key"
FACT_PROMOTION ||--o{ DIM_CAMPAIGN : "campaign_key"
FACT_SALES {
int sale_id PK
int date_key FK
int product_key FK
int customer_key FK
decimal amount
}
FACT_INVENTORY {
int inventory_id PK
int date_key FK
int product_key FK
int warehouse_key FK
int stock_level
}
FACT_PROMOTION {
int promotion_id PK
int date_key FK
int product_key FK
int campaign_key FK
decimal discount
}
四、ETL流程
4.1 ETL架构
graph TD
A[数据抽取 Extract] --> B[数据转换 Transform]
B --> C[数据加载 Load]
A --> A1[全量抽取]
A --> A2[增量抽取]
A --> A3[CDC抽取]
B --> B1[数据清洗]
B1 --> B2[数据转换]
B2 --> B3[数据整合]
C --> C1[全量加载]
C --> C2[增量加载]
C --> C3[追加加载]
D[ETL监控] --> E[任务调度]
D --> F[数据质量检查]
D --> G[日志记录]
D --> H[告警通知]
4.2 ETL实现
public class EtlPipeline
{
private readonly List _steps = new();
public void AddStep(IEtlStep step)
{
_steps.Add(step);
}
public async Task ExecuteAsync(EtlContext context)
{
var result = new EtlResult();
foreach (var step in _steps)
{
try
{
await step.ExecuteAsync(context);
result.SuccessSteps++;
}
catch (Exception ex)
{
result.FailedSteps++;
result.Errors.Add(new EtlError
{
StepName = step.Name,
ErrorMessage = ex.Message,
Timestamp = DateTime.UtcNow
});
if (step.IsCritical)
{
throw;
}
}
}
result.Success = result.FailedSteps == 0;
return result;
}
}
public interface IEtlStep
{
string Name { get; }
bool IsCritical { get; }
Task ExecuteAsync(EtlContext context);
}
public class EtlContext
{
public Dictionary Variables { get; set; } = new();
public DateTime StartTime { get; set; }
public DateTime EndTime { get; set; }
public long ProcessedRecords { get; set; }
}
public class EtlResult
{
public bool Success { get; set; }
public int SuccessSteps { get; set; }
public int FailedSteps { get; set; }
public List Errors { get; set; } = new();
}
4.3 数据抽取
public class DataExtractor
{
public async Task> ExtractFullAsync(string source, DataSourceConfig config)
{
var records = new List();
using var connection = CreateConnection(config);
var query = config.Query;
using var command = connection.CreateCommand();
command.CommandText = query;
using var reader = await command.ExecuteReaderAsync();
while (await reader.ReadAsync())
{
records.Add(ConvertToRecord(reader));
}
return records;
}
public async Task> ExtractIncrementalAsync(string source, DataSourceConfig config, DateTime lastExtractTime)
{
var records = new List();
using var connection = CreateConnection(config);
var query = $"{config.Query} WHERE {config.IncrementalColumn} > @LastExtractTime";
using var command = connection.CreateCommand();
command.CommandText = query;
command.Parameters.AddWithValue("@LastExtractTime", lastExtractTime);
using var reader = await command.ExecuteReaderAsync();
while (await reader.ReadAsync())
{
records.Add(ConvertToRecord(reader));
}
return records;
}
public async Task> ExtractCdcAsync(string source, DataSourceConfig config, string checkpoint)
{
return await _cdcService.GetChangesAsync(source, config, checkpoint);
}
private IDbConnection CreateConnection(DataSourceConfig config)
{
return config.DataSourceType switch
{
DataSourceType.MySql => new MySqlConnection(config.ConnectionString),
DataSourceType.Postgres => new NpgsqlConnection(config.ConnectionString),
DataSourceType.SqlServer => new SqlConnection(config.ConnectionString),
_ => throw new NotSupportedException("不支持的数据源类型")
};
}
}
public class DataSourceConfig
{
public string ConnectionString { get; set; }
public DataSourceType DataSourceType { get; set; }
public string Query { get; set; }
public string IncrementalColumn { get; set; }
public CdcConfig CdcConfig { get; set; }
}
public enum DataSourceType { MySql, Postgres, SqlServer, Oracle, MongoDB, Kafka }
4.4 数据转换
public class DataTransformer
{
public List Transform(List records, List rules)
{
var transformed = new List();
foreach (var record in records)
{
var transformedRecord = TransformRecord(record, rules);
transformed.Add(transformedRecord);
}
return transformed;
}
private DataRecord TransformRecord(DataRecord record, List rules)
{
var transformed = record.Clone();
foreach (var rule in rules)
{
transformed = rule.Apply(transformed);
}
return transformed;
}
public List CleanData(List records)
{
return records
.Where(r => IsValidRecord(r))
.Select(r => RemoveNullValues(r))
.Select(r => TrimStringValues(r))
.ToList();
}
private bool IsValidRecord(DataRecord record)
{
return record.Fields.Any();
}
private DataRecord RemoveNullValues(DataRecord record)
{
var clean = new DataRecord();
foreach (var (key, value) in record.Fields)
{
if (value != null)
{
clean.Fields[key] = value;
}
}
return clean;
}
private DataRecord TrimStringValues(DataRecord record)
{
var clean = new DataRecord();
foreach (var (key, value) in record.Fields)
{
clean.Fields[key] = value is string str ? str.Trim() : value;
}
return clean;
}
}
public class TransformationRule
{
public string RuleName { get; set; }
public Func Apply { get; set; }
}
五、OLAP分析
5.1 OLAP多维分析
public class OlapAnalyzer
{
public async Task AnalyzeAsync(OlapQuery query)
{
var result = new OlapResult
{
Dimensions = query.Dimensions,
Measures = query.Measures,
Filters = query.Filters
};
var data = await _dataWarehouse.QueryAsync(query);
result.Cells = BuildCells(data, query);
return result;
}
private List BuildCells(List data, OlapQuery query)
{
var cells = new List();
var grouped = data.GroupBy(r => BuildKey(r, query.Dimensions));
foreach (var group in grouped)
{
var cell = new OlapCell();
foreach (var dim in query.Dimensions)
{
cell.Dimensions[dim] = group.Key[dim];
}
foreach (var measure in query.Measures)
{
cell.Measures[measure] = CalculateMeasure(group.ToList(), measure);
}
cells.Add(cell);
}
return cells;
}
private object CalculateMeasure(List records, string measure)
{
return measure switch
{
"SUM" => records.Sum(r => Convert.ToDecimal(r.Fields["Amount"])),
"COUNT" => records.Count,
"AVG" => records.Average(r => Convert.ToDecimal(r.Fields["Amount"])),
"MAX" => records.Max(r => Convert.ToDecimal(r.Fields["Amount"])),
"MIN" => records.Min(r => Convert.ToDecimal(r.Fields["Amount"])),
_ => records.Sum(r => Convert.ToDecimal(r.Fields[measure]))
};
}
public async Task DrillDownAsync(OlapQuery query, string dimension, string value)
{
query.Filters[dimension] = value;
return await AnalyzeAsync(query);
}
public async Task RollUpAsync(OlapQuery query, string dimension)
{
query.Dimensions.Remove(dimension);
return await AnalyzeAsync(query);
}
}
public class OlapQuery
{
public List Dimensions { get; set; } = new();
public List Measures { get; set; } = new();
public Dictionary Filters { get; set; } = new();
public string CubeName { get; set; }
}
public class OlapCell
{
public Dictionary Dimensions { get; set; } = new();
public Dictionary Measures { get; set; } = new();
}
5.2 OLAP引擎实现
public class OlapCube
{
private readonly Dictionary> _data = new();
public void LoadData(string dimension, List records)
{
_data[dimension] = records;
}
public async Task> QueryAsync(OlapQuery query)
{
var allRecords = new List();
foreach (var dim in query.Dimensions)
{
if (_data.TryGetValue(dim, out var records))
{
allRecords.AddRange(records);
}
}
var filtered = FilterRecords(allRecords, query.Filters);
return filtered;
}
private List FilterRecords(List records, Dictionary filters)
{
return records.Where(record =>
{
foreach (var (key, value) in filters)
{
if (!record.Fields.TryGetValue(key, out var recordValue) ||
!recordValue.Equals(value))
{
return false;
}
}
return true;
}).ToList();
}
public async Task BuildAggregationAsync(AggregationConfig config)
{
foreach (var measure in config.Measures)
{
var aggregated = await CalculateAggregation(measure);
_data[measure.Name] = aggregated;
}
}
private async Task> CalculateAggregation(AggregationMeasure measure)
{
var baseData = _data[measure.SourceDimension];
return baseData
.GroupBy(r => r.Fields[measure.GroupBy])
.Select(g => new DataRecord
{
Fields = new Dictionary
{
{ measure.GroupBy, g.Key },
{ measure.Name, ApplyAggregation(g.ToList(), measure.Function) }
}
})
.ToList();
}
private object ApplyAggregation(List records, string function)
{
return function switch
{
"SUM" => records.Sum(r => Convert.ToDecimal(r.Fields["Value"])),
"COUNT" => records.Count,
"AVG" => records.Average(r => Convert.ToDecimal(r.Fields["Value"])),
_ => records.Sum(r => Convert.ToDecimal(r.Fields["Value"]))
};
}
}
六、数据仓库最佳实践
6.1 数据仓库设计原则
| 原则 | 描述 | 实现方式 |
|---|---|---|
| 分层设计 | 清晰的数据层级 | ODS/DWD/DWS/ADS |
| 维度建模 | 星型/雪花模型 | 事实表+维度表 |
| 增量更新 | 减少数据处理量 | CDC、增量抽取 |
| 数据质量 | 保证数据准确性 | 数据校验、监控 |
| 可扩展性 | 支持业务增长 | 分区、索引 |
6.2 OLAP最佳实践
- 合理设计维度和度量
- 使用物化视图加速查询
- 实现数据预聚合
- 支持多维分析操作
- 优化查询性能
七、总结
数据仓库与OLAP分析是数据密集型应用的核心技术。通过分层架构设计(ODS/DWD/DWS/ADS),能够实现数据的逐步清洗和整合。维度建模(星型模型、雪花模型、星座模型)是数据仓库设计的关键。ETL流程负责数据的抽取、转换和加载。OLAP多维分析支持复杂的数据分析和决策支持。遵循数据仓库最佳实践,能够构建高效、可扩展的数据分析平台。