做了几年大数据开发,大部分时间都在跟日志清洗、报表汇总打交道。真正让我觉得“大数据能直接创造业务价值”的,反而是这个基于Hive的歌曲筛选音乐推荐系统。项目本身不算宏大,但它是典型的“大数据+推荐系统”落地场景:用Hive处理上亿条用户行为记录,从中筛选出值得推荐的歌曲候选集,再交给下游推荐引擎做排序。这套流程跑通之后,我才意识到,推荐系统的天花板往往不在精排模型,而在数据入口的筛选质量。
这篇文章就把整个项目的设计思路、Hive SQL实现、性能优化和上线后的坑完整复盘一遍。如果你正在做音乐、短视频、资讯类推荐,或者准备把Hive用在推荐候选集生成这个环节,这篇文章可以直接当作参考。
1. 项目切入:推荐系统前面为什么要加一道Hive筛选关
1.1 从推荐漏斗说起
任何推荐系统,本质上都是一个漏斗:全量歌曲池 → 候选集 → 粗排 → 精排 → 最终推荐列表。我见过不少团队把精力全砸在精排模型上,却忽略了候选集这层。结果模型再花哨,喂进去的候选集如果有一堆冷门垃圾或重复歌曲,排序效果也上不去。
这个项目的核心目标,就是替代原来那套跑在业务库里的Java定时任务,用Hive构建一个离线歌曲候选集生产线。原始数据包括用户播放日志、收藏记录、搜索点击、歌曲基础信息表,日均新增数据量在亿级。过去用MySQL直接聚合,跑一个多小时是常态,而且严重拖垮线上库。迁到Hive之后,同样逻辑压缩到十几分钟,还顺手解决了历史数据回溯的问题。
1.2 为什么选型Hive而不是Spark/Flink
很多人一听说“实时推荐”就开喷,“都什么年代了还用Hive离线”。这里需要说清楚:音乐推荐对实时性要求没那么恐怖。用户听歌行为可以允许几分钟甚至小时级延迟,重点是稳定和海量数据的批处理能力。Hive的优点恰恰在这:SQL语义清晰,易维护,跑批稳,出错能重跑,离线链路排查问题也直观。
当然,实时部分不是没有。项目里实时计算用了Flink处理用户最近5分钟的播放行为,但最终融合特征时还是会把Hive产出的离线候选集作为主要底座。换句话说,Hive负责“今天该推哪些歌”,Flink负责“用户此刻正在听什么歌,微调排序”。
1.3 项目整体技术栈
这套系统的运行环境不算复杂,但足够说明一个完整的离线推荐数据链路:
| 模块 | 选型 | 说明 |
|---|---|---|
| 数据存储 | HDFS | 原始日志和Hive表统一落在HDFS |
| 计算引擎 | Hive 3.1.3 + Tez | Tez相比MR快很多,适合DAG多阶段作业 |
| 调度 | Apache DolphinScheduler | 管理每日任务依赖和重跑 |
| 结果导出 | 生成HFile + Redis | 候选集和特征推送到Redis供推荐服务读取 |
| 元数据 | MySQL | Hive metastore,初始配置时踩了不少坑 |
这里有一个比较重要的经验:Hive版本尽量用3.x,配套的Hive on Tez部署时要注意Tez的tez-site.xml配置,尤其是tez.container.size,设置不当容易导致Container频繁OOM。项目里用的是Hive 3.1.3,配合CDH调整过的Tez容器参数跑起来很稳。
2. 歌曲数据仓库建模:从原始日志到可计算的歌曲宽表
2.1 数据分层设计:ODS、DWD、DWS
这套歌曲筛选系统没有复杂的算法,但数据分层做得比较规矩。整个数仓分了三层:
- ODS层:原始播放日志、收藏日志、歌曲信息表,每天全量/增量落一份,不加工。
- DWD层:清洗、脱敏、解析后的明细事实表,按天分区,字段标准化。
- DWS层:以歌曲和歌手为粒度聚合的指标宽表,直接供筛选使用。
这样设计最大的好处是:指标口径统一。团队里如果有人问“这首歌昨天播放量到底是多少”,不用各自写一遍SQL,直接查DWS层即可。
2.2 歌曲维度表设计:唯一标识必须稳
做歌曲筛选,先得有干净的主数据。歌曲维度表设计时,我反复强调一个点:歌曲ID必须全局唯一且稳定。很多音乐平台历史上踩过坑,同一首歌因为版权方不同,在库里存在多条记录,这就导致播放日志join歌曲表时炸出重复计数。
这张维度表至少包含这些字段:
CREATE TABLE dwd_song_info_d ( song_id STRING COMMENT '歌曲全局唯一ID', song_name STRING COMMENT '歌曲名称', artist_id STRING COMMENT '歌手ID', artist_name STRING COMMENT '歌手名称', album_id STRING COMMENT '专辑ID', genre STRING COMMENT '曲风:流行/摇滚/民谣等', language STRING COMMENT '语种', duration_seconds BIGINT COMMENT '时长秒数', publish_date STRING COMMENT '发行日期', is_original TINYINT COMMENT '是否原创', status TINYINT COMMENT '歌曲状态:0下架 1可播放', etl_time STRING COMMENT 'ETL时间' ) PARTITIONED BY (dt STRING) STORED AS ORC TBLPROPERTIES ('orc.compress'='SNAPPY');这里用ORC加Snappy压缩是刻意的。ORC列式存储在读取少量列时效率极高,Snappy压缩则平衡了体积和CPU开销。三层求和下来的压缩率大约是原始文本的1/6,跑了三个月的数据量也不会把NameNode压垮。
2.3 用户行为事实表:清洗重点不在SQL,在规则
播放日志是筛选的核心输入。原始日志长什么样?基本是前端埋点上报的一堆JSON字符串。清洗过程主要做四件事:
- 去重:用户连续播放同一首歌只算一条有效记录,防止刷播放量。
- 过滤异常时长:播放时长超过歌曲时长1.5倍的记录直接丢弃,这类多是缓冲卡顿导致的重复上报。
- 有效播放判定:播放时长必须大于30秒或超过歌曲时长50%,才算有效播放。这个口径直接影响热度指标。
- 补充维度:通过song_id关联歌曲表,补上歌手ID、曲风、语种等字段,方便后续聚合。
清洗后的明细表结构大致如下:
CREATE TABLE dwd_user_play_d ( song_id STRING, user_id STRING, play_ts BIGINT, play_date STRING, play_duration BIGINT, is_valid TINYINT, artist_id STRING, genre STRING, source STRING COMMENT '播放来源:推荐/搜索/歌单/日推', dt STRING ) PARTITIONED BY (dt STRING);清洗后的数据按天分区。每天凌晨调度任务先跑ODS到DWD,再跑DWD到DWS,整个链路像流水线一样稳定推进。
2.4 歌曲特征宽表:把指标提前算好
筛选候选集需要大量指标,每次都临时聚合肯定不行。所以项目里专门构建了一张歌曲特征宽表dws_song_feature_d,按日累计更新,包含以下核心指标:
- 当日播放量、7日播放量、30日播放量
- 当日收藏量、7日收藏量
- 当日有效播放率、平均播放时长
- 完播率(播放时长/歌曲时长)
- 搜索点击量(反映主动意图)
- 推荐位点击率(反映用户对推荐结果的接受度)
提前算好这些指标,筛选SQL写起来就是简单的where条件比较。这也是Hive数仓的核心思想:把复杂的计算提前到调度链路里完成,下游消费数据时只做轻量筛选。
3. 歌曲筛选策略的Hive SQL实现:从热度到个性化
3.1 热门歌曲筛选:计算口径先统一
最朴素的筛选逻辑就是“只推热门歌”。但热门不能只看原始播放量,否则那些推广位资源多、曝光量大的歌曲永远霸榜,用户很快审美疲劳。
我用的热度分计算公式是这样的:
hot_score = w1 * 播放量忇准值 + w2 * 收藏量忇准值 + w3 * 搜索量忇准值 + w4 * 完播率忇准值每个指标先做Min-Max归一化。为什么要归一化?因为播放量动辄百万量级,完播率只有0到1,如果不归一化,完播率那个维度直接被吃掉,相当于权重失效。以下是核心SQL逻辑:
INSERT OVERWRITE TABLE dws_song_feature_d PARTITION (dt='2025-01-01') SELECT song_id, play_cnt_1d, fav_cnt_1d, search_cnt_1d, finished_rate, round( 0.4 * (play_cnt_1d / max_play) + 0.3 * (fav_cnt_1d / max_fav) + 0.2 * (search_cnt_1d / max_search) + 0.1 * finished_rate, 4 ) AS hot_score FROM ( -- 子查询计算各指标 ) t;3.2 异常值剔除:percentile_approx扛大梁
热门筛选有一个常见的坑:统计口径被异常值污染。有些歌曲因为平台活动、新闻事件播放量突然冲高,如果直接按热度分排序,这种“假热门”会挤掉真正被用户持续喜爱的歌曲。
这时候percentile_approx函数特别好用。它可以近似计算分位数,从全量数据里找出播放量的合理区间,把超过99分位的极端值单独打标,或者在归一化时用99分位代替最大值,避免被单个异常值拉偏。
这个热搜词在标题里出来了,估计不少人也踩过。用法很简单:
SELECT percentile_approx(play_cnt_1d, 0.99) AS p99_play_cnt FROM dws_song_feature_d WHERE dt = '2025-01-01';项目里我把99分位的播放量作为归一化的上限,超过这个值的一律按1算。这样既保留了一些头部爆款,又不会让它们把其余歌曲的分数压得太低。实测下来,推荐列表的长尾覆盖率明显提升。
3.3 规则筛选与候选集生成:去重、去冷门、做时间衰减
拿到特征宽表之后,筛选逻辑用一套规则SQL组合实现:
INSERT OVERWRITE TABLE ads_song_candidate_d PARTITION (dt='2025-01-01') SELECT song_id, artist_id, genre, hot_score, play_cnt_7d, finished_rate FROM dws_song_feature_d WHERE dt = '2025-01-01' AND status = 1 -- 只推可播放歌曲 AND publish_date >= date_sub('2025-01-01', 730) -- 两年内的歌 AND play_cnt_7d >= 1000 -- 7日播放量门槛 AND finished_rate >= 0.2 -- 完播率门槛 AND hot_score > 0.01 -- 基础热度门槛 SORT BY hot_score DESC;这里有两个细节值得单独说:
时间衰减。音乐消费有很强的时效性,去年爆火的歌今年不一定适合推荐。我设计了一个衰减系数:超过30天未更新的老歌,热度分按指数衰减;新发布的歌在初始两周内有加权。衰减逻辑不放在维度表里,而是放在候选集SQL中,用exp()函数按当前日期与发布日期的差值动态计算。
去重策略。同一首歌的不同版本(Live、翻唱、Remix)需要在这里去重,我按artist_name + song_name做分组取热度分最高的版本保留。用MD5加密生成分组键,避免中文和特殊字符导致的分组不一致问题。
3.4 冷启动歌曲的特殊处理
热门筛选只是基础盘,冷启动处理才能真正体现筛选系统的价值。新歌没有历史播放量,按规则筛选会被直接过滤掉,因此单独建一个新歌池子:
- 上线时间在7天内
- 歌曲质量评分(人工标注 + 音质检测)达标
- 同类型歌手历史表现作为先验信号,给一个加权初始热度
新歌池的初始热度不做硬性门槛,而是在推荐时用小流量实验的方式逐步放量:先推给少量偏好匹配度高的用户,观察点击率和完播率,达标后自动转入正式候选集。这套逻辑相当于给新歌一个“考试期”。
4. 数据倾斜与小文件治理:Hive跑批的真实性能瓶颈
4.1 音乐场景最容易触发数据倾斜的三个地方
整个系统的首版上线不太平,最头疼的就是跑批越来越慢,从15分钟恶化到2小时。排查下来都是数据倾斜。
第一个坑是Join歌手维度表。头部歌手和长尾歌手的歌曲数量差距悬殊,按artist_id做join时,头部歌手对应的ReduceTask数据量爆炸,其他Task却空转。
第二个坑是GROUP BY中的热门歌曲。聚合函数按key分发,如果某个热搜歌曲数据量是其他歌曲的几百倍,单个Reducer被拖死。
第三个坑是动态分区写入。按发布时间做动态分区时,某些日期分区涌入海量数据,容易触发布隆或短时拥塞。
4.2 数据倾斜解决方案实战
针对这些坑,我不推荐一刀切用skewjoin,而是结合场景处理:
第一个场景,join倾斜,我先用salt加盐打散热点Key:
-- 将歌手维表按歌手ID加盐,播放表同样加盐后join SELECT a.song_id, b.artist_name FROM ( SELECT song_id, artist_id, concat(artist_id, '_', ceil(rand() * 10)) AS salted_artist_id FROM dwd_user_play_d WHERE dt = '2025-01-01' ) a JOIN ( SELECT artist_id, concat(artist_id, '_', suffix) AS salted_artist_id, artist_name FROM dim_artist_d LATERAL VIEW explode(array(0,1,2,3,4,5,6,7,8,9)) t AS suffix ) b ON a.salted_artist_id = b.salted_artist_id;第二个场景,聚合倾斜,我先把热点Key单独抽出聚合,再和非热点结果合并。热点判定用play_cnt_7d > 阈值,这个阈值通过观察分布确定。
第三个场景,动态分区写入,我改成固定分区加后续的MSCK REPAIR或按日期范围分批次写入。
4.3 小文件问题:从几千个文件到几百个
热搜词里有“hive优化小文件”,这绝对是跑批系统的隐藏杀手。音乐播放表按天分区,如果每个小时都有一堆小任务写数据,一天下来一个分区可能产生几千个小文件,每个只有几MB甚至几百KB。HDFS的NameNode内存被这些文件元数据吃光,查询时Map数激增,调度开销比计算本身还大。
我的优化方案分三层:
- 写入端控制:在Hive中设置
hive.merge.mapfiles=true和hive.merge.size.per.task=256000000,合并小文件到256MB左右。 - 定期合并:每天调度里加一个专门的
_merge任务,用INSERT OVERWRITE重写大分区,按歌曲ID做分桶写入。 - 合理分桶:歌曲表按
song_id做分桶,桶数固定,避免文件越积越多。
实施之后,同样的跑批任务文件数从日均5000+降到600以内,查询耗时下降了差不多一半。
4.4 一次真实故障排查记录
这里记录一个让我印象深刻的故障。上线初期,候选集表ads_song_candidate_d每天凌晨5点应该生成,但经常拖到8点。排查链路如下:
- 先看调度DAG,发现瓶颈在DWD层播放明细表重跑任务,耗时3.5小时。
- 点进任务看日志,发现Reducer数量恒定在20个左右,但数据量是平时的2倍。
- 进一步查发现某个播放日志源重复导入了前一天的存量数据,ODS层没有做幂等去重。
- 在ODS导入任务加上
INSERT OVERWRITE前先TRUNCATE对应分区,同时在DWD清洗时加ROW_NUMBER()窗口函数严格去重,问题解决。
这次故障的教训是:数仓链路里,上游没做幂等,下游SQL写得再高效也没用。我后来给ODS所有日分区任务都加了重跑前删分区的逻辑,算是彻底堵住了这个坑。
5. Hive筛选结果如何对接推荐系统
5.1 候选集导出链路设计
Hive算完候选集和特征,最终要能暴露给线上推荐服务。项目里主流的导出路径是:
Hive表 → 生成HFile → 批量导入HBase → HBase客户端直查
之所以选HBase而不是直接Redis,是因为候选集数据量有千万级,Redis全量加载很吃内存,HBase则天然适合海量KeyValue存储和范围扫描。线上推荐服务读取候选集的时候,走的是HBase的批量Get,单次RT能控制在10ms以内。
5.2 特征文件周期调度
特征宽表和候选集都是每天凌晨定时生成。调度配置大概是:
- 02:00 ODS日志数据抽取
- 03:00 DWD清洗任务
- 04:30 DWS特征聚合
- 05:30 候选集SQL + HFile生成
- 06:30 HBase批量导入
- 07:00 线上服务自动加载新数据
时间上留足缓冲,因为Hive跑批的耗时会有波动。调度系统选用DolphinScheduler,它自带失败重试和依赖管理,对Hive任务的DAG编排非常友好。
5.3 推荐服务如何消费这批数据
推荐服务拿到HBase里的候选集后,并不是直接返回,而是有一层轻量级的rerank逻辑:
- 先根据用户最近的播放行为做粗过滤,把用户不喜欢的曲风权重调低。
- 再结合实时的收藏/跳过行为做微调。
- 最后用多样性规则保证推荐列表不出现同一个歌手的3首以上歌曲。
Hive在中间扮演的角色是“把该推的候选集全部准备好”,后续模型和规则只做最终排序。这个分工让团队里算法、数据、后端各自专注,不用互相等。
6. 项目复盘:这套系统的边界和进阶空间
6.1 值得保留的设计
如果让我重新做一次,以下设计我会原样保留:
第一,歌曲特征宽表。看着简单,实际解决了大量重复计算问题。所有下游任务包括实时任务都直接查这张表,效率和稳定性都上来了。
第二,百分比剔异常值。用percentile_approx做归一化上限是一种很实用的工程技巧。它相当于给热门指标加了软保护,让分数分布更接近真实用户偏好。
第三,分层分桶加ORC压缩。这套组合让存储和计算性能都很均衡。别小看文件格式,我见过团队用TextFile跑推荐统计,数据量翻三倍后作业直接跑不动。
6.2 如果重来一次我会改的地方
有两个地方,如果重新做,我会提前规划:
一是接入实时计算的时间点。前期只做离线跑批,后来实时特征接入时发现,离线DWS和实时特征的字段口径没有完全对齐,比如“有效播放”的判定条件离线用了30秒,实时用了50%,两边对不上,排查半天。建议一开始就统一定义指标口径,离线实时共用一套口径文档。
二是候选集生成最好留一个规则配置平台,而不是直接改SQL。运营想要临时调整歌曲门槛条件,比如“今天重点推国风歌曲”,就直接改配置,不用等开发排期。后期我做了个简单的规则配置表,把筛选条件参数化,运营自己就能操作。
6.3 对做同类系统的建议
最后说点实在的。如果你们也在做一个基于Hive的推荐候选集系统,我总结的验收标准就三条:
- 数据质量:重复播放、异常时长、下架歌曲这些最脏的数据最先处理,数据不干净,后面一切白搭。
- 任务稳定:调度一定要设计好,失败重跑机制、幂等写入、超时告警缺一不可。
- 口径统一:从第一天就开始维护指标字典,别等到算法团队拿着离线特征和实时特征对比的时候才意识到问题。
这个项目做完之后,我对Hive的看法有了实质性的改变。以前总觉得Hive就是个跑报表的老古董,真正深入进来才发现,在离线大数据链路上,它的稳定性、生态成熟度和排查便利性,依然是不可替代的选择。推荐系统的数据底座,不一定要堆一堆花哨的框架,把Hive用到位,就已经能解决大部分问题了。