1. 推荐链路设计与选题拆解1.1 为什么是“Spark 新闻推荐”这个组合近几年高校计算机毕业设计里基于Spark的新闻推荐系统属于出镜率很高的题目。它看起来像个“经典款”但实际操作中大部分同学都在同一个地方卡住——不知道怎么把 Spark 这个分布式计算引擎和“给用户推新闻”这件事自然串起来。有人把 Spark 当数据库用有人把推荐逻辑写在 Web 后端、Spark 只做个数据导入结果答辩一问“Spark 在系统里到底承担了什么”就答不上来。我的看法是这个题目真正的价值不在推荐算法本身而在于让你把一条完整的数据处理链路走通从原始日志和新闻语料到特征加工再到离线批量计算推荐结果最后通过一个轻量级 Web 应用把结果露出来。Spark 在中间承担的是“离线计算枢纽”的角色。新闻推荐相比电商推荐、视频推荐有一个很大的特点物品新闻的生命周期极短。一条新闻的热度窗口可能就是几小时到一两天用户更看重“新鲜”和“相关”而不是“历史偏好”。这就决定了新闻推荐系统更适合做基于内容的相似度推荐——你今天看了一篇关于芯片的深度报道系统就该马上给你推其他芯片相关新闻。相比之下基于用户行为的协同过滤在新闻场景下会遇到严重的冷启动问题新用户没有行为记录新新闻没有交互记录。把业务逻辑和技术选型结合起来看这套毕设的合理定位是用 Spark 完成全部离线计算和推荐生成用 Web 后端只做结果展示。数据量不需要多大但链路必须是完整的。1.2 核心需求拆解系统要做什么任何毕设课题落到实现层面都要先拆成几个可以独立验收的功能模块。这个系统我建议拆成四块第一新闻语料处理模块。你得有足够多的新闻数据至少几百条才能让推荐效果看起来“像那么回事”。数据可以从公开数据集下载也可以通过爬虫获取。处理内容包括清洗缺失字段、过滤重复新闻、统一时间格式。第二用户行为模拟与特征加工模块。新闻推荐系统没有真实用户怎么办自己构造。设计一张用户行为表模拟不同用户对新闻的浏览、点击、收藏行为。这张表的意义在于保证 Web 端有“个性化”的输入也让答辩时有话可说——这个模块对应的是真实系统中的埋点日志。第三基于 Spark 的推荐计算模块。这是核心中的核心。读取新闻特征向量通过 Spark SQL 和 Spark MLlib 计算新闻之间的余弦相似度选出每篇新闻的 Top-N 相似文章再结合用户近期浏览记录把用户没看过的相似新闻推荐出来最后用热榜新闻填充以保证推荐列表够长。第四Web 展示模块。用 Flask 搭一个极简前端显示新闻详情和对应的推荐列表。推荐结果可以预计算好存在 Redis 或 MySQL 里Web 端只负责查表。这是最稳妥的做法避免后端现场跑 Spark 作业——提交一次任务几十秒演示容易卡住。下面这张表把这些模块对应的技术选型和产出物整理清楚。功能模块技术载体核心产出语料处理Spark DataFrame SQL清洗后的新闻表行为数据Python 脚本生成用户-浏览记录表特征计算Spark MLlib分词 TF-IDF新闻特征向量表相似度计算Spark 分布式计算新闻相似度 TopN 表推荐结果汇总Spark SQL 关联每位用户的推荐列表Web 展示Flask Bootstrap可点击浏览的推荐页面这套拆法还有一个额外好处论文的“系统设计”章节可以直接按这个思路展开功能模块图和数据流向图画起来非常自然。2. 技术选型逻辑与数据设计2.1 Spark 版本与运行模式怎么选选型是最容易纠结的环节。我的建议是除非实验室有现成的集群否则一律从 Spark 3.x 单机Local模式起步。不推荐一开始就搭三台虚拟机的集群——时间成本高网络配置容易出现玄学问题而且对毕设演示来说没有任何额外收益。如果你选用 Spark 3.3 或 3.4对应的 PySpark 版本也跟着走。这里有一个我的个人偏好使用 Python 版 PySpark而不是 Scala。理由是PySpark 的 DataFrame API 和 SQL 操作几乎和 Scala 版一样但 Python 处理数据结构更灵活写 UDF 也快而且答辩现场如果有需要可以直接在 Jupyter Notebook 里演示中间结果这种东西对答辩很有说服力。Local 模式跑 Spark本质上是 Spark 在本地用多线程模拟分布式执行。你不用配置 Hadoop 集群只需要确保 Java JDK 版本匹配一般 Spark 3.x 对应 Java 8 或 11。数据量在万条级别下Local 模式的执行速度几乎是秒级足够应付毕设的所有需求。配置参数时注意 JVM 内存设置。启动 PySpark 或提交 Spark 作业时在提交命令里显式加上 driver 内存参数避免因默认内存过小导致任务失败。后面我在“踩坑记录”里会专门展开这一类问题。2.2 数据集结构新闻表和用户行为表数据是整个项目的地基。我基于公开的新闻数据集进行测试但为了演示稳定性也自己构造了一批结构规整的数据。通用的设计如下新闻表news存的是新闻正文语料。核心字段news_id -- 新闻唯一编号 title -- 新闻标题 content -- 新闻正文内容 category -- 栏目分类比如 科技、财经、体育、娱乐 publish_time -- 发布时间统一为 yyyy-MM-dd HH:mm:ss用户行为表user_behavior模拟用户与新闻的交互。核心字段behavior_id -- 行为编号 user_id -- 用户编号 news_id -- 浏览的新闻编号 behavior_type-- 行为类型click 点击 / collect 收藏 / share 分享 behavior_time-- 行为发生时间这两张表的关系很清晰行为表通过 news_id 关联新闻表拿到新闻内容是后面所有计算的入口。行为数据我建议用 Python 脚本生成逻辑可以稍微“带点规律”比如用户 A 在科技类新闻上看得多用户 B 在体育类新闻上看得多。这样生成的行为数据用于测试时推荐结果会更符合直觉演示效果也更真实。2.3 算法路线内容相似度为主、热榜兜底给这个系统选推荐算法的时候我先筛掉了几条路。不使用纯协同过滤。传统 ItemCF 或 ALS 需要稠密的“用户-物品”评分矩阵新闻场景下用户行为天然稀疏且新闻更新换代快新建的新闻没有任何交互记录协同过滤没法对这些新内容做推荐。如果面试官或评审老师问“新新闻怎么推出来”你答不上来就很被动。使用基于内容的新闻相似度推荐作为主算法。思路三步走对新闻标题和正文做分词、去停用词用 TF-IDF 把每条新闻加工成特征向量计算新闻之间的余弦相似度取 Top-N。基于 Spark 实现这个流程Spark MLlib 里现成的组件完全够用分词可以用简单的 HanLP 或者 jieba 分词通过 UDF 调用向量化用 HashingTF 和 IDF相似度计算直接用 DataFrame 的自连接操作完成。全部代码控制在 100 行以内但对 Spark 的核心机制——RDD 到 DataFrame 的转换、Transformer 和 Estimator 的管道思想、分布式 shuffle 的代价——都能够体现出来。在实际的推荐策略里主算法的结果需要做两部分补充。第一部分是热门新闻作为冷启动补充——给新用户直接推全站点击量最高的新闻这是最基础的策略。第二部分是栏目托底——如果某个用户的行为数据太少相似推荐结果不足就用对应栏目下的热门新闻填补空缺。整个推荐列表的构成相似新闻 80% 热门新闻 20%比例可以后续调整。这套方案还有一个附带好处它就是真实互联网领域经典的“兴趣 Feed”和“热门榜单”双通道策略的缩略版后续论文的创新点可以写在“双通道融合”上。3. 核心实现与关键代码实战3.1 环境准备一个能跑通 Spark 的最小环境适合毕设的部署方式单机 Linux 或 macOS安装 Java JDK 8/11、Python 3.8、Spark 3.x。Windows 也能跑但 Shell 命令和 PATH 配置会更麻烦我自己不推荐在 Windows 上做。验证环境是否可用最稳妥的“冒烟测试”是跑一个最简单的 WordCountfrom pyspark.sql import SparkSession spark SparkSession.builder \ .appName(WordCountTest) \ .master(local[4]) \ .getOrCreate() sc spark.sparkContext lines sc.parallelize([hello world, hello spark]) counts lines.flatMap(lambda line: line.split( )) \ .map(lambda word: (word, 1)) \ .reduceByKey(lambda a, b: a b) print(counts.collect())如果这个程序能正常输出结果说明 Spark 环境基本没问题。这个测试同时也是答辩时解释“Spark 基础编程模型”的最好例子——flatMap、map、reduceByKey 这几个算子把 MapReduce 思想演示得清清楚楚。注意local[4]表示用 4 个线程模拟并行执行。如果你电脑核心数少可以用local[2]。跑假的“集群并行”没有实际意义重点是让程序顺利跑通。3.2 新闻数据加载与预处理数据准备完毕之后用 Spark 加载并清洗数据。这一步是数据仓库里最常见的 ETL 操作也是你展示 Spark SQL 能力最好的地方。代码大致如下from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, length, to_timestamp spark SparkSession.builder \ .appName(NewsRecSys) \ .master(local[4]) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate() news_df spark.read.json(data/news.json) # 基础清洗去掉没有正文内容的数据去除重复数据 clean_df news_df.filter(col(content).isNotNull()) \ .filter(length(col(content)) 30) \ .dropDuplicates([news_id]) \ .withColumn(publish_time, to_timestamp(col(publish_time), yyyy-MM-dd HH:mm:ss))可以解释一下各个步骤的作用第一步筛选掉空内容第二步过滤掉太短的垃圾样本第三步按新闻 ID 去重第四步统一时间格式方便后续做时间衰减。这些是 ETL 的常规操作缺了哪一步都不够完整。这里有一个细节加载的数据如果比较大建议先做成 Parquet 列式存储格式再存一份副本后续生成特征时读取会更快clean_df.write.mode(overwrite).parquet(data/news_clean.parquet)即使在万条规模下Parquet 的压缩和列裁剪优势也不可忽视而且写 Parquet 意味着你在实践中应用了“数据分层”的思路这在答辩时也算一个小亮点。3.3 文本特征加工分词、TF-IDF、向量化特征加工是整个推荐系统的核心环节。我采用的流程是先对新闻的“标题 正文”拼接后的文本做中文分词再去停用词然后通过 TF-IDF 转化为向量。中文分词这里需要额外引入工具。使用 jieba 库作为分词器通过 PySpark UDF 在分布式中调用。需要注意UDF 中的第三方依赖需要在每个 executor 上都能找到。在 Local 模式下其实就是本机环境所以 pip 安装好 jieba 即可但如果是集群模式就必须在启动脚本里用--py-files参数把依赖库打包上传。这是很多同学初次用 Spark 时最容易忽略的底层问题。分词与向量化代码import jieba from pyspark.sql.functions import udf from pyspark.sql.types import ArrayType, StringType from pyspark.ml.feature import HashingTF, IDF, CountVectorizer def seg_cn(text): return [w for w in jieba.cut(text) if w.strip()] seg_udf udf(seg_cn, ArrayType(StringType())) df clean_df.withColumn(words, seg_udf( concat(col(title), lit(。), col(content)) )) # TF-IDF特征向量化 hashingTF HashingTF(inputColwords, outputColrawFeatures, numFeatures20000) featurizedData hashingTF.transform(df) idf IDF(inputColrawFeatures, outputColfeatures) idfModel idf.fit(featurizedData) rescaledData idfModel.transform(featurizedData)关于停止词的处理有两个方案。方案一是在 jieba 分词后过滤停用词表方案二是直接用 HashingTF 的降维能力把噪声词“冲淡”。我实际测试下来加停用词会显著改善相似度质量尤其是降低“了”“的”“在”这类词的干扰。做法是在 UDF 内部做一个全局停用词集合过滤不占用额外 DataFrame 操作stopwords set(...) # 加载停用词表 def seg_cn(text): return [w for w in jieba.cut(text) if w.strip() and w not in stopwords]特征维度numFeatures的选择影响两个维度一是内存占用二是计算精度。维度过低比如 1000会增加哈希碰撞导致不同新闻的高频词被映射到同一个位置相似度虚高维度过高则浪费资源。我实测 20000 维在一个几千条新闻的数据集上效果稳定推荐直接使用这个值。3.4 相似度计算与推荐生成特征向量准备好之后推荐核心计算就简单了计算两两新闻之间的余弦相似度垂直筛选出 TopN。Spark 里 DataFrame 自连接后再计算相似度的思路很直接但要注意自连接会产生 shuffle数据量大时会非常慢。万条数据规模下这个操作可以接受但如果上了几十万条新闻就必须改成用 BucketedRandomProjectionLSH 或聚类方案先粗筛再做精排。基础版代码from pyspark.ml.linalg import Vectors, DenseVector from pyspark.sql.functions import udf, col def cos_sim(v1, v2): dot float(v1.dot(v2)) norm float(v1.norm(2) * v2.norm(2)) return 0.0 if norm 0 else dot / norm cos_udf udf(cos_sim, DoubleType()) # 自连接计算相似度 pair_df rescaledData.select( col(news_id).alias(news_id_a), col(title).alias(title_a), col(features).alias(features_a) ).crossJoin( rescaledData.select( col(news_id).alias(news_id_b), col(title).alias(title_b), col(features).alias(features_b) ) ).filter(col(news_id_a) col(news_id_b)) sim_df pair_df.withColumn(sim_score, cos_udf(col(features_a), col(features_b))) \ .filter(col(sim_score) 0.2)跨连接和过滤必须同现。news_id_a news_id_b这行代码的目的有两个一是去掉新闻自己和自己算相似度的重复项二是让每一对新闻只保留一种排序相似度结果可以直接复用。真正的生产环境不会用纯 crossJoin因为它很昂贵。论文里如果只提到这个评审可能会质疑。为了避免这种尴尬在你的论文里加一句说明此方案在十万级数据量以内可以直接使用若数据规模超量可以改用近似近邻方法比如 Spark MLlib 的 BucketedRandomProjectionLSH在保持 MapReduce 计算思路的同时大幅减少计算量。拿到相似新闻对之后还需要做一步汇总生成最终每位用户的推荐列表。这一步的逻辑是读取该用户的最近浏览新闻列表找到这些新闻各自的 Top-8 相似新闻去掉用户已经看过的剩余结果按相似度得分排序不足部分用热门新闻补到 20 条。这个推荐列表可以预计算好写回数据库。进一步优化可以做成定时任务每天凌晨跑一次生成当天的推荐快照。这样系统的“离线计算”和“在线服务”就彻底分开了。4. 系统演示与 Web 端接入4.1 展示框架选择与计算结果存储Web 端不需要做得多复杂但必须让评委在浏览器里能直观看到推荐效果。我用 Flask 搭了一个极简应用前端页面只包含三个视图新闻列表页、新闻详情页、基于当前用户行为数据的“我的推荐”页。推荐结果的存储方案我实测推荐 Redis。理由很现实预计算好的推荐结果本质上是“以用户 ID 为 Key以新闻 ID 列表为 Value”的 KV 数据Redis 的读写速度快数据结构天然匹配。而且 Redis 可以作为独立的中间件写进系统架构图里让整个系统看起来更像真实工业级项目。存储格式建议# 用户 10086 的推荐列表按推荐顺序存储 r.rpush(rec:user:10086, n10001, n10013, n10042, ...)如果嫌 Redis 配置麻烦MySQL 里的推荐表也能替代。但 Redis 在整个技术栈中代表的是“缓存层”和“高性能访问”有这个模块系统分层更完善。4.2 Flask 接入与前端效果演示Flask 里读取 Redis 推荐列表并展示新闻这部分代码非常轻量from flask import Flask, render_template, request import redis app Flask(__name__) r redis.Redis(hostlocalhost, port6379, db0) app.route(/) def index(): news_list load_hot_news() # 热门列表兜底推荐 return render_template(index.html, news_listnews_list) app.route(/user/uid) def recommend(uid): news_ids r.lrange(frec:user:{uid}, 0, 19) news_list [load_news(news_id) for news_id in news_ids] return render_template(recommend.html, news_listnews_list)使用 Bootstrap 做样式基础首页展示新闻卡片详情页展示新闻正文和右侧的“你可能感兴趣”列表。演示时准备几个行为模式差异明显的测试用户比如用户 A 频繁点科技类新闻用户 B 只看体育类点击“我的推荐”页就能立刻看到结果差异效果很直观。这里要特别提醒一个演示陷阱现场展示时千万不要重新运行推荐计算。提前把推荐结果存好演示时只展示 Web 端页面思路是环境已经就绪结果已经预计算好Web 现场只需要读数据。真正跑 Spark 计算留到你录视频或自己调试时再做。否则现场环境稍有波动一个 Stack 报错就可能导致整个答辩效果大幅滑坡。5. 常见问题与实战踩坑记录5.1 PySpark 环境相关的“玄学”问题如果你用的是 PySpark环境问题会占据调试时间的大头。以下几个问题几乎人人都会遇到。第一个是Java 版本不匹配。Spark 3.3 和 3.4 对 Java 的版本要求比较严格有的机器装了 Java 17直接报UnsupportedClassVersionError。解决办法是装一个 Java 8 或 Java 11并在启动脚本里手动指定JAVA_HOMEexport JAVA_HOME/usr/lib/jvm/java-11-openjdk-amd64第二个是Python 解释器版本冲突。有些系统里默认 Python 是 2.7CentOS 的老问题直接把 PySpark 运行环境搅得天翻地覆。我一般习惯用 conda 建立独立环境conda create -n spark-env python3.8 conda activate spark-env pip install pyspark3.3.0 jieba flask redis第三个是Spark 作业启动后莫名失败接着看日志显示 GC overhead limit exceeded。这大多是内存不足导致的。在 SparkSession 里加上.config(spark.driver.memory, 4g) .config(spark.executor.memory, 4g)对于本地模式Driver 和 Executor 共用进程这里加上 4GB 通常就能扛住。5.2 中文分词和文本向量的特殊坑分词和向量化这块一个最容易被忽视的坑是没有过滤停用词导致所有新闻的相似度都偏高。你要知道如果“的”“了”“在”这类词没有被过滤它们的词频在所有文本里都很高会拉高任意两篇新闻的余弦相似度结果就是你推出来的“相似文章”看起来完全不搭界。推荐效果好不好停用词表质量占了很大因素。第二个坑是jieba 分词在 UDF 中重复加载词典造成的性能瓶颈。如果你在 UDF 里每次调用都执行jieba.initialize()上万篇新闻的计算时间会明显拉长。正确的做法在 UDF 外部、spark 启动初期完成分词词典加载使用全局变量jieba在 UDF 内部直接调用jieba.cut。本地模式下多个线程共用一个词典速度会快不少。注意集群模式如果是这种方式每台 worker 会单独初始化一次但那也还可以接受。import jieba jieba.initialize() def seg_cn(text): return [w for w in jieba.cut(text) if w.strip()]第三个坑和HashingTF 的维度选择有关。如果设的numFeatures太小如 500分词后的特征碰撞会很严重两篇不相关新闻因为碰撞到同一个哈希桶导致相似度虚高如果设得太大如 100 万内存浪费明显。我在实践中的建议1 万到 5 万之间万条数据规模选 2 万正合适这也是多数 Spark 文本处理实践里的常用配置。5.3 演示层面如何规避“翻车”有时候你本地一切正常一到演示现场就出问题。这类问题大多不是代码问题而是架构问题。我的几条血泪建议第一个推荐结果在演示前全部预先生成完毕数据直接存在 Redis 或 MySQL 中。现场展示时任何逻辑都必须做到秒级响应只要涉及 Spark 在线计算就一定会卡顿或报错。第二个做一份兜底的“静态演示环境”。比如录好一段系统操作视频准备一些推荐效果截图。真到现场环境无法恢复时这部分材料能马上顶上。第三个准备一份“Spark 执行计划讲解”。答辩高频问题之一是“你的推荐算法为什么准”。从 Spark 的角度你可以给出计算逻辑链IDF 特征生成 - 余弦相似度计算 - 按用户行为过滤 - 热门托底。如果还能结合数据量解释为什么本地模式够用以及生产环境中切换到集群模式需要改哪些参数这一个问题就能答得非常完整。5.4 数据量太小怎么办一个很容易被答辩老师抓住的问题就是系统里只存了 50 条新闻。推荐效果乍一看还行但老师可能反复质疑“数据量这么小有什么意义”。我的解决办法分两步。第一步公开中文新闻数据集比较多建议找一份上万条规模的或者直接爬取开源新闻网站的 RSS 数据不必全部展示但数据准备脚本要能证明“可以处理万条规模的数据”。第二步把“数据规模扩展性”写进论文的实验章节与小数据集对比说明在更大规模上 Spark 的计算优势才会真正体现然后在集群环境下用合成数据做一次万条甚至十万条规模的压测跑通一套打分和耗时记录。这套动作做完项目从“演示 Demo”升级为“完整研究过程”说服力完全不同。6. 完成度提升与课题扩展方向6.1 让自己多“一点点”竞争优势很多毕设做到上面这个程度已经可以收尾但如果你想让项目有更好的展示效果可以再加两件事。一个是给清洗和特征加工加上流程可视化。在 Jupyter Notebook 里记录 Spark DataFrame 的 Schema、分区数、数据量变化做成插图放进论文或答辩 PPT 里。这能正面回应老师最关心的“你怎么证明 Spark 干了活”。另一个是加实时更新机制。在真实新闻推荐场景中新闻的实时性很重要。你可以用 Spark Structured Streaming 模拟读取 Kafka 中的新新闻流在新闻入库后自动更新相似度结果。不需要真的搭 Kafka 集群用一万条模拟数据按时间戳回放或者直接用本地 Socket 源源不断发送 JSON 字符串就能实现。加了这一步系统的技术复杂度就从“离线批处理”提升到“离线批处理 流处理”这在毕设里属于妥妥的加分项也是 Spark 生态完整度的体现。6.2 这套代码还能迁移到哪些场景等到你把这套系统完整实现一遍会发现最宝贵的收获不是那几段代码而是这套“数据导入 - 清洗 - 特征工程 - 离线计算 - 在线服务”的套路完全可以迁移到其他项目里。比如电商商品推荐只需要把新闻表换成商品表、把新闻文本特征换成商品标题 类目 描述即可比如短视频推荐把“新闻”换成“视频标题 标签 简介”特征提取逻辑换汤不换药。甚至农产品价格分析、网约车订单数据分析这类方向也只是把处理对象换成结构化数据核心的 Spark SQL 清洗和特征加工流程是一样的。底层的 Spark 数据处理思维一旦建立后续做任何大数据方向的项目都会顺很多。我个人在实际操作中的体会是这个题目最怕的不是代码难写而是思路散——既要懂 Spark又要懂推荐系统还要会 Web 展示课题看着简单真正把三条线串起来还是费了一番功夫。如果你正在为这个题目头疼不妨从“推荐结果从哪来、怎样算、往哪存、怎么展示”这条链路入手先打通主流程再做细化。只要主流程稳定能跑后续所有优化都是锦上添花。最后一个建议保存好每一版中间结果。从原始 JSON 到清洗后 Parquet再到相似度 CSV都留一份。这不仅是为了调试方便更是因为你论文里要贴大量“处理效果对比图”这些中间文件就是最好的素材来源。
