news 2026/10/10 11:12:23

Spark+HDFS+MongoDB推荐系统全链路实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Spark+HDFS+MongoDB推荐系统全链路实战

简介:本资源是面向高校大数据课程学习者与初学者的期末实践项目,聚焦分布式电影推荐系统的完整实现,覆盖Hadoop HDFS数据存储、Spark(Scala)实时计算与MongoDB非结构化数据管理三大核心技术栈。压缩包共20个文件,含17个Scala核心业务代码(涵盖数据读写、协同过滤算法实现、特征工程与模型训练)、1个Maven配置文件(pom.xml)用于依赖管理、1个IntelliJ项目配置文件(iml)及1个Manifest文件(mf),整体仅18KB,轻量但结构完整,便于快速导入IDE运行调试。已有435人学习下载,适合作为分布式系统与推荐算法融合实践的入门范例。读者可直接复现端到端流程:从HDFS加载用户行为日志、通过Spark进行矩阵分解建模、将电影元数据与推荐结果存入MongoDB,并获得可扩展的Scala工程骨架与典型大数据组件集成方案。

1. 这不是又一个“协同过滤 Hello World”:它真把 Spark + HDFS + MongoDB 拧成了一根能跑通的推荐流水线

你肯定见过那种“用 MovieLens 数据集跑个 ALS,本地模式 spark-shell 里 print 出来三行推荐结果”的 Scala 示例——它连伪分布式都算不上,更别提和 HDFS、MongoDB 打交道。而这个期末项目 zip 包,是少数几个我亲手 unpack、编译、改配置、在单机伪集群上完整走通了「数据落盘 → HDFS 写入 → Spark 读取 → 特征计算 → 模型训练 → 结果写入 MongoDB → API 查询」全链路的实战工程。它不炫技,但每一步都踩在大数据课程考核的真实边界上:HDFS 不只是hdfs dfs -ls /,而是用FileSystemAPI 写入 Parquet;MongoDB 不是mongo shell里手动 insert,而是通过mongo-spark-connector实现 DataFrame 级别写入;Scala 不是 Java 的语法糖复刻,而是用case class建模、implicit隐式转换处理 Schema、Future异步封装服务层。适合正在啃《Hadoop 权威指南》第 3 章、刚配好spark-shell --master yarn却卡在ClassNotFoundException: org.apache.hadoop.hdfs.DistributedFileSystem的人——它不教你怎么装 Hadoop,但会告诉你core-site.xml和hdfs-site.xml的哪三处配置必须打进 jar 包的resources/目录里,否则spark-submit一运行就报No FileSystem for scheme: hdfs。这不是玩具,是能让你在答辩时打开终端,现场hdfs dfs -cat /movie/data/ratings/part-00000、再spark-submit --class movie.RecommenderApp ...,最后 curlhttp://localhost:8080/recommend/123返回 JSON 推荐列表的硬货。

2. 从解压到跑通:五步拆解项目结构与核心模块依赖

这个 zip 包表面看是标准 Maven 工程(pom.xml+src/main/scala),但它的目录结构和依赖组织,直接暴露了它对 Hadoop 生态的深度绑定。我把它拆成五个可验证步骤,每步都对应一个真实故障点——不是理论,是你马上会遇到的报错。

2.1 解压后第一眼:看清src/main/scala下的三层包结构与数据流向

解压后进入src/main/scala,你会看到三个核心包:

movie/ ├── config/ // SparkConf、MongoConfig、HDFSConfig 的统一管理,不是硬编码! ├── model/ // case class 定义:Rating(userId, movieId, rating, timestamp)、Movie(id, title, genres...) └── pipeline/ // 主干逻辑:DataLoader(从 HDFS 读)、FeatureEngineer(生成用户-电影交叉特征)、Recommender(ALS 训练+预测)、ResultWriter(写 MongoDB)

提示:pipeline/Recommender.scala是主入口,但它的main方法里没有SparkSession.builder(),而是调用config.SparkConfig.getOrCreate()—— 这意味着所有 Spark 配置(如spark.sql.adaptive.enabled)都集中在此,方便你在不同环境(local / yarn)切换。别急着 run,先看pom.xml。

2.2pom.xml里的生死依赖:Hadoop 3.x 兼容性、Mongo Connector 版本、Scala 2.12 的陷阱

打开pom.xml,重点盯这三组<dependency>,它们决定了你能不能跨过编译和运行的第一道墙:

<!-- Hadoop Client 必须显式声明,且版本要和你的 HDFS 集群一致 --> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>3.3.6</version> <!-- 注意:不是 2.x!项目用的是 Hadoop 3.x API --> </dependency> <!-- Mongo Spark Connector:必须和 Spark 版本严格匹配 --> <dependency> <groupId>org.mongodb.spark</groupId> <artifactId>mongo-spark-connector_2.12</artifactId> <version>3.4.1</version> <!-- Spark 3.3.x 对应 connector 3.4.x,Scala 2.12 --> </dependency> <!-- Spark Core & SQL:注意 scope=provided,因为集群已有 --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.12</artifactId> <version>3.3.2</version> <scope>provided</scope> </dependency>

关键点:

  • 如果你本地 Hadoop 是 2.7.7,这里填3.3.6会直接导致java.lang.NoClassDefFoundError: org/apache/hadoop/fs/FSDataInputStream;
  • 如果你用 Spark 3.4.x,但 connector 写3.4.1(它只支持 Spark 3.3.x),spark-submit时会报NoSuchMethodError: org.apache.spark.sql.Dataset.toDF();
  • scope=provided意味着mvn compile能过,但mvn package打出的 jar 不含 Spark 类——你必须用--jars指定集群上的spark-sql_2.12.jar,否则ClassNotFoundException。

2.3config/目录下的三份 XML:为什么core-site.xml必须打进 jar 包?

项目没在代码里写死fs.defaultFS=hdfs://localhost:9000,而是通过config.HDFSConfig加载core-site.xml和hdfs-site.xml。这两份文件在哪?答案是:必须放在src/main/resources/下,且名字一字不差。

# 正确路径(编译后自动进 jar 的 META-INF/resources/) src/main/resources/core-site.xml src/main/resources/hdfs-site.xml

core-site.xml关键内容:

<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> <!-- 必须和你 start-dfs.sh 启动的 NameNode 地址一致 --> </property> </configuration>

hdfs-site.xml关键内容(伪分布式最小配置):

<configuration> <property> <name>dfs.replication</name> <value>1</value> <!-- 单机伪分布式,副本数设为 1 --> </property> <property> <name>dfs.namenode.name.dir</name> <value>file:/usr/local/hadoop/data/namenode</value> </property> </configuration>

注意:如果你把core-site.xml放在/etc/hadoop/下,但没在SparkConf里.set("spark.hadoop.fs.defaultFS", "hdfs://..."),Spark 会忽略系统级配置,坚持用默认file:///—— 这就是为什么spark.read.parquet("hdfs://...")报File not found的根本原因。

2.4movie_recommend.iml文件:IntelliJ IDEA 导入时的隐藏开关

这个.iml文件不是 IntelliJ 自动生成的,而是项目作者手动配置的模块定义。它强制指定了:

  • Scala SDK:必须是 2.12.x(不是 2.11 或 2.13),否则sbt compile会报object scala is not a member of package java.lang;
  • Dependencies Scope:hadoop-client和mongo-spark-connector被标记为Compile,而spark-sql是Provided—— 这直接影响 IDEA 的代码补全和编译类路径;
  • Resources Directory:明确将src/main/resources设为资源根目录,确保core-site.xml在打包时被复制进 jar。

导入 IDEA 时,务必选择 “Import project from external model → Maven”,并勾选 “Auto-import”。如果跳过这步,IDEA 会用默认 Scala SDK,导致import org.apache.spark.sql._下划红线,但mvn compile却能过——这是新手最常翻车的玄学现场。

2.5pom.xml中的maven-shade-plugin:为什么spark-submit一定要用--class movie.RecommenderApp?

项目用maven-shade-plugin打 fat jar,但不是简单地把所有依赖塞进去。看它的配置:

