首页 / 资讯中心 / 文章详情

基于Hive的歌曲推荐候选集筛选系统实践与优化

基于Hive的歌曲推荐候选集筛选系统实践与优化 ★ FEATURED ARTICLE
做了几年大数据开发大部分时间都在跟日志清洗、报表汇总打交道。真正让我觉得“大数据能直接创造业务价值”的反而是这个基于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 TezTez相比MR快很多适合DAG多阶段作业调度Apache DolphinScheduler管理每日任务依赖和重跑结果导出生成HFile Redis候选集和特征推送到Redis供推荐服务读取元数据MySQLHive 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.compressSNAPPY);这里用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 (dt2025-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 (dt2025-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.mapfilestrue和hive.merge.size.per.task256000000合并小文件到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用到位就已经能解决大部分问题了。
阅读完成 · 觉得有帮助?
咨询建站