恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

Spark+Scala+MongoDB商品推荐系统实战:ALS离线召回与在线检索

  • 首页
  • 资讯中心
  • /
  • Spark+Scala+MongoDB商品推荐系统实战:ALS离线召回与在线检索

相关资讯

Headlamp IngressClass KubeObject API 深度解析:类结构、apiEndpoint 端点与前端接入方式 2026/9/16 10:42:37
H3C S-MLAG与服务器bond4双活接入配置实战 2026/9/16 10:42:37
MiniCPM5-2B端侧智能体实战:高智能密度Agent运行时部署指南 2026/9/16 10:37:36

最新资讯

GPT-5与可控智能体:AI落地的关键技术突破
把 Manus 式 Agent 的模型通道指到 TaoToken,任务编排才跑得动
LeetCode-Book 精讲:最长递增子序列(LIS)——从 O(N²) 动态规划到 O(NlogN) 二分优化
2026年AI论文平台核心技术解析与应用指南
AWS CLI 实战:使用 `aws codepipeline list-pipeline-executions` 查看 CodePipeline 流水线执行历史
Flutter列表性能优化与鸿蒙适配实战

今日推荐

IoT-For-Beginners 智能语音计时器:Wio Terminal 基于 DMAC 与 Flash 的音频采集实战
基于MATLAB的CRI显色指数计算:从SPD光谱到Ra的完整流程
JSP+Servlet+MySQL博客系统源码部署与优化全攻略

本周热门

AI SDK Harness 依赖更新指南:掌握 harness 包 SDK 依赖的升级、桥接同步与一致性校验
Refine v5 Ant Design NumberField 组件实战:基于 Intl 的本地化数字格式化
Flutter应用改名全指南:从Android到iOS的配置与工具实践

本月精选

自研推理加速器Redwood:两周内实现PyTorch模型高效部署的实战教程
V4L2摄像头采集实战:从camera_client.rar到出图全流程解析
从“谁发明了钢琴键”到知识问答智能体:RAG与记忆工程实践

Spark+Scala+MongoDB商品推荐系统实战:ALS离线召回与在线检索

