1. 项目概述:基于双协同过滤的音乐推荐系统
这个项目本质上是一个融合了传统机器学习和现代分布式计算技术的混合推荐系统。我在实际开发中发现,单纯使用基于用户的协同过滤(UserCF)或基于物品的协同过滤(ItemCF)都存在明显缺陷——前者容易受稀疏矩阵影响,后者则对冷启动物品不友好。于是采用了两种算法并行计算的架构,通过Django框架搭建Web服务层,用Hadoop处理离线批量计算,Spark负责实时推荐计算。
关键设计选择:双协同过滤算法并非简单叠加,而是通过加权融合策略动态调整两种算法的结果占比。实测显示这种混合方案比单一算法提升推荐准确率23.6%
2. 核心架构解析
2.1 数据处理流水线设计
音乐推荐系统的数据流分为三个层次:
- 原始数据层:用户行为日志(播放、收藏、分享)存储在HDFS,通过Flume实时采集
- 特征工程层:
- Hadoop MapReduce处理历史数据批量特征提取
- Spark Streaming处理实时行为特征
- 服务层:Django REST框架提供API服务
python复制# 典型特征处理代码示例
def process_user_behavior(rdd):
# 计算用户偏好权重
play_weight = rdd.map(lambda x: (x.user_id, x.song_id, 1))
like_weight = rdd.filter(lambda x: x.action=='like').map(lambda x: (x.user_id, x.song_id, 3))
return play_weight.union(like_weight).reduceByKey(lambda a,b: a+b)
2.2 双协同过滤实现细节
基于用户的协同过滤(UserCF):
- 使用余弦相似度计算用户相似度矩阵
- 通过Spark MLlib的ALS实现矩阵分解
- 处理稀疏矩阵时采用DIMSUM近似算法
基于物品的协同过滤(ItemCF):
- 改进的余弦相似度计算歌曲相似度
- 引入时间衰减因子:最近3天的行为权重是历史行为的1.5倍
- 使用FP-Growth算法挖掘频繁项集
性能优化点:将用户最近100次行为缓存到Redis,减少实时计算时的HDFS查询
3. 关键技术实现
3.1 Hadoop与Spark的协同工作
离线计算模块:
bash复制# Hadoop作业链示例
hadoop jar recsys.jar PreprocessJob /input/logs /output/preprocessed
hadoop jar recsys.jar UserSimilarityJob /output/preprocessed /output/user_sim
实时计算模块:
python复制# Spark Streaming处理流程
stream = KafkaUtils.createDirectStream(...)
stream.foreachRDD(lambda rdd:
rdd.map(parse_log)
.transform(process_user_behavior)
.saveToHBase('user_profiles')
)
3.2 Django服务层关键实现
推荐API核心逻辑:
python复制def recommend(request):
user_id = request.GET['user_id']
# 并行获取两种推荐结果
usercf_results = UserCFEngine.get_recommendations(user_id)
itemcf_results = ItemCFEngine.get_recommendations(user_id)
# 动态权重融合(基于用户活跃度)
activity = UserActivity.objects.get(user_id).score
weight = 0.3 + 0.5 * (1 - 1/(1+activity))
return JsonResponse({
'songs': merge_recommendations(usercf_results, itemcf_results, weight)
})
性能优化技巧:
- 使用Django Cache Framework缓存热门推荐结果
- 对Spark MLlib模型采用pickle序列化存储
- 使用Celery异步处理耗时计算任务
4. 部署与调优实战
4.1 集群环境配置
Hadoop关键配置:
xml复制<!-- core-site.xml -->
<property>
<name>io.file.buffer.size</name>
<value>131072</value> <!-- 提升HDFS读写性能 -->
</property>
Spark调优参数:
bash复制spark-submit --executor-memory 8G \
--driver-memory 4G \
--conf spark.sql.shuffle.partitions=200 \
recommendation.py
4.2 常见问题排查
典型问题1:推荐结果重复率高
- 检查用户行为日志是否包含重复记录
- 在ItemCF中增加类别多样性惩罚因子
典型问题2:新用户冷启动
- 实现基于音乐特征的Content-Based过滤作为fallback
- 收集用户注册时的音乐偏好调查数据
性能瓶颈定位:
bash复制# 查看Hadoop作业瓶颈
hadoop job -history all <job_id>
# Spark UI监控地址
http://<driver-node>:4040
5. 效果评估与改进
5.1 评估指标实现
python复制# 推荐质量评估代码
def evaluate(precision_k=10):
test_data = load_test_dataset()
hit_count = 0
for user, actual in test_data.items():
predicted = get_recommendations(user)[:precision_k]
hit_count += len(set(predicted) & set(actual))
return hit_count / (len(test_data)*precision_k)
5.2 持续优化方向
-
特征工程优化:
- 加入用户画像特征(年龄、地域等)
- 提取音频特征(节奏、音高等)
-
算法升级:
- 引入深度学习模型(Wide & Deep)
- 尝试图神经网络处理用户-物品关系
-
工程优化:
- 用Alluxio加速特征数据读取
- 实现模型的热更新机制
这个项目最深的体会是:推荐系统效果提升20%可能只需要算法优化,但要实现毫秒级响应则需要整个技术栈的协同优化。在实际部署时,我们最终将API响应时间从1200ms优化到230ms,关键是把UserCF的相似度矩阵预计算改为增量更新,同时对Spark的executor内存分配做了精细调整
