欢迎光临
我们一直在努力

【数据工程】数据仓库与ETL实战:构建企业级数据平台

【数据工程】数据仓库与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 关键要点

  • 分层设计:采用ODS->DWD->DWS->ADS分层架构
  • 增量处理:优先使用增量ETL减少数据处理量
  • 数据质量:建立完善的数据质量检查机制
  • 性能优化:合理配置资源,处理数据倾斜
  • 8.2 常见误区

  • 忽视数据血缘:缺乏追踪会导致问题难以定位
  • 全量更新:频繁全量更新影响性能
  • 缺少监控:ETL失败不能及时发现
  • 过度分区:过多分区会影响查询性能
  • 8.3 未来趋势

    • 湖仓一体:数据湖与数据仓库融合
    • 实时数据仓库:支持实时分析
    • 元数据管理:完善的元数据治理体系
    • AI辅助ETL:自动化数据清洗和转换

    参考资料:

    • Apache Hive官方文档
    • Apache Spark官方文档
    • Kimball数据仓库设计指南
    • Data Vault建模方法
    赞(0)
    未经允许不得转载:171主机测评 » 【数据工程】数据仓库与ETL实战:构建企业级数据平台
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址