ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Hive电商数仓实战:从Raw数据清洗到留存与RFM分析

Hive电商数仓实战:从Raw数据清洗到留存与RFM分析 做数据分析这些年我工作目录里躺得最多的文件就是各种“Raw”。Raw意味着未经处理、最原始的一手数据——字段乱到能逼疯强迫症值域里塞满NULL和脏字符串时间戳统一不了时区但恰恰是这些“脏东西”决定了一个分析项目能做多深。这篇过程记录就当一次完整复盘从一套电商用户行为日志的Raw数据出发走完Hive数仓分层设计、核心DDL编写、留存/复购/RFM/漏斗分析再到最后的小文件治理和结果可视化。适合正在跑电商数据分析项目、或者刚接触Hive想做一套能落地的分析体系的朋友照着这个思路走能少踩很多我踩过的坑。1. 数据进场Raw层的第一手资料长什么样1.1 一份Raw日志里到底有什么电商数据分析最常遇到的Raw数据来自三块用户行为日志、订单流水、商品与用户维表。用户行为日志通常由前端埋点上报记录的是“谁在什么时间点了什么按钮”字段类似event_id、user_id、session_id、page_url、event_time订单流水则是交易系统的导出包含order_id、user_id、order_amount、pay_time、order_status维表相对干净但也会出现两条记录重复、修改时间字段缺失这类问题。这个项目的起点就是这三类文件丢在HDFS上一堆txt和csv里没有分区、没有索引、没有统一的字段注释。第一次打开订单表我差点没绷住订单金额既有decimal格式又有字符串“39.90”混进来支付时间有的精确到毫秒有的只有日期用户ID更离谱同一个用户在不同批次日志里出现过A001和a001两种写法。这就是Raw的真实形态——它忠实地记录业务系统的混乱。1.2 字段体检正式跑数前的必要准备拿到原始文件别急着写业务SQL。我在这个项目里专门拨出一个小时做“字段体检”这个时间花得非常值因为后期所有报表都建立在这批表上。体检主要做四件事看字段完整性、看枚举值分布、看时间格式、看重复比例。用Hive建裸表把数据load进去之后先跑一组探索性的SQLSELECT COUNT(*)、COUNT(DISTINCT user_id)、MAX(event_time)、MIN(event_time)再看关键字段的NULL比例。比如订单表里order_amount的NULL比例超过5%那后续计算客单价就必须想清楚是剔除还是补0不同处理方式得出的报表数字会差不少。补充一句字段体检的SQL最好固化成一个脚本每个批次的数据进来都跑一遍防止上游业务系统改了字段格式而我们还在用老口径硬算。1.3 脏数据清单与统一清洗口径这次项目遇到的典型脏数据有六类字段类型混用、时间格式不一、枚举值大小写混乱、NULL值无意义堆积、订单重复记录、同一个用户存在多个ID映射。我整理了一张清单清洗口径直接写进ETL逻辑脏数据类型表现清洗口径金额格式不一数字和字符串混存统一转成DECIMAL(10,2)无法转换的置NULL并单独记录时间格式不一yyyy-MM-dd HH:mm:ss 与 yyyyMMdd 共存统一转成timestamp解析失败则丢弃该分区数据并告警用户ID大小写不一A001与a001统一转小写并在用户维表做ID归一映射无意义NULL渠道字段大量NULL按枚举值映射为unknown不参与渠道漏斗计算重复订单同一order_id出现多条按order_id pay_time去重保留最早一条枚举值命名混乱“已支付”和“PAID”混存通过维表映射成统一状态码统一清洗口径这件事表面上只是把数据“弄干净”本质上是立规矩。团队里每个人对“有效订单”的理解不同有的说支付成功就算有的说退款要剔除如果没有清洗层统一规则同一个月的GMV能算出三个版本。Raw层保持原样不动清洗逻辑全部放在DWD层处理是我这次做得最对的决定之一后面所有分析SQL都基于DWD谁的口径也不容易跑偏。2. 建表设计从Raw到DWDDDL里藏着的性能决策2.1 分层思路Raw、DWD、ADS各管什么事数仓分层听起来玄乎落到这个项目里其实就三句话Raw层存原样数据一个字不改出了纠纷能回溯DWD层做清洗、脱敏、标准化统一口径承载90%的复杂查询ADS层面向具体报表提前聚合好指标跑起来秒出。偶尔还有一个中间层做宽表拼接但小项目不强制。我把绝大部分精力放在DWD层。因为分析师写的大部分SQL——留存、复购、RFM、漏斗——都是直接查DWD如果DWD表结构设计不合理后面每次查询都要做大量的join和group by性能会很糟糕。举一个反面例子最初我把订单明细和用户行为日志放在同一张表里想着“反正都是用户相关的”结果订单维度的分析要过滤掉几十倍的行为日志数据白白浪费大量计算资源。后来拆成订单流和事件流两张表各自维护各自的分区逻辑清楚跑得也快。2.2 核心DDL实例订单事实表和维表订单事实表是这个项目最重要的表之一。设计上选择了按天分区分区字段dt用字符串类型格式yyyy-MM-dd。这样做的好处是既方便通过dt做时间裁剪又能配合后续的合并和回溯。CREATE EXTERNAL TABLE dwd_order_detail_di ( order_id STRING COMMENT 订单ID, user_id STRING COMMENT 用户ID, product_id STRING COMMENT 商品ID, category_id STRING COMMENT 类目ID, order_amount DECIMAL(10,2) COMMENT 订单金额, pay_time TIMESTAMP COMMENT 支付时间, channel STRING COMMENT 渠道来源, is_refund TINYINT COMMENT 是否退款 0否 1是 ) PARTITIONED BY (dt STRING) STORED AS ORC TBLPROPERTIES ( orc.compress SNAPPY );表名里我特意留下_di的后缀这个di代表每日增量提醒自己和后人这张表是按天刷新的增量表不是全量快照。类似的后缀规范很重要等表多起来之后光看表名就能大致猜出刷新策略排查问题时少走很多弯路。用户维表两端都有一张DWD层一般要有用户的基础属性、注册时间、渠道分组这些字段。订单事实表不存用户明细属性只存user_id需要用户画像时再join维表。如果一上来就把姓名、手机号、城市全塞进事实表字段冗余会让表变得臃肿而且用户属性变化时更新成本极高这也是老数仓前辈常说的“维度退化”要谨慎处理的原因。2.3 存储格式和分区策略为什么这么选很多新人问我为什么不用TextFile、为什么不用默认的压缩非要ORC加SNAPPY。三层原因第一TextFile这类行式存储在跑分析型SQL时需要扫描大量无关列而ORC是列式存储只读取涉及到的列在订单表动辄几百万行的场景下扫描量能小一个数量级。第二SNAPPY压缩比适中、解压速度快虽然压缩率不如ZSTD但对日常分析查询来说吞吐量才是关键指标。基于列式存储轻量压缩很多聚合查询甚至能做到原地读取不落盘。分区策略上订单表按dt分区是刚需因为几乎所有业务问题都会带时间范围没有分区就只能全表扫描。需要补充的是分桶我并没有在订单表上强制使用。原因很简单这个项目里join的查询大多已经通过分区裁剪把数据量降到可接受范围贸然分桶反而可能让数据分布不均。如果你的业务里有一个超大维表要频繁跟事实表join再考虑在join key上做分桶优化否则分桶属于锦上添花不做不致命。提示外部表还是内部表我的习惯是Raw层全用外部表数据文件由上游系统产出Hive只做元数据管理。DWD层用外部表还是内部表取决于团队运维方式内部表在DROP TABLE时会连带删除数据外部表则不会建议把DWD层设为外部表避免误删。3. 用户留存分析窗口函数和首日标记的实战3.1 留存的口径选择留存是电商分析里最高频的指标但“留存”这个词在不同人口中含义不同。我这边的口径是“某个时间窗口内新增访问用户在第N天仍然产生访问行为的比例。”这个定义关键在于三个词窗口期、新增、活跃。窗口期我一般取自然日周留存和月留存也可以做但计算逻辑完全一样只是粒度不同。口径一定先定好再写SQL因为不同口径结果差异很大。比如把“新增用户”定义为首次下单用户和定义为首次访问用户得到的留存率可能差好几倍。下单用户本来就经过了一层意向筛选留存天然更高。所以这份项目里的留存全部分为“新增访问留存”和“新增下单留存”两个口径分开出数避免管理层把两个数混在一起比较。3.2 给用户标“第几天”ROW_NUMBER与MIN聚合求留存的第一步是算出每个用户在一段周期内的首次活跃日期。最直接的做法是用GROUP BY加MIN聚合WITH first_active AS ( SELECT user_id, MIN(dt) AS first_dt FROM dwd_user_active_di WHERE dt 2024-01-01 AND dt 2024-01-31 GROUP BY user_id ) SELECT user_id, first_dt FROM first_active;这段SQL在窗口函数之外先拿到“每个用户第一天”思路很朴素但性能通常也不错。用不用窗口函数看场景——假设你要在明细表上直接给每一行标一个序号然后取序号为1的记录那就要用ROW_NUMBERSELECT user_id, dt AS first_dt FROM ( SELECT user_id, dt, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY dt ASC) AS rn FROM dwd_user_active_di WHERE dt 2024-01-01 AND dt 2024-01-31 ) t WHERE rn 1;两种写法结果等价。MIN聚合方式一个用户只保留一行首日信息适合后续跟活跃表做joinROW_NUMBER则是在明细数据上做标号适合你想同时保留首日这一行里的其他字段比如首日来源渠道的场景。我的经验是只要有“给每一行编号取第一”的需求窗口函数都是最趁手的工具它能让SQL语义非常直白。3.3 留存率SQL与执行注意点拿到每个用户的首日之后留存计算就是一次简单的join和分组WITH first_active AS ( SELECT user_id, MIN(dt) AS first_dt FROM dwd_user_active_di WHERE dt 2024-01-01 AND dt 2024-01-31 GROUP BY user_id ) SELECT DATEDIFF(b.dt, a.first_dt) AS day_diff, COUNT(DISTINCT a.user_id) AS retained_users FROM first_active a JOIN dwd_user_active_di b ON a.user_id b.user_id AND b.dt a.first_dt AND b.dt DATE_ADD(a.first_dt, 29) GROUP BY DATEDIFF(b.dt, a.first_dt) ORDER BY day_diff;day_diff等于1的含义是“用户首日之后第1天回来”也就是次日留存。要得到留存率再除以窗口期内的新增用户数就行。我把这个除数单独放在报表里算不在SQL里除这样看数的人能同时看到“基数”和“留存数”对可信度更有利。执行这条SQL时有几个注意点JOIN条件里必须带上b.dt的上下界裁剪否则即使有分区Hive也会为了安全扫描该用户所有分区数据COUNT(DISTINCT)在数据量大时会启动较高reduce但留存表一般按天分区后数据可控不需要刻意优化DATEDIFF用的是Hive内置函数返回的是整数天数差注意跟日期格式保持一致别混用yyyy-MM-dd和yyyyMMdd。提示做留存分析时first_active的窗口期要足够长至少覆盖“首日用户数”和“最长观察天数”。比如要看30日留存新增用户窗口就至少留出30天不能拿1月31日的用户去看2月10日的留存时间还没到。这个坑我踩过当时导出的报表里最后几天留存率跳水后来才发现是窗口截断不是业务问题。4. 复购率、RFM分群和漏斗三个必做的用户分析4.1 复购率口径与去重SQL留存回答的是“用户还来不来”复购回答的是“用户回不回来买”。这个项目里复购率我取的口径是统计周期内“下单天数大于等于2天”的用户占同期有下单行为用户的比例。注意是“下单天数”而不是“下单次数”如果一天内同一个用户买了5单那只能算一天的购买行为。这样处理是为了把一次性冲动消费和真实复购区分开否则促销节点会把复购率冲得很高。SQL写成这样WITH user_buy_day AS ( SELECT user_id, dt FROM dwd_order_detail_di WHERE dt BETWEEN 2024-01-01 AND 2024-03-31 AND is_refund 0 GROUP BY user_id, dt ) SELECT COUNT(IF(cnt 1, 1, NULL)) / COUNT(*) AS repurchase_rate FROM ( SELECT user_id, COUNT(*) AS cnt FROM user_buy_day GROUP BY user_id ) t;这个查询里最关键的是最内层GROUP BY user_id, dt它把同一个用户一天内的多笔订单折叠成一条完成“按天去重”。如果不做这一步双11当天买了3单的用户在子查询里cnt就是3直接被误判为复购人群。我的习惯是先把口径写在一行注释里SQL再简洁也要让人一眼看懂“复购率”的准确含义。4.2 RFM评分与NTILE分箱RFM是用户价值分群最经典的方法三个维度分别是最近一次消费时间、消费频率、消费金额。先算出每个用户的R、F、M原始值SELECT user_id, DATEDIFF(CURRENT_DATE, MAX(TO_DATE(pay_time))) AS recency, COUNT(DISTINCT order_id) AS frequency, SUM(order_amount) AS monetary FROM dwd_order_detail_di WHERE dt BETWEEN 2024-01-01 AND 2024-06-30 AND is_refund 0 GROUP BY user_id;拿到原始值之后最常用的做法是分位数打分。Hive里的NTILE函数可以把数据平均分成N桶非常适合给RFM三个维度各打1到5分SELECT user_id, NTILE(5) OVER (ORDER BY recency ASC) AS r_score, NTILE(5) OVER (ORDER BY frequency DESC) AS f_score, NTILE(5) OVER (ORDER BY monetary DESC) AS m_score FROM ( SELECT user_id, DATEDIFF(CURRENT_DATE, MAX(TO_DATE(pay_time))) AS recency, COUNT(DISTINCT order_id) AS frequency, SUM(order_amount) AS monetary FROM dwd_order_detail_di WHERE dt BETWEEN 2024-01-01 AND 2024-06-30 AND is_refund 0 GROUP BY user_id ) t;这里有一个反直觉的细节R值越小代表最近刚消费所以第一行NTILE的排序是ASCR值越小分越高F和M越大越好排序用DESC。打分之后再把三个分数拼接成“111-555”这种用户分箱标识进一步归成高价值、一般价值、流失风险几类。实测下来RFM分群比单看ARPU更靠谱因为ARPU会把“偶尔花大钱的用户”和“经常花小钱的用户”混在一类RFM能把这两类人拆开运营。4.3 转化漏斗行为事件的拼接与计算漏斗分析的原始数据是用户行为事件流。我这次问的是最经典的电商闭环商品浏览、加购、下单、支付。因为四种事件都在同一张DWD行为日志表里所以只需要一次扫描SELECT COUNT(DISTINCT IF(event_name view_item, user_id, NULL)) AS view_uv, COUNT(DISTINCT IF(event_name add_cart, user_id, NULL)) AS cart_uv, COUNT(DISTINCT IF(event_name order, user_id, NULL)) AS order_uv, COUNT(DISTINCT IF(event_name pay, user_id, NULL)) AS pay_uv FROM dwd_user_action_log_di WHERE dt 2024-03-31 AND event_name IN (view_item, add_cart, order, pay);漏斗分析最容易被忽略的是“串联顺序”——严格漏斗要求用户必须按浏览→加购→下单→支付依次发生倒着加购的用户不算有效转化。上面的SQL只统计了各阶段独立UV严格说叫“各环节用户数”不叫“漏斗”。要做严格漏斗可以用窗口函数判断事件序列SELECT COUNT(DISTINCT IF(has_view 1 AND has_cart 1 AND has_order 1 AND has_pay 1, user_id, NULL)) AS full_funnel_uv FROM ( SELECT user_id, MAX(IF(event_name view_item, 1, 0)) AS has_view, MAX(IF(event_name add_cart, 1, 0)) AS has_cart, MAX(IF(event_name order, 1, 0)) AS has_order, MAX(IF(event_name pay, 1, 0)) AS has_pay FROM dwd_user_action_log_di WHERE dt 2024-03-31 GROUP BY user_id ) t;这版更严格要求一个用户在每个环节都出现过。实际业务里两种口径我会并行跑宽口径看大盘趋势严口径用来定位流失环节。如果宽口径的浏览到加购转化率是40%严口径却是20%那说明大量用户跳过了浏览直接加购可能就是搜索直达加购页面的场景。这种细节比单一漏斗数字本身更能指导运营。5. 过程复盘小文件治理和“某天数据消失”排查5.1 小文件是怎么产生的危害有多大跑到第三周的时候集群响应明显变慢跑一条Hive任务经常卡在MapReduce的“Initializing Job”。排查下来发现是DWD层订单表某个分区下面堆了几百个不足几KB的小文件。小文件的来源很典型动态分区写入时如果Map或Reduce的数量比较多每个任务都会往各个分区里落文件一天下来一个分区可能被切碎成上百个小文件另外每天定时任务重跑某天数据时没有先做合并每次重跑都会往分区里追加一批小文件。小文件的危害主要体现在两个层面一是给NameNode的内存带来压力每个文件都要维护元数据二是查询时Map数量跟InputSplit挂钩一坨小文件会让Map任务暴增任务调度开销远大于实际计算开销表现出来就是“卡在启动阶段”。这个过程让我真切体会到“小文件是Hive优化的头号隐形杀手”这句话的分量。5.2 一套可行的合并与参数调优治理小文件对存量数据做合并是第一步控制增量文件是长期办法。存量合并最稳妥的方式是INSERT OVERWRITE重写目标分区SET hive.exec.dynamic.partition.modenonstrict; SET hive.merge.mapfilestrue; SET hive.merge.mapredfilestrue; SET hive.merge.size.per.task256000000; SET hive.merge.smallfiles.avgsize16777216; INSERT OVERWRITE TABLE dwd_order_detail_di PARTITION (dt) SELECT order_id, user_id, product_id, category_id, order_amount, pay_time, channel, is_refund, dt FROM dwd_order_detail_di WHERE dt BETWEEN 2024-03-01 AND 2024-03-31 DISTRIBUTE BY dt;这里的关键是DISTRIBUTE BY dt它让同一个分区的数据尽量被同一个Reduce处理配合OVERWRITE写回一个分区最终产出少量大文件。merge.size.per.task设成256MB的意思是每个Reduce尽量收集到256MB再落盘smallfiles.avgsize是触发合并的阈值低于16MB的平均小文件就合并。长期控制增量要盯住两点一是动态分区写入时按业务需求评估Reduce数别让Reduce数远大于分区数二是重跑历史的调度脚本先执行一遍清空目标分区的操作再写入新数据避免旧小文件叠加在分区里。稳定运行之后我把这套参数沉淀成项目里的“建表写入规范”只要是日新增的DWD表都默认带上新来的同事照着抄也能减少出问题。5.3 一次真实的异常排查经历项目快收尾时运营反馈某天的支付用户数比前一天跌了70%。第一反应是业务出问题了但对比同一天的新增用户数发现访问量正常说明不是流量骤降。于是直接查DWD订单明细发现当天分区居然一条数据都没有。排查链路从下往上走先看ODS层原始文件是否入库发现文件在再看清洗任务日志发现当天清洗SQL跑挂了。挂的原因很离谱业务方在订单备注字段里塞了一段超长文本触发了某个字段解析异常整个任务失败退出。因为调度脚本没有开启“失败自动重试并告警”任务就静默失败了。那次之后我把清洗任务的健壮性做了一次升级一是大字段统一改成STRING并加长度上限超长文本截断并记录异常表二是调度里配置失败重试与钉钉告警任务挂掉第一时间通知三是每个DWD表后挂一个数据质量校验SQL检查当天分区行数与上一日行数偏差超过50%就报警。被这么治过一次之后后面再没出现过“某天数据凭空消失”的事数据质量巡检SQL也被我保留下来成为每个DWD表建表的标配。6. 结果出口从Hive表到可读的图表6.1 导出到本地目录的几种方式分析SQL写完最终要出报表。相比直接在Hive控制台里看一堆数字我更习惯把结果导出到本地再交给Python或者Excel做进一步加工。最常用的导出方式有两种。第一种是直接在命令行执行查询并重定向到文件hive -e set hive.cli.print.headertrue; SELECT dt, COUNT(DISTINCT user_id) AS uv FROM dwd_user_active_di WHERE dt BETWEEN 2024-06-01 AND 2024-06-30 GROUP BY dt; /data/export/uv_daily.tsv第二种是用INSERT OVERWRITE LOCAL DIRECTORY可以在Hive内部指定分隔符和导出路径INSERT OVERWRITE LOCAL DIRECTORY /data/export/rfm_result ROW FORMAT DELIMITED FIELDS TERMINATED BY , SELECT user_id, r_score, f_score, m_score FROM ads_user_rfm_di WHERE dt 2024-06-30;第一种适合临时查数第二种适合定时脚本自动跑并落地固定目录。导出之后Python一侧用pandas读文件指定encoding和sep基本就能无缝衔接数据分析阶段。这里顺手说一句如果结果集超过几十万行别硬往Excel里塞导成CSV后用Python做聚合再输出报表页比Excel透视频繁卡死舒服很多。6.2 可视化阶段的乱码与中文问题Python做数据分析与可视化在这个项目里承担了40%的产出。Hive导出的是纯数据文件Python负责把数据变成老板能直接看的折线图、柱状图和漏斗图。这里有两个高频坑CSV中文乱码和Matplotlib中文无法显示。Hive导出的TSV/CSV默认UTF-8编码Excel直接双击打开会乱码最简单的办法是在Python读取时固定编码import pandas as pd df pd.read_csv(uv_daily.tsv, sep\t, encodingutf-8) print(df.head())如果一定要给非技术同事发CSV可以在导出后用Python先给文件加上UTF-8 BOM头再另存为xlsx这样他们用Excel打开不会乱码。Matplotlib中文问题的解法是提前设置字体import matplotlib.pyplot as plt plt.rcParams[font.sans-serif] [SimHei] plt.rcParams[axes.unicode_minus] False这两行代码能解决大多数中文字符变成方块的问题。第一次跑图满屏乱码的时候我还以为是数据编码坏了查了半天才发现是绘图库的字体配置问题。这类“工具默认配置”造成的坑文档里不会主动写但实际项目里几乎人人都会遇到。6.3 分析报告怎么组图表才不丢信息图表的最终目的是支撑决策不是展示SQL技巧。这套项目里我沉淀了一份图表组合模板按业务阅读顺序排列。第一屏放流量大盘日活用户数和GMV趋势线用来建立整体感知第二屏放漏斗转化浏览到支付的四级漏斗图标明每一步转化率第三屏放留存曲线按新用户分群看次日、7日、30日留存走势第四屏放RFM矩阵散点图或者分群占比饼图配合运营动作建议。各指标口径必须写在图表下方备注里比如“留存率为新增访问用户的次日留存分母为2024年1月新增访问用户数”。这个备注看着啰嗦却能避免汇报会上被追问“你这个数跟报表系统怎么不一样”的尴尬。所有从Hive拿出来的数字都要能追溯到一条确定的SQL追溯不到口径的数字我一般不放上报告。可视化阶段的另一个心得是控制维度。一张图最多装两到三个维度的信息多了就成蜘蛛网。比如留存曲线可以按渠道分多条线但别再叠加价格带维度。想表达的东西太多反而什么都说不清楚。这种克制做出来的图表在实际汇报中收到的正面反馈远多于那些堆满信息的高复杂度图表。这个项目做到最后我最庆幸的是在最开始花了大半天时间做字段体检和口径梳理而不是急着写业务SQL后面所有报表都建立在这些表的基础上地基一旦歪了上层全得推倒重来。如果你也在跑类似的电商数据分析项目建议先把Raw数据吃透——每个字段的枚举值、NULL比例、时间范围都过一遍再进数仓建模。还有一点很实在的体会无论用什么工具跑数核心永远是“口径先统一再谈速度”。SQL优化解决的是跑多快的问题口径统一解决的是数对不对的问题。这顺序一旦反了后面所有分析都会在无穷的返工中打转。
返回列表