基于Hadoop和大数据技术的新闻推荐系统:用户协同过滤推荐原理
基于hadoop的新闻推荐系统 用户协同过滤推荐 基于大数据的新闻推荐系统 推荐原理:以用户对新闻的喜欢和收藏行为作为基础数据集,应用hadoop通过mapreduce程序进行协同过滤计算,得出用户对新闻的预测评分,根据评分高低对新闻进行评分排序,进而推荐相应的新闻
新闻推荐系统的后台经常被问到一个问题:"每天几千万条新闻,怎么知道我会喜欢哪条?"今天咱们就扒开推荐系统的外衣,看看基于Hadoop的用户协同过滤怎么玩转新闻推荐。别被大数据吓到,核心逻辑就三句话:记录你的点击收藏行为、找到和你口味相似的用户群、把他们看过的好货推给你。

先看原始数据长啥样。用户在新闻客户端会产生三种典型行为:浏览(5秒以上计有效)、点赞、收藏。咱们用埋点日志记录这些事件:
2023-06-15 10:23:11 uid:9527 news:8848 action:read
2023-06-15 10:25:09 uid:9527 news:8848 action:like
2023-06-15 11:01:33 uid:1314 news:2266 action:collect
处理这些日志的第一关是把用户行为量化。直接上MapReduce代码片段:
// Mapper拆解原始日志
public static class BehaviorMapper extends Mapper<LongWritable, Text, Text, Text> {
public void map(LongWritable key, Text value, Context context) {
String[] parts = value.toString().split(" ");
String uid = parts[2].split(":")[1]; // 提取uid
String newsId = parts[3].split(":")[1];
String action = parts[4].split(":")[1];
context.write(new Text(uid), new Text(newsId + ":" + action));
}
}
// Reducer合并用户行为权重
public static class BehaviorReducer extends Reducer<Text, Text, Text, Text> {
public void reduce(Text key, Iterable<Text> values, Context context) {
Map<String, Integer> newsScores = new HashMap<>();
for (Text val : values) {
String[] parts = val.toString().split(":");
String newsId = parts[0];
String action = parts[1];
int score = action.equals("like") ? 2 :
action.equals("collect") ? 3 :
action.equals("read") ? 1 : 0;
newsScores.put(newsId, newsScores.getOrDefault(newsId, 0) + score);
}
// 输出格式 uid -> news1:score,news2:score
String output = StringUtils.join(newsScores.entrySet()
.stream()
.map(e -> e.getKey()+":"+e.getValue())
.collect(Collectors.toList()), ",");
context.write(key, new Text(output));
}
}
这段代码干了个脏活累活——把散落在日志里的用户行为整理成「用户-新闻评分表」。比如用户9527对新闻8848有阅读和点赞,最终这条新闻会获得1+2=3分。这里有个小技巧:收藏权重>点赞>阅读,因为主动收藏比被动阅读更能体现兴趣。
接下来进入重头戏——用户相似度计算。这里采用余弦相似度算法,公式看着唬人但其实小学生都能懂:
相似度 = 两个用户共同评分过的新闻的向量夹角余弦值
具体到MapReduce实现,需要两阶段处理。第一阶段生成用户对评分向量:
// 相似度计算Mapper
public static class SimilarityMapper extends Mapper<Text, Text, Text, Text> {
public void map(Text key, Text value, Context context) {
String[] newsScores = value.toString().split(",");
// 输出所有用户对,为后续笛卡尔积准备
context.write(new Text("ALL_USERS"), key);
// 构建评分向量
for (String ns : newsScores) {
String[] parts = ns.split(":");
context.write(new Text(parts[0]), key); // 新闻被哪些用户评分过
}
}
}
这个Mapper做了个巧妙的设计:既保留全局用户列表("ALL_USERS"键),又记录每个新闻对应的用户集合。到了Reduce阶段,系统会自动把同一个新闻下的用户聚在一起,这时候就能批量生成用户对:
// 相似度计算Reducer
public static class SimilarityReducer extends Reducer<Text, Text, Text, DoubleWritable> {
public void reduce(Text key, Iterable<Text> values, Context context) {
List<String> users = new ArrayList<>();
for (Text val : values) {
users.add(val.toString());
}
// 生成用户对
for (int i=0; i<users.size(); i++) {
for (int j=i+1; j<users.size(); j++) {
String userPair = users.get(i) + "-" + users.get(j);
// 计算余弦相似度(具体实现略)
double similarity = calculateCosineSimilarity(users.get(i), users.get(j));
context.write(new Text(userPair), new DoubleWritable(similarity));
}
}
}
}
实际生产环境这里会遇到性能瓶颈——当某个热门新闻有几百万用户评分时,用户对的组合数量会爆炸式增长。这时候得用上Combiner优化或者在Mapper端做局部聚合,避免海量数据传输。

