← 返回作品集

🏗️ 亿级行为大盘:1 亿行数据的湖仓之旅

原始 CSV 3.4GB → HDFS → MapReduce 清洗 → Hive 分层建仓 → Spark 分析 → Airflow 每日调度 → 12 项质量门禁。 本页全部数字实时抽取自运行中的湖仓(HiveServer2 直连 ADS 层),不是截图。 源码 | 数据:淘宝用户行为(2017-11-25 ~ 12-03)

1.0 亿
清洗后行为记录(DWD 层实测 100,095,231 行)
98.8 万
独立用户(987,994)
9 天
时间窗口(北京时区口径)
12 项
数据质量门禁,HARD 失败即熔断

分层架构 · ODS → DWD → ADS

ODS 原始层100,150,807 行 CSV 原样入湖
发现 318 条非法时间戳
DWD 明细层MR 清洗 + ROW_NUMBER 去重
Parquet 列存,损耗率 0.055%
ADS 应用层漏斗 / 类目 Top / 日大盘
本页数据的直供层

转化漏斗 · 89,660,688 次浏览之后

加购→购买转化 68.3%(用户口径),浏览→加购 75.1%——Hive 与 Spark 双引擎结果一致。

类目 Top 10 · 按浏览量

日大盘 · 周末效应肉眼可见

DAU(日活)购买事件

12-02(周六)峰值:DAU 970,401 / 浏览 1232 万 / 购买 25.8 万;周末流量较工作日 +30%。

三引擎同任务基准

单机伪分布式(YARN 8GB/8vcore),3.4GB CSV → 日级聚合。

引擎任务耗时
MapReduce日级行为计数156s
MapReduce用户级聚合(98.8 万用户)106s
Hive (MR)DWD 构建 CSV→Parquet228s
Hive (MR)ADS 分析(Parquet 扫描)27-69s
Spark漏斗+RFM+日大盘全套54s

单机结论:Spark 内存计算对迭代分析有数量级优势;列存让重复分析远离原始 CSV;MR 的价值在模型本身的可解释性。

质量门禁 · 12 项检查五域

检查域内容状态
完整性ODS/DWD 行数与用户数锚定PASS
合法性日期窗口 / behavior 枚举 / ts 正值PASS
唯一性DWD 复合主键零重复PASS
一致性清洗损耗率 < 5%PASS
业务合理性漏斗 buy 锚点、单调性PASS

门禁不是摆设:上线首日即抓到 318 条非法时间戳 + 49 条复合键重复——修数据不改规则,去重已内建进建仓 SQL。

RFM 用户分层(Spark ntile-5)

分层用户数人均行为人均购买
高价值活跃160,726182.12.85
有购买459,361106.82.97
流失预警153,80933.21.27
沉默浏览214,09577.80.00

全链路(HDFS/MR/Hive/Spark/Airflow/Superset)与复现文档见 仓库;调度 DAG cartlake_daily 每日 03:00 自动重跑并过门禁。