免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

Spark Streaming实时推荐系统:从商品关注度到多路召回融合实践

Spark Streaming实时推荐系统:从商品关注度到多路召回融合实践 简介基于Spark的电商商品智能分析系统面向大数据、推荐算法方向的毕业设计或课程设计学习者完整实现了流式计算商品关注度、智能推荐与关联规则分析等核心功能。资源共939个文件约5.49MB涵盖Java/Scala源码、Spark Streaming任务输出的part文件、HTML/JS前端页面、XML配置、class编译文件等目录包含数据采集、模型训练、结果展示等模块便于对照源码理解从实时行为统计到推荐生成的完整链路。已有246人学习下载适合希望快速搭建Spark电商分析项目的学生和研究者。压缩包内不仅包含可运行项目与依赖配置还保留checkpoint、日志等中间产物便于排查运行过程结合商品关注度计算、协同过滤、FP-Growth关联规则等知识点可支撑实验复现、二次开发或毕业设计文档撰写是一份兼顾学习与实战的完整资料。1. 从用户点击到商品热度为什么关注度计算要放在流上商品关注度是一个被低估的指标。我见过不少推荐项目离线任务每天凌晨跑一次热度统计结果用户上午刚看完一台手机下午就收到完全无关的耳机推荐。真正让推荐系统“活着”的是对分钟级、甚至秒级行为的响应。这个问题本质上不是算法选型问题而是架构问题点击、加购、搜索这些信号天然是一条流只有用流式计算才能表达它的时效性。这个项目正好把大数据处理的经典组件串起来了Spark Streaming 负责实时消费用户行为用窗口聚合产出商品关注度再结合协同过滤和 FP-Growth 关联规则把“大家都在买什么”变成“你可能还想要什么”。如果你正在做毕设、或者想在简历里补一个完整的大数据推荐项目这套代码的模块划分和调用链值得拆开看一遍。下面我会按“流式计算 → 特征召回 → 协同过滤排序 → 关联规则融合 → 环境调试”这条主线展开所有代码都在 Spark 2.4 / PySpark 3.x 下可运行。2. Spark Streaming 关注度计算窗口、水位与事件时间设计2.1 为什么关注度必须用流式计算而不是批处理商品关注度如果离线算只能得到“过去 24 小时的热销榜”这对推荐排序几乎没有增量价值。电商场景里用户行为是持续到达的真正有用的关注度是“最近 5 分钟被快速浏览的商品”。把批处理任务调成每 5 分钟跑一次也不现实Spark DAG 启动开销、数据分区扫描、元数据读取都会造成分钟级延迟。用 Structured Streaming 的意义在于它把流当成一张无边界的表用 window 和 watermark 做时间维度的聚合语义上和批处理一致但延迟能压到秒级。另一个容易被忽略的细节埋点系统的时间戳经常不是 Spark 收到数据的时间。用户手机离线、网络拥塞都会导致事件延迟到达。如果只用处理时间聚合晚到的点击会被算进错误的窗口。这个项目里我推荐用事件时间event_time作为处理基准配合 watermark 容忍乱序。2.2 埋点数据结构和 Kafka 入参约定关注度计算的前提是埋点有足够的维度。项目里 Kafka topicuser_behavior的消息体是 JSON核心字段如下{ user_id: u_1001, product_id: p_2034, behavior: click, stay_seconds: 12, event_time: 2025-01-10T10:23:4508:00 }behavior分为click、add_cart、purchase、collect四类stay_seconds只在点击场景有值。读取 Kafka 的 PySpark 代码可以这样写from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType, TimestampType schema StructType([ StructField(user_id, StringType()), StructField(product_id, StringType()), StructField(behavior, StringType()), StructField(stay_seconds, DoubleType()), StructField(event_time, TimestampType()), ]) spark SparkSession.builder \ .appName(ProductAttentionStream) \ .config(spark.sql.shuffle.partitions, 8) \ .getOrCreate() kafka_df spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, node01:9092,node02:9092) \ .option(subscribe, user_behavior) \ .option(startingOffsets, latest) \ .load() \ .selectExpr(CAST(value AS STRING) as json_str) parsed_df kafka_df \ .selectExpr(from_json(json_str, schema) as data) \ .select(data.user_id, data.product_id, data.behavior, data.stay_seconds, data.event_time)这里的schema用 StructType 声明避免后面聚合时因类型推断失败而断流。startingOffsetslatest表示只消费新数据调试阶段如果想重放历史消息改成earliest。spark.sql.shuffle.partitions设得太大会造成大量小文件流任务里通常 8 到 32 就够。2.3 滑动窗口聚合统计 5 分钟关注度关注度要能反映“最近热度”窗口不能是固定的一天。我用 10 分钟长度、5 分钟滑动步长的窗口作基础热度代码里用window函数表达from pyspark.sql.functions import window, when, col, count, avg, sum attention_df parsed_df \ .withWatermark(event_time, 2 minutes) \ .groupBy( window(event_time, 10 minutes, 5 minutes), product_id ) \ .agg( count(when(col(behavior) click, 1)).alias(click_cnt), count(when(col(behavior) add_cart, 1)).alias(cart_cnt), count(when(col(behavior) purchase, 1)).alias(purchase_cnt), avg(when(col(behavior) click, col(stay_seconds))).alias(avg_stay_seconds) )withWatermark(event_time, 2 minutes)的含义是允许事件时间晚于当前事件时间 2 分钟的乱序数据进入窗口超过这个阈值的消息会被标记为过期并被丢弃。窗口长度 10 分钟、滑动间隔 5 分钟意味着每 5 分钟输出一次最近 10 分钟的热度。实际生产中如果商品曝光频次高窗口可以缩到 2 分钟滑动间隔 1 分钟如果希望捕捉尖峰长度不要设太大否则跟批处理没区别。聚合结果写出去前还需要一个行为权重转换点击 1 分、加购 5 分、收藏 3 分、购买 10 分。同时要解决长停留时长带来的噪声停留 1 小时可能只是用户忘了关页面。可以用if(avg_stay_seconds 300, 1.0, avg_stay_seconds/300)做截断。2.4 热度衰减公式与输出到 Redis窗口聚合得到的是窗口内的绝对值但在推荐排序里更常用的是带时间衰减的累计热度。我用一个近似指数衰减的方式合并新旧窗口from pyspark.sql.functions import lit, exp, col decay_factor 0.5 # 每 3 分钟衰减一半 time_gap 5 # 当前窗口和上一窗口相隔 5 分钟 result_df attention_df.withColumn( attention_score, col(click_cnt) * 1 col(cart_cnt) * 5 col(collect_score) * 3 col(purchase_cnt) * 10 ) \ .withColumn( decayed_score, col(attention_score) * exp(-lit(decay_factor) * lit(time_gap) / lit(3)) )注意这里每次窗口触发时上一窗口的分数会乘上一个小于 1 的系数再累加而不是直接覆盖。这样就能让一个商品因为用户短时间大量点击而冲到推荐列表头部但也不会因为一次促销后永远霸榜。输出的落点一般选 Redis用product_id:attention做 keyscore 存给在线服务读取。3. 商品推荐召回基于内容的相似度计算与候选生成3.1 召回层为什么先用内容相似度推荐系统里协同过滤在冷启动和长尾场景经常失效新商品没行为、新用户没历史。基于内容的召回只依赖商品自身属性不依赖交互记录所以它是整个推荐链路里最稳的一路候选来源。这个项目里我把商品标题、类目、品牌、颜色、价格区间拼成一个文本特征算了相似度作为协同过滤前面的粗筛。商品属性在 MySQL 里可以用spark.read.jdbc拉取。结构大致是字段示例product_idp_2034titleRedmi K70 手机 12GB256GBcategory智能手机brandRedmiprice2499attributes白色;5000mAh;5G;OLEDattributes是用分号拼接的多值字段处理时按分号切分再展开成多行。3.2 用 TF-IDF 构建商品特征向量商品标题和类目往往很短直接用 one-hot 会得到高维稀疏向量。我用HashingTF把词映射到固定维度的索引再用IDF调整权重避免“手机”“智能”这类高频词主导相似度。PySpark 代码如下from pyspark.ml.feature import HashingTF, IDF, Tokenizer from pyspark.ml import Pipeline df spark.table(product_dim) \ .withColumn(text, concat_ws( , title, category, brand, attributes)) tokenizer Tokenizer(inputColtext, outputColwords) hashing_tf HashingTF(inputColwords, outputColrawFeatures, numFeatures2048) idf IDF(inputColrawFeatures, outputColfeatures) pipeline Pipeline(stages[tokenizer, hashing_tf, idf]) model pipeline.fit(df) feature_df model.transform(df).select(product_id, features)numFeatures2048是经验值。太小容易碰撞编码 5 万商品时建议 4096太大对内存不友好而且后面算相似度时广播代价高。用Pipeline的好处是上线新商品时可以直接model.transform不需要重算 IDF。3.3 相似度计算的两种工程做法特征向量出来后计算商品两两相似度有两种常见做法方案实现路径优点缺点离线全量计算两个向量 join 后算余弦相似度精确、可分析O(N^2)5 万商品跑不动近似近邻LSH / HNSW百万级商品可扩展有查询损失我建议毕业设计阶段直接用窄表 join 限制候选范围只对同一 category 下的商品两两计算相似度再把结果过滤到 top 50 存入 Redis。代码片段如下from pyspark.ml.linalg import Vectors from pyspark.sql.functions import udf from pyspark.sql.types import FloatType import numpy as np def cosine_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 cosine_udf udf(cosine_sim, FloatType()) candidate category_join \ .filter(product_a product_b) \ .withColumn(similarity, cosine_udf(features_a, features_b)) \ .filter(similarity 0.5)这里强制product_a product_b是为了去掉重复对让每个商品对只算一次。similarity 0.5是相似度的最低阈值低于 0.5 的所谓“相似商品”往往是噪声对召回质量没有贡献。如果数据量再大一个量级建议换 LSH 的approxSimilarityJoin。3.4 内容召回候选集的归一化内容相似度分数和后面协同过滤的打分尺度不一样直接相加会把排序结果带偏。我在融合前都会做一次 min-max 归一化把内容相似度压到 0 到 1 区间。归一化权重存到一个 broadcast 变量里在线服务读取时可以直接用省得每次请求都重新算。4. 个性化排序基于 ALS 的协同过滤模型训练4.1 行为日志如何转成偏好评分协同过滤需要一张 user-item 评分表。这个项目里没有显式评分所以我从关注度结果里映射出偏好分数from pyspark.sql.functions import col, when interaction_df parsed_df \ .filter(user_id IS NOT NULL AND product_id IS NOT NULL) \ .select( col(user_id), col(product_id), when(col(behavior) click, 1.0) .when(col(behavior) add_cart, 3.0) .when(col(behavior) collect, 2.0) .when(col(behavior) purchase, 5.0) .otherwise(0.0).alias(rating) )评分不是越高越好。点击给 1 分购买给 5 分但如果一个用户连续看了同一商品 10 次聚合时 rating 会变成 10 分这个数值会干扰 ALS 对未来偏好的预测。所以我通常在聚合前先做一次行为次数上限截断比如单个 user-product 的最大评分封顶 5 分公式是min(actual_rating, 5.0)。4.2 ALS 模型训练与参数选择PySpark MLlib 里的 ALS 可以直接吃 DataFrame 格式from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator als ALS( userColuser_id_int, itemColproduct_id_int, ratingColrating, rank20, maxIter15, regParam0.1, coldStartStrategydrop ) model als.fit(train_df)参数说明rank20隐因子数量。用户行为很稀疏时 10 到 20 足够数据量大可以试 50但模型体积和推理时间都会增加。maxIter15优化迭代次数。超过 20 次对 RMSE 的改善很少反而容易过拟合。regParam0.1正则化系数。调参时优先看 0.01、0.05、0.1 三档。coldStartStrategydrop遇到训练集里没见过的新用户/商品时预测结果会是 nulldrop 可以直接把它从评估中剔除。如果线上的推荐服务无法处理 null可以改成fill并指定coldStartStrategy补一个中性分数。训练前需要把 user_id 和 product_id 转成连续的数值 ID否则 ALS 会报userCol不支持字符串。一种简单做法是用StringIndexerfrom pyspark.ml.feature import StringIndexer uid_indexer StringIndexer(inputColuser_id, outputColuser_id_int) pid_indexer StringIndexer(inputColproduct_id, outputColproduct_id_int)4.3 模型评估RMSE 与业务命中率训练集和测试集我用randomSplit([0.8, 0.2])切分。评估指标除了 RMSE还应该加一个 top-K 召回率测试集里用户真正交互过的商品在模型给出的 top 20 推荐里出现了多少。evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) rmse evaluator.evaluate(predictions) print(fRMSE {rmse:.3f})RMSE 在 1.0 以内算正常。如果跑到 2.0 以上先检查评分分布是不是有极端值再看regParam是否过小。业务上更关注的其实是排序结果而不是预测分绝对误差。4.4 冷启动用户的后备策略ALS 给新用户预测时只能得到 null。这个项目里的做法是新用户用“热门榜 内容召回”兜底。热门榜直接取第 2 章流式计算产出的attention_score前 50 个商品。把coldStartStrategydrop的预测结果和兜底结果做一个 union再做权重合并。这样用户没有历史行为时也能收到合理推荐等行为积累到 5 条以上再切换成个性化排序。5. FP-Growth 关联规则挖掘与多路召回融合5.1 为什么选 FP-Growth 而不是 AprioriApriori 的经典问题是要反复扫描全量数据生成候选集商品量一旦上来性能会断崖式下跌。FP-Growth 用 FP 树压缩事务一次扫描就能统计频繁项集在 Spark 里属于开箱即用的实现。关联规则挖掘的目的不是做最终推荐而是处理“买了一台相机通常也会买 SD 卡和相机包”这种强关联场景。这在协同过滤里很难被捕捉因为协同过滤擅长发现 user-item 的隐空间关系却不擅长表达商品间的显式共现。5.2 会话切分与事务构造关联规则需要的是“一次会话里用户买了哪些商品”而不是全量行为。我给每条行为打上 session_id来源可以是后端埋点里的 session 字段或者按 user_id 30 分钟无操作切分。事务表的构建常见做法是SELECT session_id, collect_set(product_id) AS products FROM user_behaviors WHERE behavior purchase GROUP BY session_id这里用collect_set而不是collect_list是为了去除同一个 session 里重复购买同一商品产生的冗余项。事务只保留购买行为点击和加购产生的噪声太大会出现大量虚假关联。5.3 Spark 里 FP-Growth 的训练与规则过滤PySpark 提供了FPGrowth实现from pyspark.ml.fpm import FPGrowth fp_growth FPGrowth( minSupport0.003, minConfidence0.4, itemsColproducts, predictionColrules ) fp_model fp_growth.fit(transaction_df) fp_model.freqItemsets.show(10) fp_model.associationRules.show(10)参数说明minSupport0.003项集在事务中出现的频率下限。事务量 10 万时0.003 意味着至少出现 300 次。支持度太高会丢失有价值的尾部关联。minConfidence0.4规则置信度。取 0.4 意味着“买了 A 的用户里 40% 会买 B”这条规则才被保留。predictionColrulesFitting 完成后模型会对每个事务预测可能关联的商品这就是离线生成的补充推荐。关联规则生成后需要再用提升度lift过滤一次。置信度高不一定代表真正的关联因为 B 本身可能是热门商品比如“买手机的都会买充电线”这只是流行度在起作用。提升度大于 1.5 的规则更值得保留。5.4 多路召回融合加权分数排序融合时我把第 3 章的内容相似度、第 4 章的 ALS 预测分、第 5 章的关联规则分放在同一张候选表里召回通道分数含义权重内容相似度商品属性相似度0.2ALS 协同过滤预估偏好评分0.5FP-Growth 关联规则置信度 × 提升度0.3流式关注度当前热度0.1最终得分是加权和但要注意先做分数标准化。我采用的是基于排名的 RRFReciprocal Rank Fusion对每路召回的排序位置求倒数score sum(1/(60 rank))。这样不同档次的原始分数不会互相压制只要一路召回能给出合理的 top K都会被融合器保留下来。如果你更追求可解释性用线性加权也可以但务必先对分数做 min-max 缩放。6. 集群环境搭建与流式任务调试技巧6.1 最小可运行的 Spark 集群配置本地开发不需要一开始就上企业集群。我发现用 3 台 4 核 16G 的普通服务器就能跑通这个项目一台跑 NameNode ResourceManager另两台跑 DataNode NodeManager。如果只是学习完全可以单机 standalone 模式资源限制写在 Spark 配置里避免 OOM。spark-submit \ --master local[4] \ --conf spark.executor.memory4g \ --conf spark.sql.shuffle.partitions16 \ --py-files deps.zip \ main.pylocal[4]表示本机开 4 个线程模拟 executor适合调试。提交到 YARN 时要把--master改成yarn同时确认每台节点都安装了 PySpark且spark-defaults.conf的spark.yarn.archive指向 Spark 的归档包。6.2 流式任务的背压与状态膨胀排查Structured Streaming 跑久了最常遇到的问题有两个。第一个是背压问题。Kafka 消费速度跟不上生产速度时Streaming 内存里积压的数据会一直膨胀。解决办法是给 Kafka 消费者配maxOffsetsPerTrigger10000再加上spark.streaming.kafka.maxRatePerPartition限流。这样即使业务流量翻倍任务也不会被拉垮。第二个是状态膨胀。window 聚合会把窗口中间结果放进状态存储默认状态 TTL 越长占用的 RocksDB 空间越大。项目里我会手动观察 Streaming UI 的stateStore指标如果状态大小超过内存的 40%就把 watermark 从 2 分钟改到 1 分钟或者把窗口步长从 5 分钟改成 10 分钟允许更大的输出间隔从而减少 checkpoint 的写入频率。6.3 验证推荐结果时的三个检查点推荐系统上线前我一般先跑三个检查而不是只看离线指标。第一检查流式关注度有没有产出延迟。直接看 Redis 里attention_score的更新时间如果某个商品刚被大量点击score 应在 30 秒内变化。第二检查 ALS 预测结果里是否出现“用户根本不可能买的东西”比如给男性用户推荐卫生巾这通常是训练数据里类目维度没有过滤干净。第三检查关联规则是否有闭环“手机 → 手机壳 → 手机膜”这种规则看起来合理但连续推三个强关联商品会让用户觉得系统在重复推荐。我会加一个多样性惩罚同一条叶子类目下的商品最多出现两个。这个项目的源码和文档我在调试时遇到过不少坑比如 Scala 2.11 与 Spark 2.4 的兼容问题、Kafka 2.2 客户端与旧 broker 的协议不匹配这些在提供的环境配置里都有说明。拿到代码后先把pom.xml或requirements.txt里的版本号对齐到本地环境再按文档启动你会发现整套流程从 Kafka 到 Redis 再到推荐结果逻辑比想象中完整得多。本文还有配套的精品资源点击获取
返回列表