【数据工程】数据仓库与ETL实战:构建企业级数据平台
title: "【数据工程】数据仓库与ETL实战:构建企业级数据平台"date: 2024-05-29 09:00:00tags: ["数据仓库", "ETL", "数据工程", "Spark", "Hive"]categories: ["大数据", "数据工程"]

一、数据仓库概述
1.1 数据仓库概念
数据仓库是一个面向主题的、集成的、非易失的、随时间变化的数据集合,用于支持管理决策。
1.2 数据仓库架构
┌─────────────────────────────────────────────────────────────────┐
│ 数据仓库架构 │
├─────────────────────────────────────────────────────────────────┤
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ 源系统 │ │ 源系统 │ │ 源系统 │ │
│ │ (OLTP) │ │ (日志) │ │ (外部数据) │ │
│ └──────┬───────┘ └──────┬───────┘ └──────┬───────┘ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ ETL 层 │ │
│ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ │ │
│ │ │ Extract │→│ Transform│→│ Load │ │ │
│ │ └─────────┘ └─────────┘ └─────────┘ │ │
│ └──────────────────────────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ 数据仓库 │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │
│ │ │ ODS层 │→│ DWD层 │→│ DWS层 │→│ ADS层 │ │
│ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │
│ └──────────────────────────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ 数据应用层 │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │
│ │ │ 报表分析 │ │ 数据挖掘 │ │ 可视化 │ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ │ │
│ └──────────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘
1.3 分层架构
| ODS | 原始数据层 | 存储原始数据 | 全量、原始格式 |
| DWD | 明细数据层 | 数据清洗和转换 | 清洗后明细 |
| DWS | 汇总数据层 | 聚合计算 | 主题汇总 |
| ADS | 应用数据层 | 面向应用 | 报表数据 |
二、ETL流程设计
2.1 ETL架构
class ETLFramework:
def __init__(self):
self.stages = []
def add_stage(self, stage):
self.stages.append(stage)
def run(self, input_data):
data = input_data
for stage in self.stages:
data = stage.execute(data)
return data
class ExtractStage:
def __init__(self, source):
self.source = source
def execute(self, data):
return self._extract_from_source()
def _extract_from_source(self):
# 从源系统提取数据
pass
class TransformStage:
def __init__(self, transformations):
self.transformations = transformations
def execute(self, data):
for transform in self.transformations:
data = transform.apply(data)
return data
class LoadStage:
def __init__(self, target):
self.target = target
def execute(self, data):
self._load_to_target(data)
return data
2.2 数据抽取
# 数据库抽取
import psycopg2
import pandas as pd
class DatabaseExtractor:
def __init__(self, connection_string):
self.connection_string = connection_string
def extract(self, query):
conn = psycopg2.connect(self.connection_string)
df = pd.read_sql(query, conn)
conn.close()
return df
def incremental_extract(self, table, last_timestamp, timestamp_column='updated_at'):
query = f"""
SELECT * FROM {table}
WHERE {timestamp_column} > '{last_timestamp}'
"""
return self.extract(query)
2.3 数据转换
# 数据转换示例
class DataTransformer:
def __init__(self):
self.transformations = []
def add_transformation(self, func):
self.transformations.append(func)
def transform(self, df):
for func in self.transformations:
df = func(df)
return df
# 示例转换函数
def remove_nulls(df):
return df.dropna()
def format_dates(df, date_columns):
for col in date_columns:
df[col] = pd.to_datetime(df[col])
return df
def standardize_text(df, text_columns):
for col in text_columns:
df[col] = df[col].str.strip().str.lower()
return df
三、Spark ETL实现
3.1 Spark批处理ETL
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when, concat_ws
class SparkETL:
def __init__(self, app_name="ETL"):
self.spark = SparkSession.builder \\
.appName(app_name) \\
.enableHiveSupport() \\
.getOrCreate()
def extract(self, source_type, path):
if source_type == "csv":
return self.spark.read.csv(path, header=True, inferSchema=True)
elif source_type == "parquet":
return self.spark.read.parquet(path)
elif source_type == "jdbc":
return self.spark.read.jdbc(url=path, table="users")
def transform(self, df):
# 数据清洗
cleaned = df.dropna()
# 数据转换
transformed = cleaned \\
.withColumn("full_name", concat_ws(" ", col("first_name"), col("last_name"))) \\
.withColumn("is_active", when(col("status") == "active", True).otherwise(False))
# 数据过滤
filtered = transformed.filter(col("age") >= 18)
return filtered
def load(self, df, target_path, format="parquet", mode="overwrite"):
df.write \\
.format(format) \\
.mode(mode) \\
.save(target_path)
def run(self, source_path, target_path):
df = self.extract("csv", source_path)
transformed_df = self.transform(df)
self.load(transformed_df, target_path)
3.2 增量ETL
class IncrementalETL:
def __init__(self):
self.spark = SparkSession.builder \\
.appName("IncrementalETL") \\
.enableHiveSupport() \\
.getOrCreate()
def get_last_updated(self, target_table):
query = f"SELECT MAX(updated_at) as last_update FROM {target_table}"
result = self.spark.sql(query).collect()
return result[0][0] if result[0][0] else "1970-01-01"
def run_incremental(self, source_table, target_table):
last_update = self.get_last_updated(target_table)
# 抽取增量数据
df = self.spark.sql(f"""
SELECT * FROM {source_table}
WHERE updated_at > '{last_update}'
""")
# 转换
transformed = self._transform(df)
# 追加到目标表
transformed.write \\
.mode("append") \\
.insertInto(target_table)
四、Hive数据仓库
4.1 Hive表设计
— ODS层表
CREATE TABLE IF NOT EXISTS ods.user_logs (
user_id STRING,
action STRING,
timestamp STRING,
page STRING,
device STRING
)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY '\\t'
STORED AS TEXTFILE
LOCATION '/user/hive/warehouse/ods/user_logs';
— DWD层表
CREATE TABLE IF NOT EXISTS dwd.user_behavior (
user_id STRING,
action STRING,
event_time TIMESTAMP,
page STRING,
device STRING,
channel STRING
)
STORED AS PARQUET
LOCATION '/user/hive/warehouse/dwd/user_behavior';
— DWS层表
CREATE TABLE IF NOT EXISTS dws.user_daily_stats (
user_id STRING,
stat_date STRING,
page_views INT,
clicks INT,
purchases INT,
total_amount DECIMAL(10,2)
)
PARTITIONED BY (dt STRING)
STORED AS PARQUET
LOCATION '/user/hive/warehouse/dws/user_daily_stats';
4.2 Hive分区优化
— 创建分区表
CREATE TABLE dws.sales_by_region (
region STRING,
product_category STRING,
sales_amount DECIMAL(12,2),
order_count INT
)
PARTITIONED BY (year STRING, month STRING, day STRING)
STORED AS PARQUET;
— 添加分区
ALTER TABLE dws.sales_by_region ADD PARTITION (year='2024', month='05', day='29');
— 查询特定分区
SELECT * FROM dws.sales_by_region
WHERE year='2024' AND month='05';
五、数据质量保障
5.1 数据质量检查
class DataQualityChecker:
def __init__(self):
self.checks = []
def add_check(self, check_name, check_func):
self.checks.append({"name": check_name, "func": check_func})
def run_checks(self, df):
results = []
for check in self.checks:
try:
result = check["func"](df)
status = "PASS" if result else "FAIL"
except Exception as e:
status = "ERROR"
result = str(e)
results.append({
"check": check["name"],
"status": status,
"result": result
})
return results
# 示例检查函数
def check_no_nulls(df, column):
return df.filter(col(column).isNull()).count() == 0
def check_value_range(df, column, min_val, max_val):
return df.filter((col(column) < min_val) | (col(column) > max_val)).count() == 0
def check_unique(df, column):
return df.count() == df.select(column).distinct().count()
5.2 数据血缘追踪
class DataLineageTracker:
def __init__(self):
self.lineage = {}
def track_transformation(self, source_columns, target_columns, transformation):
for src_col in source_columns:
if src_col not in self.lineage:
self.lineage[src_col] = []
self.lineage[src_col].append({
"target_columns": target_columns,
"transformation": transformation
})
def get_lineage(self, column):
return self.lineage.get(column, [])
六、性能优化
6.1 Spark优化配置
# Spark配置优化
from pyspark.sql import SparkSession
spark = SparkSession.builder \\
.appName("OptimizedETL") \\
.config("spark.executor.memory", "8g") \\
.config("spark.driver.memory", "4g") \\
.config("spark.sql.shuffle.partitions", "200") \\
.config("spark.sql.autoBroadcastJoinThreshold", "104857600") \\
.config("spark.sql.parquet.enableVectorizedReader", "true") \\
.getOrCreate()
6.2 数据倾斜处理
# 数据倾斜处理
class SkewHandler:
def __init__(self):
pass
def salted_join(self, df1, df2, join_column):
# 添加随机盐
df1_salted = df1.withColumn("salt", (col("id") % 100).cast("string"))
df2_salted = df2.withColumn("salt", (col("id") % 100).cast("string"))
# 使用加盐后的列进行join
joined = df1_salted.join(df2_salted, ["id", "salt"])
return joined.drop("salt")
七、实战案例:电商数据仓库
7.1 数据仓库设计
class ECommerceDataWarehouse:
def __init__(self):
self.spark = SparkSession.builder \\
.appName("ECommerceDW") \\
.enableHiveSupport() \\
.getOrCreate()
def create_tables(self):
# 创建ODS层
self._create_ods_tables()
# 创建DWD层
self._create_dwd_tables()
# 创建DWS层
self._create_dws_tables()
# 创建ADS层
self._create_ads_tables()
def etl_pipeline(self):
# ODS -> DWD
self._ods_to_dwd()
# DWD -> DWS
self._dwd_to_dws()
# DWS -> ADS
self._dws_to_ads()
7.2 每日数据同步
# ETL调度脚本
#!/bin/bash
# 设置日期
DATE=$(date -d "yesterday" +%Y-%m-%d)
# 运行ETL作业
spark-submit \\
–master yarn \\
–deploy-mode cluster \\
–name "DailyETL" \\
–executor-memory 8g \\
–num-executors 10 \\
etl_pipeline.py \\
–date $DATE
# 检查执行状态
if [ $? -eq 0 ]; then
echo "ETL成功完成: $DATE"
else
echo "ETL失败: $DATE"
exit 1
fi
八、总结与最佳实践
8.1 关键要点
8.2 常见误区
8.3 未来趋势
- 湖仓一体:数据湖与数据仓库融合
- 实时数据仓库:支持实时分析
- 元数据管理:完善的元数据治理体系
- AI辅助ETL:自动化数据清洗和转换
参考资料:
- Apache Hive官方文档
- Apache Spark官方文档
- Kimball数据仓库设计指南
- Data Vault建模方法



