
简介基于Spark的电商商品智能分析系统毕业设计源码包面向大数据、软件工程相关专业学生亦适合对实时计算与推荐系统感兴趣的开发者。系统以Spark Streaming接收并处理用户浏览、点击等实时行为数据结合注意力模型实时计算商品关注度通过协同过滤、基于内容或深度学习等算法实现智能推荐并利用Apriori/FP-Growth做关联规则挖掘提升推荐准确性与时效性。压缩包共939个文件约5.49MB涵盖Java/Scala源码、XML与Properties配置、HTML/JS/CSS前端展示、147个class编译文件及Spark计算输出part等源码结构包含数据处理、模型训练、结果输出等模块便于定位所需代码与数据。目前已有246人学习下载可作为毕业设计/课程设计的完整参考也可用于理解流式计算、推荐算法与关联分析在真实电商场景中的项目落地下载后根据文档配置Hadoop与Spark环境即可运行。1. 基于 Spark 的电商商品智能分析解决的不只是「算得快」做电商数据的人大概都经历过这个场景运营晚上八点来问「现在哪个商品最热」你只能回「明天早上出报表」。用户行为日志是实时产生的但计算是隔天的等榜单出来热点早过去了。标题里这套系统——基于 Spark 的电商商品智能分析系统——就是把流式计算、智能推荐和关联分析放到一条流水线上Structured Streaming 实时算商品关注度ALS 协同过滤做个性化召回FP-Growth 挖商品组合做交叉推荐。它适合三类人电商或零售领域的数据分析师准备大数据岗位面试的开发者以及需要一套完整项目做毕设或简历项目的学生。这套方案的现实意义在于不需要专门去搭一套 Flink 集群就能把关注度计算、推荐、关联分析三件事一次跑通中小体量的电商完全够用。2. 整体链路怎么搭从行为日志到关注度、推荐、关联分析的落点很多拿到这个项目的人第一步就跑代码结果在「业务口径」上翻车——不知道关注度怎么定义、推荐哪部分是实时的、关联分析跟推荐之间到底是什么关系。这套系统拆开看其实是四个模块数据接入、流式关注度计算、离线推荐、离线关联分析。先把链路画清楚后面每一章都是在往里填肉。2.1 输入电商行为日志长什么样怎么接入 Spark电商商品分析系统最直接的输入是埋点行为日志不是订单表。订单表只有下单那一刻的数据而「关注度」需要的是点击、收藏、加购这些过程信号。每条日志通常是一个 JSON字段不需要多核心就四个用户、商品、行为类型、时间戳。{ userId: 1001, itemId: sku_7788, behavior: click, categoryId: 手机数码, ts: 2025-06-01 10:23:45 }behavior 字段常见取值是 click、fav、cart、buy 四类有的系统会细分出 browseTime停留时长、addToWishlist但那是锦上添花。接入方式上我一般会给两条路本地演示用 Structured Streaming 的 file source直接监听一个日志目录生产环境走 Kafka把接入参数从目录路径换成 Kafka broker 地址。这样一套代码两种部署不用改逻辑。spark.readStream.format(json)就能直接读 JSON 格式的日志流Spark 会自动把每行 JSON 解析成一行数据。这里容易被忽略的是时间字段的格式如果ts是字符串必须在后续处理前把它转成 TimestampType否则window和watermark全部失效。这也是很多 spark 数据分析案例里数据清洗占了大头的原因——农产品价格、网约车轨迹这类项目清洗部分都远超计算部分电商行为日志也一样第一步先把时间洗干净。2.2 关注度指标怎么定PV、UV、加购、下单的加权口径「商品关注度」不是一个标准化指标不同业务口径完全不一样。最常见的做法是给行为事件赋权算一个加权分数score 0.2 * pv 0.3 * fav 0.4 * cart 0.5 * buy_cnt但权重要看品类冲动消费品类加购权重高高客单品类收藏和浏览时长更有意义。这个公式本身不需要多复杂真正关键的是给分数加一个时间窗口约束。是算过去 10 分钟的关注度还是过去 1 小时还是当天累计窗口口径决定了榜单的时效性和波动程度。窗口口径适用场景特点10 分钟滑动窗口大促实时热榜、秒杀监控波动大能抓住瞬时热点1 小时窗口10 分钟滑动日常运营趋势榜平滑且相对实时当日累计日报、周报基础数据稳定但实时性差还要分清 PV 和 UV。PV 是用户点了多少次UV 是去重后有多少人点。流式计算里做countDistinct(userId)是开销大户尤其用户基数大时一个窗口内要去重的 userId 可能有几百万。这个点后面避坑章会专门讲先记住结论能近似就不精确能用 HyperLogLog 就别countDistinct。2.3 存储与调度Redis、MySQL、HDFS 各管哪一段整套链路里数据落在不同层每层的选型逻辑不一样模块输入输出选型理由原始日志埋点 JSONHDFS / OSS供离线训练和回溯分析关注度结果流式聚合输出Redis ZSET热榜接口高并发读天然支持按分数排序ALS 训练样本HDFS 上的历史行为MySQL / Hive离线批量训练关联规则订单/加购明细MySQL规则量小支持运营查询checkpoint流式作业状态HDFS / OSS作业重启恢复必须持久化关注度结果写 Redis 而不是 MySQL原因有两个一是热榜接口的 QPS 很高Redis 扛得住二是 ZSET 这个数据结构天生就是干这个的zrevrange一条命令就能取出 Top 100MySQL 要写一条带 order by 的 SQL还要加索引性能差一个量级。2.4 为什么这套组合比 Flink 全家桶更适合中小电商选 Spark 而不是 Flink不是 Spark 比 Flink 强而是这套系统里真正实时的部分只有关注度这一个指标粒度是分钟级。Spark Structured Streaming 的吞吐量和分钟级延迟对日活百万、日志量两亿条以内的站点完全够用。更重要的是Spark 一套代码能同时覆盖离线训练和流式计算ALS 和 FP-Growth 本来就在离线跑用 Stream 算完关注度直接落到 Redis不需要维护两套集群。这是实际做项目时最该算清楚的一笔账——很多团队在这个体量上引入 Flink ClickHouse 向量数据库最后运维成本比业务收益还高。3. 用 Structured Streaming 实时计算商品关注度最小可跑链路流式计算是这套系统的地基。关注度算不出来推荐和关联分析都成了无源之水。这一章从一段能跑的 PySpark 代码开始把窗口、水位线、触发器这些参数讲透最后落到怎么把结果安全地写进 Redis。3.1 先跑通readStream 读 JSON window 聚合 console 输出我习惯先跑一条最简通路再逐步加复杂度。下面这段代码读取日志目录里的 JSON按 10 分钟窗口、5 分钟滑动统计每个商品的 PV、UV、加购数和下单数先输出到控制台验证链路。from pyspark.sql import SparkSession from pyspark.sql.functions import window, col, count, countDistinct, when, expr from pyspark.sql.types import StructType, StructField, StringType, LongType, TimestampType schema StructType([ StructField(userId, LongType()), StructField(itemId, StringType()), StructField(behavior, StringType()), StructField(categoryId, StringType()), StructField(ts, TimestampType()), ]) spark (SparkSession.builder .appName(ecommerce-attention-analysis) .master(local[2]) .config(spark.sql.shuffle.partitions, 4) .getOrCreate()) raw (spark.readStream .format(json) .schema(schema) .option(cleanSource, archive) .load(logs/)) attention (raw .withWatermark(ts, 10 minutes) .groupBy( col(itemId), window(col(ts), 10 minutes, 5 minutes) ) .agg( count(col(userId)).alias(pv), countDistinct(col(userId)).alias(uv), count(when(col(behavior) cart, 1)).alias(cart_cnt), count(when(col(behavior) buy, 1)).alias(buy_cnt) ) .withColumn(score, expr(0.3 * pv 0.4 * cart_cnt 0.5 * buy_cnt)) ) query (attention .writeStream .outputMode(append) .format(console) .trigger(processingTime1 minute) .option(truncate, false) .start()) query.awaitTermination()这段逻辑拆开说readStream是流式读取入口format(json)处理 JSON 文件流schema提前定义好字段类型尤其是ts必须是TimestampType否则后面的窗口和水位线都算不出来。withWatermark(ts, 10 minutes)告诉 Spark 允许事件迟到 10 分钟超过这个时间的数据会被丢弃。groupBy( col(itemId), window(...) )是窗口聚合的核心写法注意窗口函数返回的是一个结构体列包含 window_start 和 window_end 两个字段。count(when(...))是条件计数PySpark 里没有countIf这是最常见的替代写法。outputMode(append)在带窗口的聚合中表示窗口结束后输出最终结果不会更新旧窗口正好适合关注度榜单这种场景。3.2 窗口、水位线和触发器参数不是抄作业是按延迟分布调的很多人在这一步把 watermark 当玄学调——随便设个 10 分钟结果移动端慢网络下的事件大量迟到被丢弃关注度数字偏得一塌糊涂。这三个参数的设置逻辑完全不同window的窗口长度和滑动间隔决定榜单粒度想做大促实时监控就用短窗口日常趋势榜用 1 小时窗口 10 分钟滑动。watermark决定允许事件迟到多久应该根据日志事件时间与到达时间的延迟分布来设取 P90 或 P95 的分位数而不是拍脑袋。trigger(processingTime...)决定 Spark 多久拉取一次数据。注意它不是窗口时长的替代品在 file source 下它只是目录扫描间隔。水位线还有一个隐蔽约束必须和窗口聚合配合使用且withWatermark必须写在groupBy之前。如果代码顺序写反了Spark 会直接抛异常不会给你任何回旋余地。触发器的设置上我一般用processingTime1 minute起步。设太短比如 5 秒在日志量小的场景下反而浪费资源每个批次可能只有几百条数据白白调度一次。设太长比如 10 分钟会让关注度结果看起来非常迟钝运营那边接受不了。3.3 把关注度写到 RedisforeachBatch 与幂等键控制台验证通过后下一步把结果落到 Redis。writeStream的输出端里foreachBatch最灵活——它把每个微批的 DataFrame 交给一个自定义函数处理我们可以在里面做任意写操作。import redis r redis.Redis(hostredis-host, port6379, db0, decode_responsesTrue) def write_attention_to_redis(batch_df, batch_id): if batch_df.isEmpty(): return rows batch_df.select(window_start, itemId, score).collect() pipe r.pipeline(transactionTrue) keys set() for row in rows: key fattn:{row[window_start]} pipe.zadd(key, {row[itemId]: float(row[score])}) keys.add(key) for key in keys: pipe.expire(key, 6 * 3600) pipe.execute() query (attention .writeStream .outputMode(append) .foreachBatch(write_attention_to_redis) .option(checkpointLocation, hdfs://namenode:8020/ckpt/attention) .trigger(processingTime1 minute) .start()) query.awaitTermination()foreachBatch里每一批只处理本批窗口的数据Redis 的 key 用window_start区分value 用商品 ID 做 member、关注度分数做 score。zadd天然支持按分数排序expire设置 6 小时过期避免历史窗口的 key 无限堆积。这段代码有两个容易翻车的细节一是collect()把整个批次拉回 Driver微批行数小没问题但日志量大时这种做法会把 Driver 内存打爆正确做法是用mapPartitions在 Executor 侧写 Redis或者用foreachPartition二是幂等性zadd是覆盖写不是累加所以任务重放不会导致分数翻倍——如果写成incrbycheckpoint 恢复时就会重复计数这是所有流式写 Redis 最常见的坑。4. 智能推荐怎么做ALS 离线召回 关注度在线融合关注度解决的是「现在什么火」推荐解决的是「这个用户可能喜欢什么」。两个不能互相替代热度榜对所有人一样推荐必须个性化。但只做 ALS 又会冷启动失灵、对突发热点反应慢。常见做法是离线 ALS 召回 在线关注度加权融合。4.1 离线召回行为日志转评分矩阵用 ALS 训练ALS交替最小二乘法是 Spark MLlib 里最成熟的协同过滤算法把用户对商品的隐式偏好分解成两个低维矩阵。第一步是把行为数据转成评分。行为本身不是评分需要加权换算from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, date_sub, current_date from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator spark SparkSession.builder.appName(als-recommendation).getOrCreate() behavior spark.sql( SELECT userId, itemId, behavior, ts FROM behavior_log WHERE ts date_sub(current_date(), 30) ) ratings (behavior .withColumn(raw_rating, when(col(behavior) click, 1.0) .when(col(behavior) fav, 1.5) .when(col(behavior) cart, 2.0) .otherwise(3.0)) # 时间衰减越近的行为权重越高 .withColumn(rating, col(raw_rating) * (1 - datediff(current_date(), col(ts)) / 30)) .select(userId, itemId, rating)) (training, test) ratings.randomSplit([0.8, 0.2], seed42) als ALS( userColuserId, itemColitemId, ratingColrating, rank12, maxIter10, regParam0.08, coldStartStrategydrop ) model als.fit(training) predictions model.transform(test).na.drop() rmse RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ).evaluate(predictions) print(fRMSE: {rmse:.4f})raw_rating的权重映射来自电商常识下单的意图强度远大于点击加购介于两者之间收藏略低于加购。时间衰减是必要的——30 天前的一次加购对今天的推荐意义已经很小了。这里用线性衰减做近似不需要更复杂的指数衰减。rank是隐因子数量一般从 8 到 20 之间试。rank 太小欠拟合用户和商品的特征表达不够rank 太大过拟合训练时间翻倍但推荐效果不再提升。regParam是正则化系数0.01 到 0.1 之间调先设 0.08 看 RMSE 趋势。coldStartStrategydrop表示测试集中出现训练时没见过的商品时直接丢弃该预测否则预测结果会出现空值RMSE 计算直接报错。4.2 在线融合ALS 预测分 关注度分数的加权策略ALS 模型输出的预测分是用户对商品偏好程度的估值但两个问题解决不了冷启动用户没有足够历史行为、新上架商品没有交互记录。而关注度恰好能补上——新商品如果突然有大量点击流式计算几分钟内就能把它推到热度榜前列。我把两者做成加权融合# 从 Redis 取出当前窗口热度 Top 200 hot_scores r.zrevrange(attn:2025-06-01-10:00:00, 0, 199, withscoresTrue) hot_map {item_id: float(score) for item_id, score in hot_scores} # 对单个用户的 ALS 候选集打分 def rerank(user_cands): # user_cands: [(itemId, als_pred), ...] result [] for item_id, als_pred in user_cands: als_norm 1.0 / (1.0 abs(als_pred)) # 简单 min-max 归一化替代 hot_score hot_map.get(item_id, 0.0) hot_norm hot_score / max_hot_score # 热度归一化到 0-1 final 0.6 * als_norm 0.3 * hot_norm 0.1 * global_ctr.get(item_id, 0.01) result.append((item_id, final)) return sorted(result, keylambda x: x[1], reverseTrue)[:50]ALS 预测分和关注度分数不是一个量纲直接相加没有意义。所以先归一化ALS 分数用1/(1|x|)压到 0-1热度分数除以当前窗口最大热度得到 0-1 的相对值。权重上个性化占 60%热度占 30%商品平均点击率占 10%。这个比例不是算出来的是拍板后拿 A/B 测试调的——先按这个跑观察推荐位点击率和转化率再逐步调。4.3 冷启动没有交互记录的商品和用户怎么办冷启动是推荐系统最现实的问题。新商品没有评分ALS 根本不会把它放进任何人的候选集。三个补救手段按工程成本从低到高排列第一用类目热门替代。新商品所在类目下当前关注度最高的 Top 50 直接作为它的临时召回结果等有了真实交互再交给 ALS。这个商品还在热度榜上所以相当于给新商品一个「初始流量」。第二用内容相似度做 item2item。商品标题和类目标签是现成的文本特征用 TF-IDF 算相似商品新商品可以借用相似商品的协同过滤结果。第三运营规则兜底。关联分析的结果在这里就能用上——买过 A 的用户把关联规则里与 A 搭配的高置信度商品 B 推给新用户这是不需要用户历史行为也能成立的推荐逻辑。5. 关联分析怎么做FP-Growth 挖商品组合反哺推荐关联分析在电商里最典型的场景是「买了 A 的用户还会买 B」——啤酒和尿布、手机和贴膜。Spark MLlib 内置了 FP-Growth 算法比 Apriori 在内存和速度上都快得多不需要反复扫描全量数据。这一章从构造成组样本开始到规则过滤最后落到怎么跟推荐联动。5.1 从行为日志构造成组样本session 聚合FP-Growth 的输入是「交易集合」——每一行是一笔订单或一次加购会话包含的商品列表。直接从订单表拿数据当然可以但订单表覆盖不到加购未下单的场景信息损失太大。常见做法是把点击日志里的加购和下单行为按用户和日期聚合成一个商品集合SELECT user_id, date_format(ts, yyyy-MM-dd) AS dt, collect_set(item_id) AS items FROM behavior_log WHERE behavior IN (cart, buy) GROUP BY user_id, date_format(ts, yyyy-MM-dd) HAVING size(items) 2collect_set自动去重同一商品在一天内被加购三次只会出现在集合里一次。HAVING size(items) 2过滤掉只有一个商品的集合——单个商品构不成关联规则。这里用「用户 天」近似一个购物会话严格一点应该按「用户 30 分钟无行为间隔」切分 session那是另一个开窗函数的活但做项目时「用户 天」的粒度通常已经能挖出有效规则。5.2 训练 FP-Growth 与规则过滤素材准备好后直接喂给FPGrowthfrom pyspark.ml.fpm import FPGrowth fp_growth FPGrowth( itemsColitems, minSupport0.01, minConfidence0.2, maxPatternLength5 ) model fp_growth.fit(transactions_df) # 频繁项集 model.freqItemsets.show(10) # 关联规则 rules model.associationRules rules.filter(lift 1.0 AND confidence 0.2).orderBy(lift, ascendingFalse).show(10)minSupport0.01表示商品组合至少出现在 1% 的交易里。数据量大就往上调数据稀疏就往下调。聚合出的交易总条数只有几万条时我一般降到 0.005 否则一条规则都挖不出来。minConfidence0.2表示买了 A 的用户至少有 20% 会买 B。maxPatternLength5限制规则最多包含 5 个商品防止出现超长组合。associationRules输出四列antecedent是前件买了什么consequent是后件还会买什么confidence是条件概率lift是提升度。lift 是这三兄弟里最重要的一个——lift 大于 1 说明 A 对 B 有正向影响等于 1 说明两者独立小于 1 说明 A 反而抑制 B。很多人只看 confidence 就上规则结果挖出一堆「买手机的人都买手机壳」的废话规则confidence 高达 80% 但 lift 接近 1没有任何增量价值。过滤条件落到实处就两条lift 1.0 保正向关联confidence 0.2 保证规则可用性。5.3 关联结果怎么和推荐联动关联规则不直接推到前端做展示它有三个更合理的去处第一购物车实时弹窗。用户在购物车加入 A 时从规则表查 A 的关联商品 B、C在结算页推荐「搭配购」这是电商转化率最高的玩法。第二推荐候选补充。ALS 候选集只包含用户历史偏好的相似商品可能漏掉「买了相机还会买 SD 卡」这种跨类目但有强因果的组合。把高 lift 规则里的 consequent 直接插入推荐候选集相当于给推荐加了一条「可解释的关联路径」。第三运营手动配置活动。规则表导出给运营做捆绑促销和品类陈列参考。这里特别提醒一个常见的错误规则表不要全量写 Redis。几百上千条规则确实不多但大多数是无用或者弱相关的。我会把规则按 lift 降序截断到前 100~200 条同时要求 consequent 的商品在近 7 天有最低曝光量否则冷门商品之间的虚高关联会污染推荐结果。6. 避坑与验证这套系统里最容易翻车的六个点6.1 checkpoint 放本地 /tmp重启丢状态现象任务重启后从 Kafka 重放全部数据Redis 里的关注度分数没有翻倍但延迟严重。原因checkpoint 默认目录写在本地/tmp被系统清理或换了机器后流式状态全部丢失。解决spark-submit时加--conf spark.sql.streaming.checkpointLocationhdfs://.../ckpt放进 HDFS 或云上持久目录。6.2 数据倾斜一个爆款 sku 拖垮整个批次现象窗口聚合阶段某个任务跑了 50 分钟没结束Spark UI 上看单个 Executor 的 GC 时间飙升。原因groupBy(itemId)时爆款商品的数据量远大于普通商品单 key 倾斜。解决两阶段聚合先按(itemId, rand(10))加盐聚合一次再去掉盐汇总极端情况下对爆款走单独的近似去重通道。6.3 watermark 设太短晚到流量全被丢弃现象关注度数字比后台离线报表低两成而且总是那些慢网络场景下的移动端流量对不上。原因withWatermark(ts, 5 minutes)设太短实际事件延迟达 20 分钟。解决先跑一天日志统计事件时间与处理时间的延迟分布把 watermark 设到 P95 分位比如 20 分钟。时效性和准确度在这里必须做取舍。6.4 foreachBatch 写 Redis 用累加重试后翻倍现象从 checkpoint 恢复后 Redis 里的关注度分数整体接近翻倍。原因foreachBatch里用了incrby累加微批重放导致重复计数。解决改成zadd覆盖写幂等键包含window_startRedis key 加过期时间。这个坑几乎每个流式写 Redis 的人都会踩一次。6.5 ALS 默认隐式反馈参数推荐结果全是爆款现象所有用户拿到的推荐列表一模一样都是全站热门。原因implicitPrefsTrue时alpha默认 1.0隐式反馈的置信度权重没起来模型坍缩到热门商品。解决alpha调到 20~40 再训练或者像前面那样按显式评分构造新增用户和冷门商品至少还有挽回余地。6.6 关联规则 lift 虚高冷门组合被当黄金规则现象置信度 100%、lift 高达 8 的规则细看是某个冷门商品只被同一个人买过两次。原因样本量太小统计显著性不足。解决规则过滤加最小出现次数硬门槛比如 consequent 至少出现在 50 个交易里或者用 Laplace 平滑收缩置信度。这套系统跑通不难跑好全是这些边角。我现在的习惯是每张关键结果表都留一份输出行数与时间范围的快照上线前先用 Spark UI 的 Streaming 页盯 Input Rate 和 Scheduling Delay 两小时。Scheduling Delay 不断上涨说明资源给少了不是代码有 bug——别把 Spark 当黑匣子每一层落下来都有数可对。希望帮到你。本文还有配套的精品资源点击获取