Spark 离线报表 ETL 综合案例
概述
本文用一个电商离线报表项目把前两阶段知识串起来:Spark Core 的 RDD/算子、Spark SQL 的 DataFrame/AQE、外部数据源读写、Hive 数仓分层。目标是构建一条 ODS → DWD → DWS → ADS 的离线 ETL 流水线。
一、项目背景
1.1 需求
| 需求 | 说明 |
|---|---|
| 数据源 | MySQL 业务库(订单/用户/商品)、日志文件 |
| 目标 | 每日产出 GMV、订单量、用户活跃等报表 |
| 存储 | Hive 数仓分层表 |
| 输出 | 报表结果写回 MySQL 供 BI 查询 |
1.2 技术栈
采集:Sqoop(MySQL → HDFS)/ Flume(日志 → HDFS)
计算:Spark SQL(Hive on 数据湖)/ Spark Core(复杂逻辑)
存储:Hive(ODS/DWD/DWS/ADS)+ MySQL(报表库)
调度:Oozie/Airflow 每日定时二、数仓分层设计
| 层 | 表 | 说明 |
|---|---|---|
| ODS | ods_orders | 原始订单(全量/增量) |
| DWD | dwd_order_detail | 清洗明细(去脏、标准化) |
| DWS | dws_order_daily | 订单日汇总 |
| ADS | ads_gmv_daily | GMV 报表(写 MySQL) |
2.1 分层职责
ODS:原样接入,保留原始
DWD:清洗、脱敏、维度补齐
DWS:按日/主题轻度汇总
ADS:面向应用的最终报表三、数据接入(ODS 层)
3.1 MySQL 增量同步(Sqoop)
bash
sqoop import \
--connect jdbc:mysql://db/orderdb \
--table orders \
--incremental append \
--check-column id \
--last-value 100000 \
--target-dir /data/ods/orders \
--split-by id \
-m 83.2 日志接入(Flume → HDFS)
Flume Agent:
source: taildir(订单日志文件)
channel: memory
sink: hdfs(按天目录 /data/ods/order_log/dt=yyyyMMdd)3.3 建 ODS 表
sql
CREATE EXTERNAL TABLE IF NOT EXISTS ods.ods_orders (
id BIGINT, user_id BIGINT, product_id BIGINT,
amount DOUBLE, status INT, create_time STRING
) PARTITIONED BY (dt STRING)
STORED AS PARQUET LOCATION '/data/ods/orders';四、清洗与标准化(DWD 层)
4.1 加工逻辑
| 清洗项 | 规则 |
|---|---|
| 去重 | 按订单 id 去重,保留最新 |
| 去脏 | status 合法范围、金额 > 0 |
| 标准化 | 时间格式化、空值补默认 |
| 维度补齐 | join 用户/商品维度 |
4.2 Spark SQL 实现
scala
val spark = SparkSession.builder()
.enableHiveSupport()
.config("spark.sql.adaptive.enabled", "true")
.getOrCreate()
spark.sql("""SET spark.sql.sources.partitionOverwriteMode=dynamic""")
val sql =
"""
INSERT OVERWRITE TABLE dwd.dwd_order_detail PARTITION(dt)
SELECT
o.id, o.user_id, o.product_id, o.amount,
CASE WHEN o.status IN (0,1,2) THEN o.status ELSE 99 END AS status,
DATE_FORMAT(o.create_time, 'yyyyMMdd') AS dt,
u.city_id AS user_city
FROM ods.ods_orders o
LEFT JOIN dim.dim_user u ON o.user_id = u.user_id
WHERE o.dt = '${dt}' AND o.amount > 0
"""
spark.sql(sql)4.3 用 DataFrame 实现去重(Core 场景)
scala
val orders = spark.table("ods.ods_orders").filter($"dt" === dt)
val dedup = orders
.groupBy("id")
.agg(max("create_time").as("create_time")) // 保留最新
.join(orders, Seq("id", "create_time"), "inner")五、轻度汇总(DWS 层)
5.1 需求
按天、商品、用户统计订单指标,供上层报表复用。
5.2 实现
scala
val dws = spark.sql("""
INSERT OVERWRITE TABLE dws.dws_order_daily PARTITION(dt)
SELECT
product_id,
COUNT(DISTINCT user_id) AS user_cnt,
COUNT(*) AS order_cnt,
SUM(amount) AS gmv
FROM dwd.dwd_order_detail
WHERE dt = '${dt}'
GROUP BY product_id
""")| 参数 | 说明 |
|---|---|
| AQE | 自动合并分区、治倾斜 |
| 预聚合 | COUNT/SUM 天然聚合 |
| 分区覆盖 | dynamic 只写当日分区 |
六、报表输出(ADS 层)
6.1 计算 GMV 报表
scala
val ads = spark.sql("""
SELECT
dt, SUM(gmv) AS total_gmv,
SUM(order_cnt) AS total_orders,
COUNT(DISTINCT product_id) AS product_cnt
FROM dws.dws_order_daily
WHERE dt BETWEEN '${start}' AND '${end}'
GROUP BY dt
""")6.2 写回 MySQL
scala
val props = new java.util.Properties()
props.setProperty("user", "report")
props.setProperty("password", "***")
ads.write.mode("overwrite")
.option("batchsize", "1000")
.jdbc("jdbc:mysql://report-db:3306/report", "ads_gmv_daily", props)6.3 幂等保障
| 手段 | 说明 |
|---|---|
| 覆盖写 | 按 dt 覆盖,重复跑结果一致 |
| 任务重跑 | 先删分区再写 |
| 依赖检查 | 上游分区就绪再启动 |
七、调度与运维
7.1 每日流水线
00:10 Sqoop 增量同步 ODS
00:20 Flume 日志落 ODS
00:30 DWD 清洗任务
00:45 DWS 汇总任务
01:00 ADS 报表 + 写 MySQL
01:10 数据质量校验(对账 GMV)7.2 关键监控
| 指标 | 监控点 |
|---|---|
| 数据量 | ODS 增量行数波动 |
| 时间 | 各任务耗时 |
| 结果 | GMV 与昨日对比波动 |
| 质量 | 空值率、重复率 |
八、优化要点回顾
1. 分区裁剪:按 dt 过滤,只读当日
2. 预聚合:DWS 先行汇总,ADS 轻量
3. 广播 join:维度表广播
4. AQE:倾斜与分区自动处理
5. 动态分区:覆盖只写当日
6. 写 MySQL:批量 + 控制并行度
7. 复用缓存:多报表共用 DWS 缓存常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 重复数据 | 增量 checkpoint 设置不当,加去重 |
| 报表与库对不上 | 时区/分区边界不一致,统一口径 |
| 写 MySQL 慢 | 批量调大、减少并行 |
| 任务积压 | 上游延迟,检查 Sqoop/日志量 |
| 分区覆盖到历史 | partitionOverwriteMode 设错,用 dynamic |