最后预测评分阶段就简单了,找出与当前用户最相似的N个用户,把他们看过但当前用户没看过的新闻捞出来,按加权评分排序:
// 预测评分Mapper
public static class PredictMapper extends Mapper<Text, Text, Text, Text> {
public void map(Text key, Text value, Context context) {
String[] parts = key.toString().split("-");
String userA = parts[0];
String userB = parts[1];
double similarity = Double.parseDouble(value.toString());
// 只处理高相似度用户对
if (similarity > 0.6) {
context.write(new Text(userA), new Text(userB + ":" + similarity));
}
}
}
// 预测评分Reducer
public static class PredictReducer extends Reducer<Text, Text, Text, DoubleWritable> {
public void reduce(Text key, Iterable<Text> values, Context context) {
Map<String, Double> similarUsers = new HashMap<>();
for (Text val : values) {
String[] parts = val.toString().split(":");
similarUsers.put(parts[0], Double.parseDouble(parts[1]));
}
// 获取相似用户看过的新闻(需访问用户-新闻评分表)
Map<String, Double> candidateNews = getCandidateNews(similarUsers.keySet());
// 计算加权评分
for (Map.Entry<String, Double> entry : candidateNews.entrySet()) {
double totalScore = 0;
double totalSimilarity = 0;
for (Map.Entry<String, Double> userEntry : similarUsers.entrySet()) {
Double userScore = getUserNewsScore(userEntry.getKey(), entry.getKey());
if (userScore != null) {
totalScore += userScore * userEntry.getValue();
totalSimilarity += userEntry.getValue();
}
}
double predictScore = totalScore / totalSimilarity;
context.write(new Text(key + ":" + entry.getKey()), new DoubleWritable(predictScore));
}
}
}
这里有个工程实现细节——getCandidateNews方法需要访问之前的用户-新闻评分表。在Hadoop生态里,通常会用HBase或者Hive临时表来存储中间结果,避免重复计算。
整个流程跑完,你会得到每个用户对应的新闻预测评分列表。这时候只需要做个简单的排序取TopN,就能生成个性化推荐结果。不过实际操作中还要考虑新闻时效性——总不能推荐上周的旧闻吧?这时候可以在最终排序时加入时间衰减因子,让新文章获得加权优势。

说几个踩过的坑:
- 用户冷启动问题:新用户没有任何行为时,可以退回到热门推荐,但别总推娱乐八卦,根据设备信息做个粗粒度分类更靠谱
- 数据倾斜处理:80%的用户行为集中在20%的热门新闻上,解决方案是在Shuffle阶段采用二次排序,把热门新闻分散到不同Reducer处理
- 实时性短板:Hadoop批处理适合天级别更新,想要分钟级推荐得接上Flink或者Spark Streaming
这套系统在千万用户量级时,用20台hadoop节点大概能在2小时内跑完全天数据。虽然现在深度学习推荐模型大行其道,但对于中小型新闻平台来说,基于协同过滤的方案仍然是性价比最高的选择。
更多推荐
所有评论(0)