发布时间:2026/9/16 10:42:37
Spark+Scala+MongoDB商品推荐系统实战:ALS离线召回与在线检索 简介基于SparkScalaMongoDB的大数据实战资源聚焦商品推荐系统设计与实现适合毕业设计、课程设计及推荐算法入门者。资源覆盖数据收集、预处理、特征工程、模型训练、评估与实时推荐等完整链路结合Spark MLlib中的协同过滤方法可快速搭建个性化推荐流程。MongoDB用于存储用户行为与商品信息Spark负责清洗整合与分布式计算Scala代码简洁易维护。压缩包共146个文件以xml配置、java与scala源码、properties配置、class编译文件为主附带csv数据、js前端及MongoDB脚本整体约7.92MB。已有130人学习适合具备一定Scala基础、希望掌握Spark与MongoDB整合应用的开发者。通过源码可学习从数据读取、预处理、特征提取到模型训练与结果输出的完整实现同时了解工程目录结构和项目排错思路是兼顾理论与实践的高性价比参考。1. 别把推荐系统做成离线统计这个项目真正要过的三道关把 Spark、Scala、MongoDB 三个词放进一个标题里看上去是一个典型的“大数据毕设套餐”但真按这个组合动手做商品推荐系统时绝大多数人卡住的不是算法而是数据从哪来、算完往哪写、线上怎么取。这个项目标题真正要回答的问题不是“ALS 怎么训练”而是用户行为日志怎么变成评分矩阵Spark 算完的召回结果怎么落到 MongoDB 里让应用层秒查以及这套离线管线在增量数据面前怎么不重建全量模型。适合谁适合已经写过 MapReduce 或 DataFrame 基础操作、想完整走通一套“离线计算 在线检索”架构的人。新手能照着搭出最小可运行版本熟手则能从中看到数据倾斜、内存溢出、MongoDB 索引设计这些真正影响产线的细节。先把结论放在前面这个系统里 Spark 只负责算MongoDB 才决定线上体验而 Scala 决定你加班到几点。2. 明确组件边界Spark、Scala、MongoDB 在推荐系统里各自解决什么问题2.1 为什么是 Spark 而不是 Flink离线批处理才是商品推荐的常态商品推荐系统最常见的落地形态不是实时计算而是离线召回 在线检索。用户行为数据攒一天凌晨跑一轮 Spark 任务算出每个用户的 TopN 商品列表写入 MongoDB白天应用层直接查表返回。这个链路里 Spark 的核心价值是分布式内存计算它能在一个集群上把千万级用户 × 百万级商品的评分矩阵拆成 RDD/DataFrame 分片并行训练。而 Flink 虽然擅长实时流处理但基于 ALS 矩阵分解这类迭代式算法离线批处理的成熟度和生态完整性远高于流计算。从选型角度说Spark 的 MLlib 自带 ALS 推荐算法实现DataSet API 天然适配 Scala 的类型系统。相比之下 Python 版 PySpark 虽然也能跑但在算子调试、类型推断和 JVM 内存调优上会隔一层出问题不好定位。2.2 Scala 的角色不是语法选择是类型安全与性能边界的取舍很多初做这个项目的人会问明明可以用 PySpark为什么标题强调 Scala这里有个真实原因ALS 训练循环里有大量对象序列化和 RDD 算子闭包传递Scala 的编译期类型检查能提前暴露字段名拼写错误、类型不匹配问题而 Python 要到运行时才报错。更关键的是Spark 原生算子 API 是为 Scala 设计的DataSet[T] 的强类型编码器在 Scala 下直接映射 JVM 对象省掉 PySpark 的 Python-JVM 通信开销。但这不意味着 Scala 没有代价。Scala 版本和 Spark 版本有严格对应关系比如 Spark 3.x 需要 Scala 2.12/2.13用错版本直接报 NoSuchMethodError。而且 Scala 的隐式转换和集合 API 对新手来说有一段陡峭学习曲线——不过在这个项目里你只需要掌握 case class、map/flatMap、groupBy 这几个操作就能覆盖 80% 的编码量。2.3 MongoDB 不是简单缓存它承担了推荐结果的存储、检索和更新三重职责推荐结果的写入端选择 MongoDB 是合理且常见的做法。相比 MySQLMongoDB 的文档模型能直接存嵌套数组结构——每个用户的推荐列表天然是一个 JSON 数组不需要拆成多张关联表再做 JOIN。相比 RedisMongoDB 能持久化全量结果且支持复杂查询适合推荐列表这种“写入一次、读取多次、定期重算”的数据形态。在这个项目里MongoDB 里至少需要建三张核心集合集合名存储内容关键索引典型查询user_product_score用户-商品评分稀疏矩阵{ user_id: 1, product_id: 1 } 复合唯一索引按用户查评分向量als_recall_resultALS 输出的 TopN 推荐结果{ user_id: 1 } 单字段索引按用户取推荐列表product_feature商品特征向量可选用于相似推荐{ product_id: 1 }按商品查向量有人会质疑为什么不把推荐结果放 Elasticsearch理由是这个项目的查询模式极度简单——给定 user_id 取数组MongoDB 的 _id 哈希索引或 user_id 普通索引足够支撑几千 QPS引入 ES 徒增运维成本。还有个容易忽略的点MongoDB 的聚合框架可以在返回推荐列表时直接做字段裁剪把商品价格、库存状态联表查出来省掉应用层二次回表。3. 从日志到召回结果基于 SparkScala 的 ALS 推荐管线完整实现3.1 数据抽象用 case class 定义输入输出的三层数据模型动手写 Spark 任务前先定义数据的物理形态。输入是用户行为日志通常包含 userId、productId、behaviorType浏览、收藏、加购、下单、timestamp 四个字段。这里行业内的通用做法是行为转评分——浏览 1 分、收藏 3 分、加购 5 分、下单 10 分时间衰减因子可选。转换逻辑写在一个 Scala object 里方便测试和复用// 输入日志的 case class字段名和 schema 严格对应 case class UserBehavior( userId: String, productId: String, behaviorType: String, // browse, favorite, cart, order timestamp: Long ) // 转评分后的数据集 case class UserProductRating( userId: String, productId: String, rating: Double ) // ALS 训练数据需要数值型 ID维护用户/商品到索引的映射 case class Rating( userIndex: Int, productIndex: Int, rating: Double )这段代码的核心设计意图有三点。第一case class 编译期校验字段名日志格式调整时 IDE 直接报错而不是等 Spark 任务跑一半报 NullPointerException。第二Rating 里显式用 Int 类型而非 String因为 ALS 算法要求输入必须是数值型 ID后续用 StringIndexer 转换会多一步序列化开销。第三时间戳保留在原始类型里方便后面做时间衰减或按天切片。3.2 数据装载与预处理处理数据倾斜和无效行为真实日志里必然有“少数用户贡献大量行为”和高频无意义浏览比如爬虫连续翻页直接喂给 ALS 会让热用户的主导评分淹没长尾商品。常见做法是加两层过滤行为总数小于 5 的用户直接丢弃数据太少训练无意义单个用户的点击序列超过 500 条时采样截断。// 过滤低频用户限制高频用户行为数 val validBehavior behaviorDF .groupBy(userId) .count() .filter(count 5) .join(behaviorDF, userId) .drop(count) // 按用户分组后每组最多取 500 条行为减少倾斜 val sampledBehavior validBehavior .rdd .map(row (row.getAs[String](userId), row)) .groupByKey() .flatMap { case (_, iter) iter.toSeq.sortBy(_.getAs[Long](timestamp)).reverse.take(500) }这里的逻辑不只是数据清洗第一次 groupBy 过滤掉冷启动噪声用户第二次按时间倒序截断保证保留的是“最新”的 500 个行为。注意采样必须在时间维度上倒序截断而不是随机采样否则会破坏用户兴趣随时间的连续性。过滤后的数据通过 rating 映射函数转换成 (userId, productId, rating) 三元组交给 ALS 前还需要用 StringIndexer 把字符串 ID 转成连续索引——这一步做在 DataFrame 层比 RDD 层更稳妥因为 DataFrame 的 Catalyst 优化器会处理常量折叠。3.3 训练 ALS 模型参数怎么设才能避免模型跑飞ALS交替最小二乘是 Spark MLlib 里最成熟的协同过滤实现核心思想是交替固定用户矩阵和物品矩阵求解最小化平方误差。训练代码本身不长真正的挑战在参数设置import org.apache.spark.ml.recommendation.ALS val als new ALS() .setRank(20) // 隐向量维度影响模型表达力 .setMaxIter(15) // 迭代次数过大容易过拟合 .setRegParam(0.1) // 正则化系数控制防止过拟合 .setAlpha(1.0) // 隐式反馈置信度参数 .setUserCol(userIndex) .setItemCol(productIndex) .setRatingCol(rating) .setColdStartStrategy(drop) // 预测时遇到新用户/商品直接丢弃避免 NaN val model als.fit(trainingData)每个参数都有实际影响逐个说清楚参数常见取值范围设置逻辑错误表现rank10-50商品品类数多就调大否则欠拟合推荐结果集中在头部商品maxIter10-20超过 20 收益递减训练时间翻倍损失函数下降曲线变平regParam0.01-0.1数据稀疏时调小稠密时调大验证集 AUC 波动大alpha0.5-2.0隐式反馈场景下控制置信度冷门商品完全不被召回coldStartStrategydrop/none必须设 drop否则 predict 返回 NaN预测结果出现 null训练过程中需要盯住的指标不是 loss 而是 topN 命中率——ALS 的 loss 下降不代表推荐质量好。业内会取一部分用户做时间切分验证前 7 天训练后 3 天验证用 RecallK 和 PrecisionK 评估。3.4 写入 MongoDB批量 upsert 而不是逐条 insert推荐结果写入 MongoDB 是整个环节里最容易被低估性能瓶颈的地方。初学者最容易犯的错误是用了 for 循环逐条 insert5 万用户的推荐列表要插入半小时。正确做法是批量 upsert用 bulkWrite 一次性提交import org.bson.Document import com.mongodb.MongoClient import com.mongodb.client.model.UpdateOneModel import com.mongodb.client.model.UpdateOptions val mongoClient new MongoClient(localhost, 27017) val db mongoClient.getDatabase(recommend) val coll db.getCollection(als_recall_result) // 每个用户生成一个文档字段是 userId 商品数组 val docs model.recommendForAllUsers(20) .collect() .map { row val userId row.getAs[Int](userIndex) val recommendations row.getAs[Seq[Row]](recommendations) .map(item item.getAs[Int](productIndex)) new Document(user_id, userId) .append(product_list, recommendations) } // bulkWrite 批量更新按 user_id 匹配则覆盖不存在则插入 val writes docs.map { doc new UpdateOneModel[Document]( new Document(user_id, doc.get(user_id)), new Document($set, doc), new UpdateOptions().upsert(true) ) } coll.bulkWrite(java.util.Arrays.asList(writes: _*))bulkWrite 的关键在于 upsert(true)——推荐结果是整体重算而不是增量追加用 upsert 保证同一用户重复计算时只覆盖不堆积。还有个细节是 recommendForAllUsers 返回的 recommendations 是一个数组字段MongoDB 文档模型能直接存为 BSON 数组查询时一次取出不需要反范式或 JOIN。这里注意 .collect() 会把全量结果拉到 Driver 端如果结果集超过 Driver 内存就报 OOM生产环境会改成分区写入但单机项目 20 万用户以内直接 collect 没问题。4. 本地跑通最小可运行版本从环境搭建到 Spark 提交的完整闭环4.1 环境选配版本组合是最大的隐性坑基于标题的典型场景是本地开发或小集群实验版本选配决定你是否能一次跑通。常见的稳定组合是 Hadoop 3.3.x Spark 3.3.x Scala 2.12 MongoDB 6.0对应关系是 Spark 3.3 预编译包自带 Scala 2.12你本地装 Scala 2.12.15 即可不要用 Scala 2.13 否则编译报错。具体到 MongoDB社区版 6.0 在 Windows 和 CentOS 上安装都比较干净Windows 安装失败最常见原因是缺少 Visual C Redistributable装完再执行安装包。这里有张环境清单组件版本建议验证命令常见坑JDK1.8 or 11java -versionSpark 3.3 不支持 JDK 17Scala2.12.15scala -version版本过高和 Spark 冲突Spark3.3.xspark-shell --version下载选 pre-built for Apache Hadoop 3.3MongoDB6.0.xmongod --versionWindows 安装失败查 VC 运行库4.2 把完整流程串起来一份可提交的 Spark 作业代码下面是一个完整的最小可运行对象把数据生成、ALS 训练、MongoDB 写入串成一个可执行的作业。这里用一个假想的数据生成器模拟日志完整代码可以直接复制到本地跑但注意 MongoDB 连接地址按实际环境修改import org.apache.spark.sql.SparkSession import org.apache.spark.ml.evaluation.RegressionEvaluator import org.apache.spark.ml.recommendation.ALS object RecommendJob { def main(args: Array[String]): Unit { val spark SparkSession.builder() .appName(ProductRecommend) .master(local[*]) .config(spark.executor.memory, 2g) .getOrCreate() import spark.implicits._ // 第 1 步模拟生成用户行为数据真实场景换成读 HDFS 或 Kafka val behaviorData Seq( (u1, p100, 10.0), (u1, p101, 5.0), (u1, p103, 8.0), (u2, p100, 6.0), (u2, p102, 9.0), (u3, p101, 7.0) ).toDF(userId, productId, rating) // 第 2 步StringIndexer 把字符串 ID 转数值索引 val userIndexer new StringIndexer() .setInputCol(userId) .setOutputCol(userIndex) .fit(behaviorData) val productIndexer new StringIndexer() .setInputCol(productId) .setOutputCol(productIndex) .fit(behaviorData) val indexedData productIndexer.transform(userIndexer.transform(behaviorData)) // 第 3 步拆分训练集和测试集 val Array(training, test) indexedData.randomSplit(Array(0.8, 0.2)) // 第 4 步ALS 训练coldStartStrategy 必须设为 drop val als new ALS() .setMaxIter(10) .setRegParam(0.1) .setUserCol(userIndex) .setItemCol(productIndex) .setRatingCol(rating) .setColdStartStrategy(drop) val model als.fit(training) // 第 5 步评估模型 RMSE val predictions model.transform(test) val evaluator new RegressionEvaluator() .setMetricName(rmse) .setLabelCol(rating) .setPredictionCol(prediction) val rmse evaluator.evaluate(predictions) println(sRMSE $rmse) // 第 6 步为每个用户生成 Top10 推荐写入 MongoDB val topN model.recommendForAllUsers(10) // ... 接前面提到的 bulkWrite 写入逻辑 spark.stop() } }这段代码覆盖了完整的闭环ID 编码、ALS 训练、模型评估、结果生成。有个细节值得注意ALS 的评估不要只看 RMSE因为评分预测的 RMSE 低不代表推荐列表质量高更贴近业务的做法是看推荐列表里的商品是否覆盖了测试集中的正反馈样本。要观察这个指标可以自己写一个函数统计每个用户 TopN 召回里命中的商品比例。4.3 Spark 提交和集群部署local 模式和生产模式的差异本地跑spark-submit和集群跑的唯一区别在主类参数和资源设定上。本地验证时用--master local[*]就能跑通但一旦上了 Spark on YARN就有几个容易遇到的问题。第一个是spark.executor.memory设得太大但 YARN 容器没给够配额导致 executor 起不来第二个是 MongoDB 驱动 jar 包要一起打包到--jars里否则运行时报 NoClassDefFoundError第三个是 Scala 版本必须和 Spark 预编译版本一致否则序列化阶段直接抛 ClassCastException。# 本地调试 spark-submit \ --class com.example.RecommendJob \ --master local[4] \ --jars /path/to/mongo-spark-connector.jar \ recommend-job-1.0.0.jar # 提交到 YARN 集群时注意 --num-executors 和 --executor-cores 的配比 spark-submit \ --class com.example.RecommendJob \ --master yarn \ --deploy-mode cluster \ --num-executors 8 \ --executor-cores 4 \ --executor-memory 8g \ --jars /path/to/mongo-spark-connector.jar \ recommend-job-1.0.0.jar参数说明--num-executors和--executor-cores的比例建议遵循“单 executor 的 vcore 数不超过 5”的原则这是业内共识——每个 executor 分配的 CPU 太多会导致线程争抢磁盘、内存带宽反而降低吞吐。--jars里除了 MongoDB 驱动还要把 Spark 没有默认捆绑的依赖全部放进去比如mongo-java-driver。至于为什么用cluster模式而不是client模式前者 Driver 运行在 YARN 容器里提交命令的终端断开不影响任务后者 Driver 跑在提交机器上终端一关任务就挂生产环境很少有人用 client 模式。5. 推荐系统的生产化改造冷启动、增量更新与 MongoDB 查询优化5.1 冷启动的兜底策略用规则召回填补模型盲区ALS 模型对历史行为丰富的老用户效果好但对新注册用户、新上架商品完全无能为力。不要试图在模型层面解决冷启动——这个问题的标准解法是分层推荐ALS 结果作为第一层规则召回作为第二层两层结果按比例混合输出。规则召回可以是全局热门榜、分类热榜、最近 24 小时上升最快商品榜。这套思路的实现成本很低定期用 Spark SQL 跑一个 groupBy userId 的统计写入 MongoDB 的hot_rank集合线上接口按 user_id 查不到 ALS 结果时直接返回热门榜同时埋点记录冷启动曝光数据用于后续模型训练。5.2 增量更新与全量重算的取舍很多新手以为增量更新就是每天重跑一遍全量 ALS。真实工业界里这确实是最常见的做法——ALS 这类矩阵分解算法本质是全局迭代求解增量训练在数学上很难做到精确等价。业内的做法是每天定时全量重算 TopN 列表同时用近 1 小时的用户实时行为做一个简单的 Item-CF 实时补偿两个列表在接口层合并。全量重算用 Spark 定时调度实时补偿在应用层做不需要引入复杂流计算框架。如果你负责的是一个流量不大的中长尾电商直接每天全量重算 热门榜兜底完全够用。5.3 MongoDB 查询优化三条索引设计规则推荐结果存进 MongoDB 后线上查询性能就看索引设计够不够讲究。访问模式越简单索引设计就越重要——推荐接口的 QPS 往往很高一条慢查询会导致整个服务雪崩。这里总结三条实用规则第一条user_id字段建升序单字段索引支撑按用户取推荐列表的精确查询。第二条需要在推荐结果里按商品 ID 过滤现状时用product_list字段建 multikey 索引但只在确实有这种查询需求时才建——多键索引占空间大写入成本高。第三条如果查询里还要排序比如按推荐分数降序展示把分数字段一起做成复合索引避免内存中排序导致SORT操作卡在 32MB 内存限制上。用 Compass 或 MongoDB Shell 可以快速验证索引是否生效// 给 user_id 建索引确保查询走 IXSCAN 而不是 COLLSCAN db.als_recall_result.createIndex({ user_id: 1 }) // 查看查询执行计划 db.als_recall_result.find({ user_id: u12345 }).explain(executionStats)explain(executionStats)输出里的totalDocsExamined如果等于 1说明索引完全命中如果等于集合总数说明全表扫描了。需要注意的一个经验值当totalDocsExamined与nReturned的比例超过 10:1 时要检查过滤条件是否缺索引SORT阶段如果出现usedDisk: true说明排序超过 32MB 内存上限必须加索引或调整查询逻辑。5.4 推荐质量验证离线指标与线上反馈闭环部署上线后必须有一个可量化的验证步骤否则推荐系统沦为一个“能跑但不知道好不好”的黑盒。常见且低成本的验证方案是在 MongoDB 的推荐结果表里给每条推荐记录加一个recommend_date字段埋点系统记录用户对推荐位的曝光和点击次日离线算 CTR点击/曝光。如果 CTR 低于人工运营位优先怀疑不是模型问题而是召回链路过期或排序逻辑没生效。更进一步可以做一个线上 A/B 对比新模型计算出的结果写入als_recall_result_v2和线上 v1 各放 50% 流量对比一周围观 CTR 与转化率。等 v2 数据持续胜出再把查询切到新集合。这套验证闭环的优点是不需要额外搭建复杂的实验平台只要 MongoDB 里多建一个集合并调整应用层查询配置就能在低成本下完成模型迭代验证。到了这一步基于 SparkScalaMongoDB 的商品推荐系统才算真正闭环——数据、计算、存储、评估、迭代五个环节缺一不可。本文还有配套的精品资源点击获取

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号