简介:基于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:45+08: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 声明,避免后面聚合时因类型推断失败而断流。startingOffsets=latest表示只消费新数据,调试阶段如果想重放历史消息,改成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做 key,score 存给在线服务读取。
3. 商品推荐召回:基于内容的相似度计算与候选生成
3.1 召回层为什么先用内容相似度
推荐系统里协同过滤在冷启动和长尾场景经常失效:新商品没行为、新用户没历史。基于内容的召回只依赖商品自身属性,不依赖交互记录,所以它是整个推荐链路里最稳的一路候选来源。这个项目里,我把商品标题、类目、品牌、颜色、价格区间拼成一个文本特征,算了相似度,作为协同过滤前面的粗筛。
商品属性在 MySQL 里,可以用spark.read.jdbc拉取。结构大致是:
| 字段 | 示例 |
|---|---|
| product_id | p_2034 |
| title | Redmi K70 手机 12GB+256GB |
| category | 智能手机 |
| brand | Redmi |
| price | 2499 |
| attributes | 白色;5000mAh;5G;OLED |
attributes是用分号拼接的多值字段,处理时按分号切分再展开成多行。
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(inputCol="text", outputCol="words") hashing_tf = HashingTF(inputCol="words", outputCol="rawFeatures", numFeatures=2048) idf = IDF(inputCol="rawFeatures", outputCol="features") pipeline = Pipeline(stages=[tokenizer, hashing_tf, idf]) model = pipeline.fit(df) feature_df = model.transform(df).select("product_id", "features")numFeatures=2048是经验值。太小容易碰撞,编码 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( userCol="user_id_int", itemCol="product_id_int", ratingCol="rating", rank=20, maxIter=15, regParam=0.1, coldStartStrategy="drop" ) model = als.fit(train_df)参数说明:
rank=20:隐因子数量。用户行为很稀疏时 10 到 20 足够;数据量大可以试 50,但模型体积和推理时间都会增加。maxIter=15:优化迭代次数。超过 20 次对 RMSE 的改善很少,反而容易过拟合。regParam=0.1:正则化系数。调参时优先看 0.01、0.05、0.1 三档。coldStartStrategy="drop":遇到训练集里没见过的新用户/商品时,预测结果会是 null,drop 可以直接把它从评估中剔除。如果线上的推荐服务无法处理 null,可以改成"fill"并指定coldStartStrategy补一个中性分数。
训练前,需要把 user_id 和 product_id 转成连续的数值 ID,否则 ALS 会报userCol不支持字符串。一种简单做法是用StringIndexer:
from pyspark.ml.feature import StringIndexer uid_indexer = StringIndexer(inputCol="user_id", outputCol="user_id_int") pid_indexer = StringIndexer(inputCol="product_id", outputCol="product_id_int")4.3 模型评估:RMSE 与业务命中率
训练集和测试集我用randomSplit([0.8, 0.2])切分。评估指标除了 RMSE,还应该加一个 top-K 召回率:测试集里用户真正交互过的商品,在模型给出的 top 20 推荐里出现了多少。
evaluator = RegressionEvaluator( metricName="rmse", labelCol="rating", predictionCol="prediction" ) rmse = evaluator.evaluate(predictions) print(f"RMSE = {rmse:.3f}")RMSE 在 1.0 以内算正常。如果跑到 2.0 以上,先检查评分分布是不是有极端值,再看regParam是否过小。业务上更关注的其实是排序结果,而不是预测分绝对误差。
4.4 冷启动用户的后备策略
ALS 给新用户预测时只能得到 null。这个项目里的做法是:新用户用“热门榜 + 内容召回”兜底。热门榜直接取第 2 章流式计算产出的attention_score前 50 个商品。把coldStartStrategy="drop"的预测结果和兜底结果做一个 union,再做权重合并。这样用户没有历史行为时也能收到合理推荐,等行为积累到 5 条以上再切换成个性化排序。
5. FP-Growth 关联规则挖掘与多路召回融合
5.1 为什么选 FP-Growth 而不是 Apriori
Apriori 的经典问题是要反复扫描全量数据生成候选集,商品量一旦上来,性能会断崖式下跌。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( minSupport=0.003, minConfidence=0.4, itemsCol="products", predictionCol="rules" ) fp_model = fp_growth.fit(transaction_df) fp_model.freqItemsets.show(10) fp_model.associationRules.show(10)参数说明:
minSupport=0.003:项集在事务中出现的频率下限。事务量 10 万时,0.003 意味着至少出现 300 次。支持度太高会丢失有价值的尾部关联。minConfidence=0.4:规则置信度。取 0.4 意味着,“买了 A 的用户里 40% 会买 B”,这条规则才被保留。predictionCol="rules":Fitting 完成后,模型会对每个事务预测可能关联的商品,这就是离线生成的补充推荐。
关联规则生成后,需要再用提升度(lift)过滤一次。置信度高不一定代表真正的关联,因为 B 本身可能是热门商品,比如“买手机的都会买充电线”,这只是流行度在起作用。提升度大于 1.5 的规则更值得保留。
5.4 多路召回融合:加权分数排序
融合时,我把第 3 章的内容相似度、第 4 章的 ALS 预测分、第 5 章的关联规则分放在同一张候选表里:
| 召回通道 | 分数含义 | 权重 |
|---|---|---|
| 内容相似度 | 商品属性相似度 | 0.2 |
| ALS 协同过滤 | 预估偏好评分 | 0.5 |
| FP-Growth 关联规则 | 置信度 × 提升度 | 0.3 |
| 流式关注度 | 当前热度 | 0.1 |
最终得分是加权和,但要注意先做分数标准化。我采用的是基于排名的 RRF(Reciprocal 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.memory=4g \ --conf spark.sql.shuffle.partitions=16 \ --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 消费者配maxOffsetsPerTrigger=10000,再加上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 再到推荐结果,逻辑比想象中完整得多。
本文还有配套的精品资源,点击获取