ARTICLE DETAIL

资讯详情

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

基于Hive离线数仓的歌曲筛选音乐推荐系统项目实践

基于Hive离线数仓的歌曲筛选音乐推荐系统项目实践 做大数据项目最怕什么不是集群挂了也不是磁盘满了而是数据堆了一堆却不知道拿它干什么。这两天整理了一个很有意思的离线项目——基于Hive的歌曲筛选音乐推荐系统正好可以回答这个问题。这个项目用一个很典型的离线数仓思路把用户行为日志、歌曲基础信息、播放记录全部落到Hive里通过分层建模、SQL清洗、特征提取和排序筛选最终输出一份“给用户推荐哪些歌”的候选列表。也就是说这既是一个歌曲筛选工具也是一个简化版的离线音乐推荐系统。我不是在写什么框架级的大工程咱们聊的是能落地、能跑通、能写到简历上的那种数据项目。如果你是大数据方向的在校生或者刚接触数仓没做过完整任务的初级工程师这篇文章应该能帮你把Hive在推荐场景里的位置一下子串起来——从“Hive就是写SQL查数据”这种散装认知过渡到“我能独立设计一套离线推荐数据流程”的完整项目视角。1. 项目整体设计与技术选型思路1.1 为什么推荐系统要选Hive而不是Spark或Flink先说一个很多新人纠结的问题音乐推荐不是应该上Spark、Flink这种实时计算引擎吗怎么用Hive这种离线批处理工具答案是取决于推荐延时的要求。如果我们做的是用户点击App后“毫秒级返回相似歌曲”那确实得靠Flink Redis这类实时链路。但绝大多数音乐推荐场景尤其是冷启动、每日歌单、相似歌手推荐、猜你喜欢这种更新频率低、候选集大的模块分钟级甚至小时级更新完全够用。Hive的优势就在于部署简单、SQL门槛低、跑一次全量离线数仓成本可控、与HDFS生态天然契合。在这个项目里推荐列表是T1更新的Hive是不二之选。从另一个角度看Spark和Flink的学习曲线陡峭很多而且如果数据量没有到PB级Hive跑的这批任务根本不存在“晚了几分钟”的体验差异。做技术选型最忌讳追新而不看场景。我经常看到有人拿Flink去算日活这种离线指标纯属大材小用集群成本和运维复杂度都上去了。把Hive这套数仓设计思路吃透后面切Spark SQL其实只是换引擎的问题SQL模型几乎可以平移这才是打基础的价值所在。1.2 数仓分层架构从原始日志到推荐候选集这个项目我按照标准数仓四层来设计ODS层原始数据落地层。用户行为日志播放、收藏、搜索、分享、歌曲基础信息表、歌手信息表原封不动地同步到Hive。DWD层数据清洗与明细层。去掉无效字段、统一时间格式、过滤机器人行为、打标歌曲风格生成一份干净的用户行为明细表和歌曲明细表。ADS层应用数据服务层。输出歌曲特征宽表、用户画像表和最终推荐结果表。为什么一定要分层直接对原始日志跑SQL不香吗第一次做项目的人确实会这么想但跑几天就发现两个问题第一上游日志字段变了下游所有逻辑都跟着崩第二清洗逻辑和推荐逻辑耦合在一起调试一次SQL得把几亿行数据重新过一遍。分层之后每一层有明确职责清洗在DWD内收敛特征在ADS内统一排查问题时可以精准定位到层。这个项目里ADS层不是简单复制明细数据而是做了几个关键动作歌曲特征宽表每首歌一个字段集合包含所属风格、时长、热度分、质量分、平均播放完成度、收藏次数。用户画像表每个用户一个JSON字段存收藏歌曲的标签分布。推荐候选结果表用户ID 候选歌曲ID 推荐分数 排序位置。1.3 技术栈选型与集群环境假设整个项目用的技术栈很收敛都是大数据圈子里最主流的离线组件组合组件用途版本建议Hadoop HDFS底层分布式存储Hadoop 3.xHive离线数仓计算引擎Hive 3.xTezHive底层计算引擎Tez 0.9.xMySQL存储推荐结果供业务读取8.xDataX采集业务库数据入仓库3.xShell/Crontab调度离线任务-集群规模模拟的是3主3从共6台机器的小集群每台16GB内存、4核CPU500GB磁盘。这个配置跑几千万行日志没问题也能撑起主要业务逻辑。如果读者用自己的测试环境甚至可以用1台机器部署Hadoop伪分布式跑通流程没问题就是看日志时要多注意区分任务耗时和数据量。为什么要加MySQL这一环因为Hive自身不擅长高并发的单条查询业务方不可能直接通过JDBC连Hive给用户展示“每日推荐歌单”。标准的做法是把Hive计算完的推荐结果通过Sqoop或DataX同步到MySQL对外提供接口查询。这个细节很关键很多毕设项目做到Hive出结果就完事了没有打通“计算结果 - 关系型数据库 - 对外服务”的最后一环面试被问起来就露馅了。2. 歌曲数据清洗与Hive表结构设计2.1 ODS层原始数据接入要点音乐平台每天产生的数据最核心的其实就是两张表歌曲基础信息表song_info和用户行为日志表user_behavior_log。前者来自业务库的MySQL通过DataX全量或增量同步到HDFS后者来自App埋点日志通过Flume或直接上传日志文件落地。这里有一个毕设项目很常见的坑直接拿全部字段建Hive表然后把所有日志一股脑塞进去。结果就是ODS层表字段特别多分区杂乱后面DWD层清洗时不知道该相信哪个字段。我的习惯是ODS层的表结构设计以“保留原始内容”为准则最多加一个etl_time字段不做任何业务判断不提前过滤。这样后续重新清洗时还有后悔药可吃。歌曲信息表的建表语句我用了ORC格式加Snappy压缩CREATE TABLE ods_song_info ( song_id STRING COMMENT 歌曲唯一ID, song_name STRING COMMENT 歌曲名称, singer STRING COMMENT 歌手, album_id STRING COMMENT 专辑ID, genre STRING COMMENT 风格标签逗号分隔, duration INT COMMENT 歌曲时长单位秒, upload_time STRING COMMENT 上传时间 yyyyMMdd, source STRING COMMENT 来源渠道, etl_time STRING COMMENT 入库时间 ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY \t STORED AS ORC TBLPROPERTIES (orc.compressSNAPPY);注意几个细节分区字段dt用字符串类型取值形如20241201这样后续按天运维非常直观。歌曲时长duration用秒不用毫秒避免不同埋点平台精度不一致的问题。genre支持多标签用逗号分隔虽然违反了第一范式但在数仓里很实用——我们后面做筛选时用like或split函数都能处理不用为了建表“优雅”而拆成子表给自己制造join负担。2.2 DWD层清洗策略去重、过滤、格式统一DWD层是整个项目最花时间的地方也是数据质量的生命线。我对user_behavior_log做的清洗主要有五步去除无效用户user_id为空或长度为0的直接过滤。去重播放记录同一用户、同一首歌、同一秒内重复的播放事件只保留一条。这个用row_number()开窗轻松搞定。过滤异常时长播放时长大于歌曲本身时长2倍的记录大概率是客户端Bug或刷量行为。统一时间格式日志里可能出现yyyyMMdd、yyyy-MM-dd HH:mm:ss、时间戳三种格式全部转成标准格式。打宽字段把歌曲ID关联上歌曲风格为后续用户画像做准备。这一步我不建议在Hive里用update做去重Hive本身不擅长行级更新常规做法是“覆盖写”先把清洗逻辑写成一个大的INSERT OVERWRITE语句结合row_number()去重再把结果写回DWD分区表。这种“读原表、算新表、覆盖分区”的方式既是Hive的强项也便于任务失败后回溯重跑。清洗效果的验证也有一招在DWD层对每一个业务规则计数比如“过滤了多少无效用户”、“去重减少了多少播放记录”。把这些运行指标写到一个单独的结果表里后续看数据质量是否下降直接查这张表就行不用每次都从头逻辑排查。2.3 ADS层歌曲特征宽表与用户画像表设计ADS层是直接面向推荐业务的表结构设计一定要让下游SQL简单。我们的歌曲特征宽表最终长这样字段名字段含义计算逻辑song_id歌曲ID基础字段genre风格标签从ODS拆分解析hot_score热度分播放次数归一化 收藏次数归一化加权求和quality_score质量分平均播放完成度映射到1-10分duration时长直接使用play_cnt_30d30天播放次数窗口聚合fav_cnt_30d30天收藏次数窗口聚合avg_play_ratio平均播放完成度播放时长 / 歌曲时长取均值用户画像表的设计也值得一提我用一个包含用户ID和一个JSON字符串的宽表存储该用户收藏歌曲的风格权重。为什么用JSON而不是多行因为推荐排序时希望一行取到所有特征信息避免多次自连接。JSON字段在Hive里可以直接用get_json_object函数解析配合后面的UDF使用非常顺手。3. 歌曲筛选规则与推荐逻辑落地3.1 筛选规则梳理业务语义优先于技术花活做歌曲筛选最忌讳上来就写SQL先得跟业务对齐“到底什么样的歌该进入推荐候选池”。我在这个项目里把筛选规则收敛成五条全部用可解释的SQL表达基础合法性歌曲状态在线不能被下架。质量底线质量分 6得分来自播放完成度。时长区间90秒到300秒之间太短像片段太长用户没耐心。时效性优先近一年内上架的新歌历史老歌只在热度极高时保留。热度下限30天播放次数 某个阈值根据整体数据分布定通常取P50百分位。这些规则没有一条是技术驱动全部来自产品经验。也不要小看这种“笨办法”很多推荐系统冷启动阶段根本没有机器学习模型就是靠规则筛掉大量长尾垃圾数据后再做排序。规则筛选的最大价值是压缩候选集合原来100万首歌筛完可能只剩5万后续计算量大幅降低。你别嫌暴力线上真实推荐流程基本都是“召回 - 粗排 - 精排”歌曲筛选就是在做召回。3.2 用户画像构建收藏行为里的标签权重画像我选收藏行为而不是播放行为原因很简单播放行为噪声太大随机播放、自动连播都会产生播放记录但用户主动点“收藏”就代表着明确的喜好表达。构建逻辑分三步统计用户收藏列表把DWD层的收藏记录按用户聚合生成用户收藏歌曲集合。关联歌曲风格标签JOIN歌曲特征表统计每个用户收藏歌曲中风格标签的出现频次。生成归一化权重用出现频次除以该用户总收藏数得到每个风格标签的占比存到画像表。比如用户A收藏了10首歌其中5首是摇滚、3首是民谣、2首是电子那画像字段就是{摇滚:0.5,民谣:0.3,电子:0.2}。这个数据看起来简单但它在后面排序阶段非常有用能够解释“为什么这个用户会优先看到某类歌”。整套思路就是从“协同过滤”简化而来的属于基于物品流行度和用户偏好向量的一种基础召回排序策略。3.3 排序加权策略不止是看相似度候选歌曲筛出来后做排序时我用了一个综合打分公式把不同的业务目标揉进去recommend_score 0.4 * user_match_score 0.3 * song_quality_score 0.2 * song_hot_score 0.1 * novelty_score其中user_match_score由用户画像里的风格权重计算song_quality_score就是歌曲特征宽表里的质量分hot_score是热度分novelty_score是“近30天播放次数少但质量高”的潜力歌曲加分。权重怎么确定的不是拍脑袋瞎调我先把各项分数分别归一化到0-1区间再拿历史两周的数据做离线AB测试看调整权重后top榜单的歌曲用户播放完成度是否有提升最终才定下这套权重。做这个项目的同学如果时间有限可以先采用默认权重跑通全流程后续再慢慢调优。4. 核心SQL实现、UDF/UDAF开发与Hive调优实录4.1 歌曲筛选核心SQL完整拆解下面这段SQL是歌曲筛选的枢纽放在ADS层的计算流程里。它做了两件事把DWD层的歌曲表和30天播放聚合结果关联然后按规则过滤。INSERT OVERWRITE TABLE ads_song_screened PARTITION(dt2024-12-01) SELECT s.song_id, s.song_name, s.singer, s.genre, s.duration, s.quality_score, s.hot_score, s.new_flag FROM dwd_song_info s JOIN ( SELECT song_id, COUNT(*) AS play_cnt_30d, COUNT(DISTINCT user_id) AS play_uv_30d FROM dwd_user_behavior_log WHERE dt 2024-11-01 AND dt 2024-11-30 AND action_type play GROUP BY song_id ) b ON s.song_id b.song_id WHERE s.song_status online AND s.quality_score 6 AND s.duration BETWEEN 90 AND 300 AND s.upload_time 2023-12-01 AND b.play_cnt_30d 1000;小细节提醒COUNT(DISTINCT user_id)在数据量大时性能比较慢如果只是过滤阈值可以先用COUNT(*)因为播放总量和UV通常强相关。如果确实需要UV口径再用COUNT(DISTINCT user_id)但要提前做好数据倾斜防护。我在这个项目里因为控制过数据量所以直接用了UV读者复现时根据自己集群能力来选择就行。再看看推荐结果表的生产SQL这里用了row_number()做全局排序再截断top50INSERT OVERWRITE TABLE ads_recommend_result PARTITION(dt2024-12-01) SELECT uid, song_id, recommend_score FROM ( SELECT p.uid, s.song_id, (0.4 * get_user_match(p.profile_json, s.genre) 0.3 * s.quality_score 0.2 * s.hot_score 0.1 * s.novelty_score) AS recommend_score, row_number() OVER (PARTITION BY p.uid ORDER BY recommend_score DESC) AS rn FROM dwd_user_profile p JOIN ads_song_screened s ) t WHERE rn 50;这里出现了一个自定义函数get_user_match我在4.2里展开讲。4.2 Hive自定义UDF开发计算用户与歌曲的匹配分get_user_match这个函数主要作用是输入用户的画像JSON字符串和歌曲的风格标签输出一个0到1的匹配分。逻辑很简单解析JSON得到用户每个风格的权重再判断歌曲风格是否在里面。如果歌曲是“摇滚民谣”用户画像里摇滚权重0.5、民谣权重0.3匹配分就取max(0.5, 0.3)再乘以一个系数。我用的Hive版本支持Java编写UDF核心代码大致是package com.example.hive.udf; import org.apache.hadoop.hive.ql.exec.Description; import org.apache.hadoop.hive.ql.exec.UDF; import org.apache.hadoop.io.Text; import com.alibaba.fastjson.JSONObject; Description(name get_user_match, value _FUNC_(json, genre) - compute match score between user profile json and song genre) public class UserMatchUDF extends UDF { public double evaluate(Text profileJson, Text genreList) { if (profileJson null || genreList null) return 0.0; JSONObject profile JSONObject.parseObject(profileJson.toString()); double maxScore 0.0; for (String genre : genreList.toString().split(,)) { Double weight profile.getDouble(genre.trim()); if (weight ! null weight maxScore) { maxScore weight; } } return maxScore; } }然后在Hive里临时添加函数ADD JAR hdfs:///path/to/udf.jar; CREATE TEMPORARY FUNCTION get_user_match AS com.example.hive.udf.UserMatchUDF;注意两点临时函数只在当前会话有效如果每天定时跑要把ADD JAR和CREATE TEMPORARY FUNCTION写进调度脚本或者把JAR放到Hive的auxlib目录下做成永久函数。UDF内部要用Text而不是String做入参类型否则会有序列化性能开销这是老Hive程序员都知道的细节。4.3 Hive UDAF开发的取舍到底什么时候该自己写再说说热词里提到的UDAF。我在这个项目中其实遇到了一个场景需要计算每个用户收藏歌曲风格的分布直方图。如果只依赖内置聚合函数可以用collect_set加侧写SQL硬凑但代码非常啰嗦而且扫描好几轮。于是我自己写了一个UDAF实现“遍历用户所有收藏风格累加权重”。UDAF开发的思路和UDF完全不一样必须实现四个方法init初始化状态、iterate每行累加、terminatePartial局部合并、terminate输出最终结果。核心代码如下public class GenreWeightUDAF extends AbstractGenericUDAFResolver { Override public GenericUDAFEvaluator getEvaluator(TypeInfo[] parameters) throws SemanticException { return new GenreWeightEvaluator(); } public static class GenreWeightEvaluator extends GenericUDAFEvaluator { private PrimitiveObjectInspector inputOI; private MapObjectInspector outputOI; Override public ObjectInspector init(Mode m, ObjectInspector[] parameters) throws HiveException { if (m Mode.PARTIAL1 || m Mode.COMPLETE) { inputOI (PrimitiveObjectInspector) parameters[0]; } outputOI ObjectInspectorFactory.getStandardMapObjectInspector( PrimitiveObjectInspectorFactory.javaStringObjectInspector, PrimitiveObjectInspectorFactory.javaDoubleObjectInspector); return outputOI; } Override public AggregationBuffer getNewAggregationBuffer() { return new GenreBuffer(); } Override public void iterate(AggregationBuffer agg, Object[] parameters) { // 解析风格字符串累加权重 } Override public Object terminate(AggregationBuffer agg) { // 正常化并返回Map } } }不过说实话我在项目里亲自写了这个UDAF之后最大的感悟是UDAF能不用就不用。除非真有复杂的多行聚合场景否则直接拆成多个SQL步骤清晰度和维护性都更好。自定义函数意味着你要维护代码、要处理类型兼容性、要考虑引擎差异这对一个以数据开发为主的项目来说是不小的额外成本。UDAF更适合做成通用工具沉淀而不是每个项目都现场写一遍。如果你在毕设里写了一个UDAF请一定把代码和说明都写完整这确实是加分项但别为了加分支写一堆完全不必要的函数。4.4 Hive调优实录小文件治理、数据倾斜与并行执行项目跑数据时第一批任务就把我卡住了几十个map任务在1分钟内结束reduce任务动辄十几分钟整个Job看起来又碎又慢。后来定位到两个主要原因一是HDFS上小文件太多二是部分热门歌曲的key发生数据倾斜。小文件问题的解法是“先合并、再计算”。我在生产流程里加了这样几步数据进入ODS时通过控制Flume的滚动周期让文件大小尽量接近128MB的块尺寸。在DWD层做INSERT OVERWRITE时设置输出参数合并小文件。周期性地用小文件治理SQL重刷分区。常用的Hive合并参数如下SET hive.merge.mapfilestrue; SET hive.merge.mapredfilestrue; SET hive.merge.size.per.task256000000; SET hive.merge.smallfiles.avgsize16000000;这里解释一下如果Task产出的平均文件小于16MB就触发合并流程目标是每个reduce产出接近256MB的大文件。合并后HDFS的NameNode内存压力明显下降查询速度也有提升。数据倾斜的处理比小文件更折腾。音乐热歌榜非常典型——周杰伦的某首歌播放量是长尾歌曲的上万倍按song_id分组时单个reduce处理的数据量远超其他reduce任务卡在99%不结束就是倾斜的典型症状。常规解法有几种热点key加随机前缀打散再二次聚合。多阶段作业把倾斜key单独拆分出来处理。不开map join时改用存文件过滤减少shuffle数据量。我在这个项目里用的是拆key方案。具体操作是把热门歌曲比如播放量超过500万次的歌筛选到一个单独的小表和普通歌曲分开聚合最后UNION ALL结果。实现不复杂但效果立竿见影job从25分钟降到7分钟。需要注意优化完要验证前后数据一致性别为了性能把结果算错了。另外开启动态并行执行也很有效SET hive.exec.paralleltrue; SET hive.exec.parallel.thread.number16;这几个参数能让Hive同时运行多个没有依赖关系的stage而不是傻等串行。比如前面4.1里的SQL先跑子查询聚合播放量、再JOIN歌曲表这两个阶段本来没有依赖关系并行后效率提升相当明显。实测在一个3节点小集群上整体数据任务耗时减少约25%。但并行开的stage多了同一时刻占用的资源也多了生产环境要根据集群实际负载来定别一把梭。5. 常见问题、踩坑记录与项目扩展建议5.1 典型问题速查表我在开发和运行这个项目的过程中整理了一份高频问题清单基本涵盖了离线数仓项目会遇到的常见坑现象根本原因解决思路按天分区跑数据时某个分区出现乱码名字串了分区字段规范不一致字符串编码或格式错误用ALTER TABLE DROP PARTITION先删除问题分区再重新排查数据来源统一分区字段为yyyyMMdd并用严格校验JOIN查询结果少数据或字段错位ORC表字段顺序与SELECT顺序不一致新旧表结构变更后未验证schema建完表后立即执行DESCRIBE验证字段变更表结构务必用ALTER TABLE ADD COLUMNS不要直接覆盖字段时间字段条件过滤不出数据表中时间存的是字符串yyyyMMdd查询用了yyyy-MM-dd统一格式from_unixtime(unix_timestamp(dt, yyyyMMdd), yyyy-MM-dd) 做转换一个reduce卡住99%数据倾斜热点key集中在一个节点定位热点key加随机前缀拆key或单独处理热点数据后再union任务输出大量几十KB的小文件默认并行度过高reduce数量太多合并小文件参数调优并重刷数据控制reduce数量与输出文件大小get_json_object解析不了中文key中文外层JSON转义问题建表时指定UTF-8字符集解析前先确保源字符串编码正确如果UDF处理统一在Java端设置response编码5.2 如果想让项目更像“生产级”还能怎么扩展文章开头提到这套推荐系统最大的优势是向下兼容、向上可扩展。如果做完T1的离线推荐后还有精力我建议按照下面三个阶段往上走第一接入调度平台。Crontab只能算玩具把Hive任务迁移到Apache DolphinScheduler或Azkaban加上任务失败重试、血缘追踪、日志可视化整个项目的工程化水平会明显提升。调度平台的价值在做大数据项目时怎么强调都不过分。第二增加实时链路。Hive负责T1的每日推荐如果需要实时热门歌曲榜可以增加Flink消费Kafka的播放日志把热门歌曲实时写入Redis对外提供毫秒级查询。此时架构就变成“Lambda架构”实时链路服务热数据、离线链路服务个性化。第三引入OLAP引擎。像热词里提到的Hive与Doris的对比基本就是离线批处理和实时OLAP的经典取舍。如果后续需要分析师自助查询歌曲指标Hive响应太慢可以在Doris里建一张汇总表通过T1同步把Hive的结果灌入Doris分析师用Doris查秒级出结果。这样整个架构的层次就非常完整Hive做离线加工、Doris做数据服务、Redis做实时缓存、MySQL做业务结果下发。5.3 做完这个项目后我的个人体会如果只选一个词总结这个项目的核心我会选“数据链路完整性”。很多人以为音乐推荐系统必须要有深度学习模型但实际工作中把数据仓库搭好、把清洗流程做扎实、把推荐候选集稳定产出来这件事本身已经解决了80%的业务诉求。模型只是锦上添花而数据是地基。地基不稳模型调得再花哨也是空中楼阁。再补一句实在话这个项目里的所有代码单拎出来任何一个SQL都不算难难的是把七零八碎的表、脚本、参数、函数组合成一个能每天自动产出结果的系统。你会碰到各种奇奇怪怪的数据问题这是正常的。一个分区丢了、一个UDF序列化报错、一次JOIN倾斜每一次排查都是在加深对Hive原理的理解。真正做完一遍你会发现自己的成长远超“会写SQL”这个层面你已经理解了整条离线数仓到底是怎么转起来的。
返回列表