基于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,就能生成个性化推荐结果。不过实际操作中还要考虑新闻时效性——总不能推荐上周的旧闻吧?这时候可以在最终排序时加入时间衰减因子,让新文章获得加权优势。

说几个踩过的坑:

  1. 用户冷启动问题:新用户没有任何行为时,可以退回到热门推荐,但别总推娱乐八卦,根据设备信息做个粗粒度分类更靠谱
  2. 数据倾斜处理:80%的用户行为集中在20%的热门新闻上,解决方案是在Shuffle阶段采用二次排序,把热门新闻分散到不同Reducer处理
  3. 实时性短板:Hadoop批处理适合天级别更新,想要分钟级推荐得接上Flink或者Spark Streaming

这套系统在千万用户量级时,用20台hadoop节点大概能在2小时内跑完全天数据。虽然现在深度学习推荐模型大行其道,但对于中小型新闻平台来说,基于协同过滤的方案仍然是性价比最高的选择。

更多推荐