ARTICLE DETAIL

资讯详情

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

基于大数据计算引擎的新闻推荐系统:Flink+Spark批流分工实战

基于大数据计算引擎的新闻推荐系统:Flink+Spark批流分工实战 简介这是一份面向高校学生与大数据入门者的新闻推荐系统课程设计实战包围绕大数据计算引擎与用户画像两条主线解决个性化新闻推送从数据到算法的落地问题。包内共1789个文件以315个Python源码、595个JavaScript前端脚本、74个HTML页面及318个pyc编译文件为主另含Java、Scala、PHP等实现与css、jpg、xml等静态资源压缩包约25.6MB目录结构完整便于按模块检索。资源覆盖大数据计算引擎、用户画像构建、基于内容与协同过滤及矩阵分解的推荐算法、数据预处理与NLP分词、Flink或Spark Streaming实时推荐、准确率召回率等评估指标与A/B测试优化并兼顾多样性与隐私等用户体验考量。已有517人学习下载适合作为理论结合实践的参考案例帮助读者理解推荐系统各环节的具体实现与工程组织方式。1. 新闻推荐系统为什么值得用大数据计算引擎重做一遍如果你手头有一个新闻类产品日活刚过万推荐逻辑还停留在「按热度排序 人工运营标签」那你大概率已经撞上这堵墙了用户刷了两屏就退出点击率卡在 3% 上不去运营每天手动配头条配到凌晨。问题不在推荐算法本身而在于数据量一旦上来单机脚本跑一次全量用户画像要四十分钟实时行为日志根本来不及进模型。基于大数据计算引擎的新闻推荐系统解决的正是这个断层——把离线批处理、实时流计算和推荐召回排序串成一条流水线让新闻feed 在秒级内完成「用户此刻想看什么」的判断。这套方案适合有 Java 或 Scala 基础、手头有 3 台以上机器或云主机的后端与数据开发也适合正在做大数据技术毕业设计、需要一套能跑通全链路项目的学生。下面我按自己搭过两版的路径把选型、代码和踩过的坑一次讲清。2. 计算引擎选型与新闻推荐链路拆解批流怎么分工2.1 为什么新闻推荐不能只用离线批处理新闻场景和电商推荐最大的区别是时效性。一条突发新闻发布后 15 分钟内如果没有推给可能感兴趣的用户它的点击价值就衰减大半。纯离线批处理比如每天凌晨跑一次 Spark 任务更新用户兴趣向量在电商场景勉强够用因为商品生命周期以周计但新闻的生命周期以小时甚至分钟计。我第一版系统就是纯离线架构Hive 存用户行为日志Spark 每天凌晨算一次 ItemCF 相似度矩阵结果第二天推荐的全是昨天已经看过三遍的旧闻。用户投诉「推荐的都是旧闻」时我才意识到离线链路只能解决「用户长期兴趣」解决不了「当下热点匹配」。所以新闻推荐系统的计算引擎必须做批流分工链路计算引擎处理内容延迟要求离线批处理Spark on YARN用户长期画像、物品相似度矩阵、模型训练小时级实时流计算Flink用户实时行为解析、热点新闻统计、在线特征拼接秒级在线服务Redis 推荐 API召回结果缓存、排序打分毫秒级这个分工不是拍脑袋定的。Spark 擅长大规模全量计算一次处理 TB 级历史日志没问题Flink 的原生流处理模型事件时间 水位线能处理乱序到达的点击日志保证「用户 3 秒前点了一条体育新闻」这件事在 1 秒内反映到推荐结果里。2.2 新闻推荐系统的四层链路从日志到 feed 流把整条链路拆开看从用户行为产生到推荐结果展示一共经过四层第一层数据采集与接入。客户端埋点上报用户曝光、点击、停留时长、滑动跳过等事件通过 Kafka 统一接入。Kafka 的 topic 按事件类型分news_click、news_expose、news_dwell。分区数按峰值吞吐定我一般设 6 个分区起步日活十万以内够用。第二层实时特征计算。Flink 消费 Kafka 中的行为事件做三件事一是维护每个用户的实时兴趣队列最近 50 次点击的新闻类别分布二是统计每篇新闻的实时 CTR曝光点击比用于热点发现三是把用户实时特征写入 Redis供在线排序使用。第三层离线模型训练。Spark 每天凌晨跑三个任务从 Hive 读前一天的全量行为日志训练 ALS 协同过滤模型产出用户-物品隐向量计算新闻内容的 TF-IDF 向量用于冷启动召回更新用户长期兴趣标签。模型文件写入 HDFS在线服务定时加载。第四层在线召回与排序。推荐 API 收到请求后并行从三路召回实时热点召回Redis 中的高 CTR 新闻、协同过滤召回Faiss 向量检索、内容相似召回TF-IDF 倒排。三路结果合并去重后用 LR 或 GBDT 排序模型打分返回 Top N。2.3 最小可跑通的集群规划与资源参数如果你只有 3 台机器可以这样分配Master 节点1 台8C16GNameNode ResourceManager Kafka BrokerWorker 节点2 台8C32GNodeManager Flink TaskManager Redis关键参数配置以 Flink 为例# flink-conf.yaml 核心参数 jobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 8192m taskmanager.numberOfTaskSlots: 4 parallelism.default: 4 state.backend: rocksdb state.checkpoints.dir: hdfs://master:9000/flink/checkpoints execution.checkpointing.interval: 60000taskmanager.memory.process.size设 8G 是因为 RocksDB 状态后端在维护用户实时兴趣队列时会占用较多堆外内存。numberOfTaskSlots设 4 意味着每个 TaskManager 可以并行跑 4 个算子链2 个 Worker 总共 8 个 slot对应parallelism.default: 4留了一倍余量。execution.checkpointing.interval设 60 秒是权衡太短会增加 HDFS 写入压力太长则故障恢复时丢数据多。注意如果机器内存小于 32G把taskmanager.memory.process.size降到 4096m同时把state.backend改成filesystem用内存换稳定性。3. 用 Flink 做实时特征计算从 Kafka 到 Redis 的完整代码3.1 用户实时兴趣队列的 Flink 实现实时特征的核心是维护一个滑动窗口内的用户兴趣分布。下面这段代码消费 Kafka 点击事件按用户分组用 30 分钟滑动窗口统计类别分布写入 Redis// UserInterestJob.java StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(60000); env.setStateBackend(new RocksDBStateBackend(hdfs://master:9000/flink/checkpoints)); // 1. 消费 Kafka 点击事件 DataStreamString clickStream env.addSource( new FlinkKafkaConsumer(news_click, new SimpleStringSchema(), KafkaConfig.getProps()) ); // 2. 解析 JSON提取 userId、category、timestamp DataStreamClickEvent events clickStream .map(new ClickEventParser()) .assignTimestampsAndWatermarks( WatermarkStrategy.ClickEventforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((e, ts) - e.getTimestamp()) ); // 3. 按用户分组30 分钟滑动窗口每 5 分钟触发一次 DataStreamUserInterest interest events .keyBy(ClickEvent::getUserId) .window(SlidingEventTimeWindows.of(Time.minutes(30), Time.minutes(5))) .aggregate(new CategoryCountAgg(), new CategoryWindowFunction()); // 4. 写入 Redis interest.addSink(new RedisSink(RedisConfig.getConf(), new UserInterestRedisMapper()));forBoundedOutOfOrderness(Duration.ofSeconds(5))允许事件迟到 5 秒超过 5 秒的迟到数据会被丢弃。这个值根据客户端上报延迟分布来定我实测移动端 4G 网络下 P99 延迟在 3 秒左右设 5 秒是安全值。SlidingEventTimeWindows.of(Time.minutes(30), Time.minutes(5))表示窗口长度 30 分钟、滑动步长 5 分钟意味着每 5 分钟输出一次「最近 30 分钟用户兴趣分布」既保证实时性又避免过于频繁写 Redis。3.2 热点新闻实时统计与 Redis 写入热点新闻统计逻辑类似但按新闻 ID 分组统计 5 分钟内的曝光和点击// HotNewsJob.java DataStreamExposeEvent exposeStream env.addSource( new FlinkKafkaConsumer(news_expose, new SimpleStringSchema(), props) ).map(new ExposeEventParser()) .assignTimestampsAndWatermarks(WatermarkStrategy .ExposeEventforBoundedOutOfOrderness(Duration.ofSeconds(3)) .withTimestampAssigner((e, ts) - e.getTimestamp())); DataStreamClickEvent clickStream env.addSource( new FlinkKafkaConsumer(news_click, new SimpleStringSchema(), props) ).map(new ClickEventParser()) .assignTimestampsAndWatermarks(WatermarkStrategy .ClickEventforBoundedOutOfOrderness(Duration.ofSeconds(3)) .withTimestampAssigner((e, ts) - e.getTimestamp())); // 双流 join按新闻 ID 关联曝光和点击 DataStreamNewsCtr ctrStream exposeStream .keyBy(ExposeEvent::getNewsId) .intervalJoin(clickStream.keyBy(ClickEvent::getNewsId)) .between(Time.seconds(-10), Time.seconds(10)) .process(new CtrCalculateProcess()); ctrStream.addSink(new RedisSink(RedisConfig.getConf(), new NewsCtrRedisMapper()));intervalJoin的between(Time.seconds(-10), Time.seconds(10))表示曝光事件前后 10 秒内的点击算作有效点击。这个窗口不能太大否则会把「用户看了 5 分钟后才点」也算进去虚高 CTR也不能太小否则漏掉正常点击。10 秒是新闻场景的经验值。3.3 离线 Spark 任务ALS 模型训练与向量产出离线部分用 Spark MLlib 的 ALS 训练协同过滤模型# als_train.py from pyspark.ml.recommendation import ALS from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(NewsALS) \ .config(spark.sql.shuffle.partitions, 200) \ .getOrCreate() # 读取 Hive 中的用户-新闻交互表 df spark.sql( SELECT user_id, news_id, SUM(CASE WHEN event_typeclick THEN 3 WHEN event_typedwell AND dwell_time10 THEN 2 ELSE 1 END) AS rating FROM dwd_user_behavior WHERE dt date_sub(current_date(), 1) GROUP BY user_id, news_id ) # 训练 ALS als ALS( rank64, maxIter15, regParam0.01, userColuser_id, itemColnews_id, ratingColrating, coldStartStrategydrop, nonnegativeTrue ) model als.fit(df) # 产出用户和物品向量 user_vec model.userFactors item_vec model.itemFactors user_vec.write.parquet(hdfs://master:9000/model/user_vec) item_vec.write.parquet(hdfs://master:9000/model/item_vec)rank64是隐向量维度新闻场景下 64 维足够表达用户兴趣再高容易过拟合且增加在线检索耗时。regParam0.01是正则化系数防止评分矩阵稀疏导致过拟合。coldStartStrategydrop表示预测时丢弃冷启动用户这些用户走内容召回兜底。4. 避坑与排查新闻推荐系统上线后最容易翻车的 5 个点4.1 现象推荐结果全是旧闻实时特征没生效原因Flink 任务写入 Redis 的 key 没有设 TTL或者 TTL 设得太长比如 24 小时导致用户兴趣队列一直不更新。另一个可能是 Flink 的 checkpoint 失败后任务重启状态从旧 checkpoint 恢复丢失了最近几小时的数据。解决Redis key 的 TTL 设为窗口长度的 2 倍30 分钟窗口设 60 分钟保证过期数据自动清理。同时监控 Flink checkpoint 成功率低于 95% 要查 HDFS 写入是否正常。4.2 现象热点新闻 CTR 虚高运营手动推的新闻霸榜原因曝光事件和点击事件的 join 窗口太大或者曝光事件重复上报。客户端在列表页滑动时同一篇新闻可能被上报多次曝光导致分母虚高、CTR 虚低反过来如果点击事件重复上报CTR 虚高。解决在 Flink 解析层做去重按(userId, newsId, timestamp/10)做 key 去重10 秒内同一用户对同一新闻的重复曝光只算一次。同时把 join 窗口从 10 秒收紧到 5 秒。4.3 现象Spark ALS 训练任务 OOMExecutor 频繁挂掉原因spark.sql.shuffle.partitions设得太小比如默认 200数据量大时每个 partition 数据倾斜严重或者rank设得过高用户-物品评分矩阵在 ALS 迭代时内存爆炸。解决把spark.sql.shuffle.partitions调到 500 到 1000同时开启spark.sql.adaptive.enabledtrue让 Spark 自动处理倾斜。rank从 64 降到 32 试试如果效果没明显下降就用 32。4.4 现象在线推荐 API 响应时间超过 500ms用户滑动卡顿原因三路召回串行执行每路都要查 Redis 或 Faiss累加起来就慢了。或者 Faiss 索引没有加载到内存每次查询都读磁盘。解决三路召回改成并行执行用CompletableFuture或线程池并发调用总耗时取决于最慢的一路。Faiss 索引在服务启动时全量加载到内存用IndexFlatIP而不是IndexIVFFlat后者虽然省内存但查询慢。4.5 现象新用户冷启动推荐全是随机新闻点击率极低原因新用户没有行为数据协同过滤召回为空内容召回又没有用户画像只能随机推。解决新用户走「地域 时间 热度」兜底策略根据用户 IP 定位城市推当地新闻根据当前时间段推对应类别早上推财经、中午推体育、晚上推娱乐再混入全局热度 Top 20 的新闻。等用户产生 5 次以上点击后逐步切换到个性化召回。5. 推荐效果验证与一个提升 CTR 的排序技巧5.1 离线评估AUC 和 RecallK 怎么看模型训练完后不要直接上线先用离线指标验证。把最后一天的行为数据按 8:2 切分成训练集和测试集在测试集上算两个指标指标含义合格线我的实测值AUC排序模型区分正负样本的能力 0.700.74Recall50前 50 个推荐里命中用户实际点击的比例 0.350.41Coverage推荐结果覆盖的新闻占全量新闻的比例 0.600.67AUC 低于 0.70 说明特征区分度不够检查用户实时兴趣特征是否正常写入。Recall50 低于 0.35 说明召回层漏了太多考虑增加一路召回或放宽召回数量。Coverage 低于 0.60 说明推荐结果太集中检查热门新闻的权重是否过高。5.2 在线 A/B 实验流量切分与指标观测离线指标达标后用 A/B 实验验证在线效果。把用户按 user_id 哈希取模分成两组对照组走旧推荐逻辑实验组走新系统。实验组流量从 5% 开始观察三天点击率CTR实验组应比对照组高 15% 以上人均停留时长实验组应比对照组高 10% 以上滑动跳过率实验组应比对照组低 10% 以上如果 CTR 涨了但停留时长没涨说明推荐结果标题党多、内容质量差需要在排序模型里加内容质量分特征。如果滑动跳过率没降说明推荐多样性不够在排序层加一个类别打散逻辑。5.3 一个提升 CTR 的排序技巧实时兴趣衰减加权最后分享一个我实测有效的排序技巧。用户兴趣是有衰减的10 分钟前点击的新闻类别权重应该比 2 小时前点击的高。在排序模型的特征里加一个「时间衰减因子」# 实时兴趣特征加权 import math def time_decay_weight(click_time, current_time, half_life1800): click_time: 用户点击新闻的时间戳秒 current_time: 当前推荐请求时间戳秒 half_life: 半衰期默认 30 分钟 delta current_time - click_time return math.exp(-delta * math.log(2) / half_life)这个因子乘在用户实时兴趣向量的每个维度上让最近点击的类别权重更高。我上线这个技巧后CTR 从 4.1% 涨到 4.7%涨幅约 15%。half_life设 1800 秒是新闻场景的经验值电商场景可以设 86400 秒一天因为商品兴趣衰减慢。提示half_life不要设得太小否则推荐结果会过度集中在用户刚点过的类别上多样性下降。我试过 600 秒CTR 虽然涨了但停留时长掉了得不偿失。这套系统我从零搭到上线跑了三个月最大的教训是不要一上来就追求模型多复杂先把 Kafka 到 Flink 到 Redis 的实时链路跑通保证数据不丢不重再逐步加召回和排序。很多翻车不是因为算法不行而是因为实时特征根本没写进去模型拿到的全是过期数据。希望帮到你。本文还有配套的精品资源点击获取
返回列表