<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.4.1</version> <executions> <execution> <phase>package</phase> <goals><goal>shade</goal></goals> <configuration> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>movie.RecommenderApp</mainClass> <!-- 这是 spark-submit 的入口 --> </transformer> </transformers> <filters> <filter> <artifact>*:*</artifact> <excludes> <exclude>META-INF/*.SF</exclude> <exclude>META-INF/*.DSA</exclude> <exclude>META-INF/*.RSA</exclude> </excludes> </filter> </filters> </configuration> </execution> </executions> </plugin>

这意味着:

  • mvn clean package生成的target/movie-recommend-1.0-SNAPSHOT.jar是一个可执行 jar;
  • spark-submit --class movie.RecommenderApp movie-recommend-1.0-SNAPSHOT.jar能直接运行;
  • 但如果你写成--class movie.pipeline.Recommender(少了个App),会报ClassNotFoundException—— 因为RecommenderApp.scala是真正的 main class,而Recommender.scala只是算法实现类。

3. HDFS 数据准备:从原始 CSV 到 Parquet 分区表的四步实操

项目不会帮你把 MovieLens 数据扔进 HDFS。它假设你已准备好/movie/data/ratings.csv和/movie/data/movies.csv。但“准备好”不是hdfs dfs -put就完事——HDFS 上的数据格式、分区、压缩,直接决定 Spark 作业的性能和稳定性。我按生产环境习惯,拆成四步。

3.1 原始数据清洗:用awk和sed处理 MovieLens 的字段错位

MovieLens 1M 数据集的ratings.dat是::分隔,但项目DataLoader.scala期望的是 CSV 格式。直接tr '::' ','会出问题:电影标题里有逗号(如"Toy Story (1995)")。正确做法是用awk精准切分:

# 将 ratings.dat 转为标准 CSV(userId,movieId,rating,timestamp) awk -F"::" '{print $1 "," $2 "," $3 "," $4}' ratings.dat > ratings.csv # 将 movies.dat 转为 CSV(movieId,title,genres),注意 genres 是 | 分隔,需转义 awk -F"::" '{ gsub(/\|/, "\\|", $3) # 将 genres 中的 | 替换为 \|,避免后续 split 错乱 print $1 "," "\"" $2 "\"" "," "\"" $3 "\"" }' movies.dat > movies.csv

提示:movies.csv的 title 字段必须加双引号,否则 Spark 读 CSV 时遇到The Matrix (1999)会把括号当字段分隔符,导致列数错乱。这是血泪经验——我曾因此 debug 3 小时,发现df.count()返回 0。

3.2 上传前校验:用hdfs fsck确保目标路径可写且无残留

别急着hdfs dfs -put。先确认 HDFS 状态和路径权限:

# 检查 NameNode 是否健康(返回 HEALTHY 表示 OK) hdfs fsck / -files -blocks -locations | head -20 # 创建目标目录,并设为 777(伪分布式开发环境,生产环境请用 proper ACL) hdfs dfs -mkdir -p /movie/data hdfs dfs -chmod -R 777 /movie/data # 清空旧数据(避免重复写入导致 Spark 读到脏数据) hdfs dfs -rm -r /movie/data/ratings hdfs dfs -rm -r /movie/data/movies

注意:hdfs dfs -chmod 777 /movie/data在 Hadoop 3.x 是允许的,但如果你用的是启用了 Kerberos 的集群,这条命令会失败——此时必须用hdfs dfs -chown youruser:supergroup /movie/data。

3.3 上传并转存为 Parquet:为什么不用 CSV 直接读?

项目DataLoader.scala里,loadRatingsFromHDFS()方法明确指定:

val ratingsDF = spark.read .option("header", "true") .option("inferSchema", "true") .parquet("hdfs://localhost:9000/movie/data/ratings") // 注意:是 parquet,不是 csv!

所以你必须把 CSV 转成 Parquet。本地转存脚本(convert_to_parquet.sh):

#!/bin/bash # 用 spark-submit 调用一个临时转换脚本 spark-submit \ --master local[*] \ --class movie.util.CSVToParquetConverter \ target/movie-recommend-1.0-SNAPSHOT.jar \ "file:///path/to/local/ratings.csv" \ "hdfs://localhost:9000/movie/data/ratings" \ "file:///path/to/local/movies.csv" \ "hdfs://localhost:9000/movie/data/movies"

对应的CSVToParquetConverter.scala核心逻辑:

def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("CSV to Parquet Converter") .master("local[*]") .getOrCreate() // 读 CSV,强制指定 schema 避免 inferSchema 的类型错误 val ratingsSchema = StructType(Array( StructField("userId", IntegerType, nullable = false), StructField("movieId", IntegerType, nullable = false), StructField("rating", DoubleType, nullable = false), StructField("timestamp", LongType, nullable = false) )) val ratingsDF = spark.read .option("header", "true") .schema(ratingsSchema) // 关键!避免 timestamp 被 infer 为 string .csv(args(0)) .write .mode("overwrite") .parquet(args(1)) // 写入 HDFS 的 Parquet 目录 spark.stop() }

提示:inferSchema=true在大数据量下极慢,且易出错(如某行 timestamp 是空字符串,整列 infer 为 string)。项目虽没写死 schema,但你在DataLoader里应该补上——这是性能优化的第一步。

3.4 分区与压缩:给 Parquet 加上snappy压缩和date分区

项目当前没做分区,但实际中,ratings表按date分区能极大加速时间范围查询(如“最近 30 天评分”)。修改CSVToParquetConverter.scala:

// 在 write 前加一行:按日期分区(假设 timestamp 是毫秒,转为 yyyy-MM-dd) val ratingsWithDate = ratingsDF .withColumn("date", date_format(from_unixtime(col("timestamp") / 1000), "yyyy-MM-dd")) ratingsWithDate .write .mode("overwrite") .option("compression", "snappy") // 比默认 uncompressed 小 3x,CPU 开销可控 .partitionBy("date") // 生成 /movie/data/ratings/date=2023-01-01/ 目录 .parquet(args(1))

验证分区是否生效:

hdfs dfs -ls /movie/data/ratings | head -10 # 应看到: # drwxr-xr-x - hadoop supergroup 0 2023-10-01 10:00 /movie/data/ratings/date=2023-01-01 # drwxr-xr-x - hadoop supergroup 0 2023-10-01 10:00 /movie/data/ratings/date=2023-01-02

4. Spark 推荐引擎落地:ALS 模型训练、实时预测与 MongoDB 写入的闭环

movie.pipeline.Recommender是整个项目的黑匣子,但它不是魔法。我把它的 ALS 训练、预测、写入流程拆成三步可调试的环节,每步都附带spark-shell交互式验证命令——你不需要跑完整 App,就能确认模型是否真的在工作。

4.1 ALS 训练参数调优:rank、maxIter、regParam的取舍逻辑

项目Recommender.scala中的 ALS 配置:

val als = new ALS() .setMaxIter(10) // 迭代次数:10 是平衡精度和时间的经验值 .setRegParam(0.01) // 正则化系数:防止过拟合,MovieLens 1M 数据,0.01 较稳妥 .setRank(10) // 隐语义维度:10 是经典起点,<5 效果差,>20 易过拟合且内存暴涨 .setUserCol("userId") .setItemCol("movieId") .setRatingCol("rating") .setPredictionCol("prediction")

为什么不是rank=50?因为内存消耗公式是:O(rank * (numUsers + numItems))。MovieLens 1M 有 6000 用户、4000 电影,rank=10时模型矩阵约 10*(6000+4000)=100,000 个浮点数;rank=50就是 500,000,单机 JVM 很容易 OOM。我在本地测试过,rank=10的 RMSE 是 0.87,rank=20是 0.85,提升仅 2.3%,但训练时间从 42s 增至 118s。

验证模型是否正常训练:

# 进入 spark-shell,手动加载数据并训练 spark-shell --master local[*] --jars /path/to/mongo-spark-connector_2.12-3.4.1.jar scala> val ratings = spark.read.parquet("hdfs://localhost:9000/movie/data/ratings") scala> import org.apache.spark.ml.recommendation.ALS scala> val als = new ALS().setMaxIter(5).setRank(10).setRegParam(0.01) scala> val model = als.fit(ratings) scala> model.userFactors.show(1) // 查看用户因子矩阵第一行,确认非空

4.2 实时预测:model.transform()vsmodel.recommendForAllUsers()的场景选择

项目用的是model.recommendForAllUsers(10),即为每个用户生成 Top10 推荐。但这是离线批量操作,耗时长。如果你需要实时响应(如用户点开个人页立即显示推荐),必须用model.transform():

// 构造单个用户的评分记录(冷启动?用平均分填充) val singleUserDF = spark.createDataFrame(Seq( (123, 1, 0.0), (123, 2, 0.0), (123, 3, 0.0) // userId=123 对 movieId=1,2,3 的隐式评分 )).toDF("userId", "movieId", "rating") val predictions = model.transform(singleUserDF) predictions.select("userId", "movieId", "prediction").show()

注意:transform()要求输入 DataFrame 必须包含userId和movieId列,且movieId必须在训练集中出现过(否则 prediction 为 NaN)。冷启动问题(新用户/新电影)项目没处理,你需要加 fallback 逻辑,比如用热门电影或基于内容的推荐。

4.3 写入 MongoDB:mongo-spark-connector的WriteConfig关键参数

ResultWriter.scala用DataFrame.write.format("com.mongodb.spark.sql.DefaultSource")写入,但核心是WriteConfig:

val writeConfig = WriteConfig(Map( "uri" -> "mongodb://localhost:27017", "database" -> "movie_db", "collection" -> "recommendations", "replaceDocument" -> "false", // true 会覆盖同 _id 文档,false 则追加 "spark.mongodb.output.ignoreNulls" -> "true" // 避免 null 字段写入 )) recommendationsDF.write .mode("append") // 注意:不是 overwrite,否则每次重跑清空历史 .format("com.mongodb.spark.sql.DefaultSource") .options(writeConfig.asOptions) .save()

验证写入是否成功:

# 进入 mongo shell $ mongo > use movie_db > db.recommendations.find({userId: 123}).limit(1).pretty() # 应返回类似: # { # "_id" : ObjectId("65a1b2c3d4e5f67890123456"), # "userId" : 123, # "movieIds" : [ 456, 789, 101 ], # "scores" : [ 4.8, 4.5, 4.3 ], # "timestamp" : ISODate("2023-10-01T10:00:00Z") # }

提示:replaceDocument=false是安全底线。我曾误设为true,导致userId=123的推荐结果被新批次覆盖,前端查不到历史推荐——这种 bug 在日志里完全不报错,只能靠人工比对 MongoDB 文档。

5. 避坑指南:五个让 90% 新手卡住的致命细节与解决方案

这个项目看似结构清晰,但每一个技术栈的交接处都是深坑。以下是我在三台不同配置的虚拟机(Ubuntu 20.04 / CentOS 7 / macOS M1)上反复踩过的五个具体问题,每个都按「现象 → 原因 → 解决」给出可执行方案。

5.1 现象:spark-submit报java.lang.ClassNotFoundException: org.apache.hadoop.fs.FileSystem

原因:hadoop-client依赖在pom.xml中 scope 是compile,但spark-submit运行时,Spark 集群的 classpath 里没有 Hadoop 的 JAR。尤其当你用--master yarn时,YARN NodeManager 的HADOOP_HOME环境变量未正确指向 Hadoop 安装目录,导致FileSystem类找不到。

解决:

  1. 确认HADOOP_HOME已设(echo $HADOOP_HOME应输出/usr/local/hadoop);
  2. 在spark-submit命令中显式添加 Hadoop JAR:
spark-submit \ --master yarn \ --jars $HADOOP_HOME/share/hadoop/common/hadoop-common-3.3.6.jar,$HADOOP_HOME/share/hadoop/hdfs/hadoop-hdfs-3.3.6.jar \ --class movie.RecommenderApp \ target/movie-recommend-1.0-SNAPSHOT.jar
  1. 更彻底的方案:把hadoop-client的scope改为compile,并在maven-shade-plugin中排除hadoop-*的重复类(避免Duplicate class错误)。

5.2 现象:MongoDB 写入后,db.recommendations.count()返回 0,但spark-submit日志显示Write completed

原因:mongo-spark-connector默认使用WriteMode.Append,但如果 MongoDB 集合不存在,它不会自动创建集合,而是静默失败。更隐蔽的是,如果uri中的数据库名拼写错误(如movie_db写成movie-db),connector 会连接到一个空数据库,写入无报错但数据不可见。

解决:

  1. 提前在 MongoDB 中创建集合并插入一条测试文档:
$ mongo > use movie_db > db.recommendations.insertOne({test: "init"}) > db.recommendations.countDocuments({}) # 应返回 1
  1. 在WriteConfig中强制指定spark.mongodb.output.createCollectionOptions:
"spark.mongodb.output.createCollectionOptions" -> "{collation: {locale: 'en'}}"
  1. 检查mongo.log,搜索insert关键字,确认是否有WriteResult。

5.3 现象:spark.read.parquet("hdfs://...")报java.io.IOException: Failed on local exception: java.io.IOException: Response too long

原因:HDFS 的dfs.client.socket-timeout默认是 60000ms(60秒),当 Parquet 文件过大(>1GB)或网络延迟高时,客户端等待超时。这不是数据问题,是 RPC 超时。

解决:

  1. 修改hdfs-site.xml,增加超时配置:
<property> <name>dfs.client.socket-timeout</name> <value>300000</value> <!-- 5分钟 --> </property> <property> <name>dfs.client.read.shortcircuit.streams.cache.size</name> <value>4096</value> </property>
  1. 在 Spark 代码中设置:
spark.conf.set("spark.hadoop.dfs.client.socket-timeout", "300000")
  1. 重启 HDFS:stop-dfs.sh && start-dfs.sh。

5.4 现象:mvn compile成功,但 IDEA 中import org.apache.spark.sql.functions._报红,且spark变量无代码提示

原因:IDEA 的 Scala SDK 未正确关联spark-sql_2.12的源码和文档。pom.xml中spark-sql的 scope 是provided,IDEA 默认不下载其依赖,导致索引缺失。

解决:

  1. 在 IDEA 中,File → Project Structure → Libraries,点击+→Java,导航到$SPARK_HOME/jars/spark-sql_2.12-3.3.2.jar;
  2. 右键该 jar →Download Sources and Documentation;
  3. 在Project Settings → Modules → Dependencies中,找到spark-sql_2.12,将其Scope临时改为Compile,应用后立刻恢复为Provided—— 这能强制 IDEA 重新索引。

5.5 现象:ALS 训练时java.lang.OutOfMemoryError: GC overhead limit exceeded

原因:ALS的fit()方法在 Driver 端构建模型时,会将用户因子和物品因子矩阵全部加载进内存。rank=10时内存尚可,但若数据集扩大(如 MovieLens 10M),或rank设为 50,Driver 内存必然爆。

解决:

  1. 增加 Driver 内存:spark-submit --driver-memory 4g;
  2. 关键:启用ALS的checkpointInterval,将中间 RDD 持久化到 HDFS,减少内存压力:
spark.sparkContext.setCheckpointDir("hdfs://localhost:9000/tmp/checkpoint") val als = new ALS().setCheckpointInterval(2) // 每2次迭代 checkpoint 一次
  1. 最终方案:改用ALS.trainImplicit()(隐式反馈),它比显式评分训练内存占用低 40%,且更适合点击/播放时长等行为数据。

6. 从离线批处理到轻量实时:用 Spark Streaming 接入 Kafka 日志流的改造技巧

项目当前是纯离线批处理:每天凌晨跑一次 ALS,更新全量推荐。但真实业务需要“用户刚打五星,10 秒内出现在好友推荐页”。我把它升级为 Lambda 架构——离线层(ALS 全量)+ 实时层(Streaming 增量)。核心改动只有三处,且完全兼容原项目结构。

6.1 新增 Kafka Producer:模拟用户实时评分事件

不改原有ratings.csv,新增kafka-producer.py,向 topicmovie-ratings发送 JSON:

from kafka import KafkaProducer import json import time import random producer = KafkaProducer(bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8')) # 模拟用户评分事件 events = [ {"userId": 123, "movieId": 456, "rating": 5.0, "timestamp": int(time.time())}, {"userId": 789, "movieId": 101, "rating": 4.5, "timestamp": int(time.time())} ] for event in events: producer.send('movie-ratings', value=event) time.sleep(1) producer.flush()

6.2 新增 Streaming Job:StreamingRecommender.scala,用foreachBatch写入 MongoDB

新建src/main/scala/movie/streaming/StreamingRecommender.scala:

object StreamingRecommender { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("Streaming Recommender") .master("local[*]") .config("spark.sql.adaptive.enabled", "true") .getOrCreate() // 从 Kafka 读取 val kafkaDF = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "movie-ratings") .option("startingOffsets", "latest") .load() // 解析 JSON val ratingsDF = kafkaDF .selectExpr("CAST(value AS STRING)") .select(from_json(col("value"), ratingSchema).alias("data")) .select("data.*") // 增量写入 MongoDB(注意:用 foreachBatch 避免 foreach 写入的序列化问题) ratingsDF.writeStream .foreachBatch { (batchDF: DataFrame, batchId: Long) => batchDF.write .format("com.mongodb.spark.sql.DefaultSource") .mode("append") .options(WriteConfig(Map( "uri" -> "mongodb://localhost:27017", "database" -> "movie_db", "collection" -> "realtime_ratings" )).asOptions) .save() } .start() .awaitTermination() } }

6.3 离线与实时的融合:在RecommenderApp中加入realtime_ratings的权重衰减

原RecommenderApp只读hdfs://.../ratings。现在,我们让它同时读 HDFS 全量 + MongoDB 实时流,并用时间衰减加权:

// 读实时评分(过去 1 小时) val recentRatings = spark .read .format("com.mongodb.spark.sql.DefaultSource") .option("uri", "mongodb://localhost:27017") .option("database", "movie_db") .option("collection", "realtime_ratings") .load() .filter(col("timestamp") > unix_timestamp() - 3600) // 过去 1 小时 .withColumn("weight", lit(0.3)) // 实时数据权重 0.3 // 读 HDFS 全量(权重 0.7) val fullRatings = spark.read.parquet("hdfs://.../ratings") .withColumn("weight", lit(0.7)) // 合并并加权 val mergedRatings = recentRatings.unionByName(fullRatings) .withColumn("rating", col("rating") * col("weight")) // 用 mergedRatings 训练 ALS...

从那以后我每次做推荐系统,都强制走一遍“离线全量 + 实时增量”的双链路验证:先spark-submit跑通 ALS,再spark-submit --class movie.streaming.StreamingRecommender启动流任务,最后用mongo查realtime_ratings和recommendations两个集合,对比同一用户的推荐结果是否随实时行为动态变化。这招帮我避开了 80% 的线上推荐不准问题。希望帮到你。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/10 11:12:02

硅碳相变:大模型微调到底值不值得做?开发者技术解析与决策清单

硅碳相变&#xff1a;大模型微调到底值不值得做&#xff1f;开发者技术解析与决策清单 后台工程师、算法同学、AI 应用负责人&#xff0c;你们问得最多的一句话我替你们说了&#xff1a;手上这个业务&#xff0c;到底该不该上微调&#xff1f;我做了两年多模型接入和推理服务的…

作者头像 李华
网站建设 2026/10/10 11:10:50

LeetCode 128最长连续序列:从排序到O(n)哈希表全解析

很多刷题的人第一次见到 LeetCode 128“最长连续序列”这道题&#xff0c;第一反应是排序&#xff1a;排完序扫一遍不就完了嘛。要是你面试时真这么答&#xff0c;面试官多半会追问一句&#xff1a;“能不能做到 O(n)&#xff1f;” 这道题之所以被列为经典中的经典&#xff0c…

作者头像 李华
网站建设 2026/10/10 11:10:42

Cherry Studio接入DeepSeek:API配置、RAG知识库与避坑指南

简介&#xff1a;《解锁AI新体验&#xff1a;Cherry Studio安装指南及DeepSeek完美融合》是一份面向AI工具爱好者与开发者的操作型图文文档。它聚焦Cherry Studio桌面客户端与DeepSeek大模型的集成应用&#xff0c;从安装部署、API密钥配置到功能实操均有覆盖&#xff0c;适合希…

作者头像 李华
网站建设 2026/10/10 11:10:41

Java 17调用Responses图像输入:商品问答核心链路与结果边界实践

做了大半年商品问答项目&#xff0c;踩了不少坑之后&#xff0c;把Java 17调用Responses图像输入的核心链路整理出来了。这篇文章重点聊聊商品问答场景里&#xff0c;怎么把商品图片喂给语言模型&#xff0c;怎么拿到结构化结果&#xff0c;以及最容易被忽视的"结果边界&q…

作者头像 李华
网站建设 2026/10/10 11:10:39

Linux Socket 编程实战:从 TCP/epoll 到高并发优化

简介&#xff1a;面向 Linux 网络编程入门与进阶的读者&#xff0c;这份 Word 文档以图文方式系统梳理 Socket 编程的核心知识&#xff1a;网络中进程如何通信、Socket 的本质与设计理念&#xff0c;以及 socket、bind、listen、connect、accept、read、write、close 等基础接口…

作者头像 李华
网站建设 2026/10/10 11:10:37

基于Spark的地铁客流分析系统:架构、实践与避坑指南

简介&#xff1a;面向计算机专业毕业设计的完整项目&#xff0c;基于Spark的地铁大数据客流分析系统&#xff0c;以城市地铁客流数据为分析对象&#xff0c;覆盖数据采集、清洗、存储、分析、可视化与客流预测等环节&#xff0c;适合大数据方向学生用于课程设计、毕设参考或技术…

作者头像 李华