简介:这是一套面向毕业设计或课程设计场景的Spark电商商品智能分析系统源码,适合想掌握流式计算、推荐算法与关联规则分析的大数据专业学生和入门开发者。项目基于Spark Streaming采集用户浏览、点击等实时行为数据,完成商品关注度计算,并集成协同过滤、基于内容推荐及FP-Growth等关联规则挖掘,覆盖数据清洗、模型训练、结果输出等完整模块。压缩包共939个文件,大小5.49MB,主要包含Java、Scala、Spark Streaming中间结果part文件、checkpoint数据、HTML/JS前端页面及配置文件等,目录结构清晰,便于按模块查看训练日志与运行结果。项目已吸引246人学习下载,适合用于课设答辩、毕业设计演示或Spark实战入门参考。
1. 用 Spark 把电商关注度算清楚,推荐和关联分析才站得住
“基于spark的电商商品智能分析系统”这个名字里,真正承重的是中间那截“流式计算电商商品关注度”。关注度算得不准,后面的智能推荐、关联分析都只能算赶时髦。很多初次搭这套系统的工程师会把 PV、UV 直接当关注度,发现推荐结果越跑越偏:用户刚加购的商品,过十分钟热度还在涨,等推荐系统反应过来,用户早下单走了。电商关注度必须是一个带时间属性的滑动窗口值,不是数据库里一个累加计数器。
这篇文章按这套系统的完整链路来拆:先讲为什么 Spark 适合同时承载流式计算、推荐和关联分析,再各用一章把关注度窗口计算、ALS 推荐召回、FP-Growth 关联规则做成可复现的代码,最后一章落到内存参数、数据倾斜和 Spark UI 排错。适合两类人看:一类是拿这套系统做数据平台初始化,另一类是已经跑通任务但发现结果不对、想搞清楚参数怎么调的开发。
2. 电商实时分析系统为什么用 Spark:选型逻辑与架构分层
2.1 一套引擎同时处理实时热度、离线推荐与关联规则
电商商品智能分析系统的数据链路很典型:行为日志从 Nginx 或 App 端进入 Kafka,一部分走实时计算,一部分落到数据湖做离线训练。如果实时和离线分别维护两套技术栈,比如 Flink 做流、Spark 做批,代码至少有 30% 是重复的解析逻辑,小团队很难维护。这个项目标题把“流式计算、智能推荐、关联分析”三个能力放在一起,用 Spark 统一承载是最省力的架构。
选择 Spark 的理由有三层。第一,Structured Streaming 的 DataFrame API 与离线 DataFrame API 同构,白天跑实时关注度,凌晨跑同一个逻辑的离线全量重算,代码差异只在水位线和窗口参数上。第二,MLlib 直接提供 ALS(协同过滤)和 FP-Growth(关联规则),不需要额外引入独立的推荐引擎或规则引擎。第三,Spark 的部署生态足够成熟,Yarn、K8s 都能调度,前面接 Kafka,后面接 Redis 或 MySQL,运维边界清晰。
2.2 Kafka 到 Redis 再到推荐服务的数据流向
整个系统按功能可以切成四层,理解了这四层,后面读源码时才能分清哪个模块在干什么。
第一层是数据接入层,客户端把商品浏览、收藏、加购、下单四个动作统一上报为一条 JSON,写入 Kafka 的behavior_topic。第二层是流式计算层,Spark Structured Streaming 以KafkaSource方式消费,计算出商品在滑动窗口内的关注度分数,写入 Redis 的 Sorted Set 和 MySQL 的item_hot_rank表。第三层是推荐层,ALS 模型离线训练用户和商品向量,实时模块把用户最近的点击行为拼接成偏好信号,从 Redis 中召回候选商品。第四层是分析层,FP-Growth 每日对订单明细跑一次关联规则,输出“买了 A 的人还会买 B”的关系表,推荐接口兜底时直接读这张表。
这样分层的好处是:关注度计算和推荐逻辑解耦。推荐服务只消费关注度结果,不感知 Spark 作业内部的窗口机制;Spark 作业也不关心推荐怎么过滤和排序。
2.3 行为事件与存储模型设计
在建流处理逻辑前,先把数据模型定下来。Kafka 里的行为消息统一用如下 JSON 结构,字段宁可多不要少,因为后面推荐冷启动大概率需要补字段:
{ "userId": "u_10032", "itemId": "p_88231", "behavior": "cart", "ts": 1717228800123, "channel": "app_home", "sessionId": "s_99821" }| 字段 | 类型 | 用途 |
|---|---|---|
| userId | String | 用户标识,推荐与去重依赖 |
| itemId | String | 商品标识,关注度分组键 |
| behavior | String | view / cart / favorite / order |
| ts | Long | 事件时间毫秒,窗口计算的依据 |
| channel | String | 流量来源,排查异常流量 |
| sessionId | String | 加盐去重和会话级关注度参考 |
Redis 侧建议用三个键分别存不同时效的数据:hot:rank:10min存近期热度 Top200,rec:user:{userId}存用户个性化推荐结果,rel:item:{itemId}存关联商品列表。实时链路写入频率高,不能把 MySQL 当主存储,MySQL 只做分钟级刷盘和报表查询。
2.4 Structured Streaming 消费 Kafka 的最简起点
下面是不带业务逻辑的最小读取段,先把数据接进来再看窗口聚合:
import org.apache.spark.sql.types._ val schema = StructType(Array( StructField("userId", StringType), StructField("itemId", StringType), StructField("behavior", StringType), StructField("ts", LongType), StructField("channel", StringType), StructField("sessionId", StringType) )) val rawStream = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "kafka1:9092,kafka2:9092") .option("subscribe", "behavior_topic") .option("startingOffsets", "latest") .option("maxOffsetsPerTrigger", 100000) .load() .selectExpr("CAST(value AS STRING) as json") .select(from_json(col("json"), schema).as("data")) .select("data.*")startingOffsets只在首次启动时决定从哪开始读,任务重启后优先以 checkpoint 为准,这里设latest是为了避免测试环境一启动就消费大量历史日志把 Redis 打满。maxOffsetsPerTrigger建议生产环境根据单条消息大小调整,每条 300 字节左右时 10 万条基本对应每秒 5 万到 10 万事件,是单作业比较稳的起点。真正执行窗口聚合前,把这段代码跑通并观察RateController的背压行为,比先写复杂聚合要省时间。
3. Spark 流式计算商品关注度:窗口设置、去重与热度分输出
3.1 关注度不是 PV 累加,是加权事件分数
把“关注”映射成某个商品的实时热度,业内常见的做法是给不同行为分配不同权重,再放到同一个窗口里做加权累计。简单的公式如下:
score = 1.0 * view + 0.4 * favorite + 1.5 * cart + 3.0 * order这里权重不是拍脑袋定的,而是根据转化率倒推。比如从浏览到加购的转化率是 5%,加购对成交的贡献大约是浏览的 15 到 20 倍,取 1.5 到 3 之间都合理。权重建议放到 Redis 或配置中心,别写死在代码里,因为运营活动期间加购权重往往要临时调高。
只算总分会有一个明显的坑:同一用户一小时刷新页面 50 次,会把商品热度顶上去,推荐系统随后会把商品推给这个用户本人,制造出“自己推自己”的循环。因此有效关注度必须对 user 去重,至少做到窗口内同一用户同一商品只计一次行为。
3.2 滑动窗口里的去重聚合与水位线
下面是一段可直接放进作业的窗口聚合核心逻辑:
import org.apache.spark.sql.functions._ val hotStream = rawStream .withWatermark("ts_ms", "2 minutes") .withColumn("ts_ms", (col("ts") / 1000).cast("timestamp")) .groupBy( window(col("ts_ms"), "10 minutes", "5 minutes"), col("itemId") ) .agg( sum(when(col("behavior") === "view", 1.0) .when(col("behavior") === "favorite", 0.4) .when(col("behavior") === "cart", 1.5) .otherwise(3.0)).as("score"), approx_count_distinct("userId").as("uv") )withWatermark("ts_ms", "2 minutes")表示容忍事件乱序 2 分钟,超过水位线的迟到数据会被丢弃。这个值的设置要看客户端上报链路:一般 App 端日志从上报到进入 Kafka 在秒级到分钟级,2 到 5 分钟足够;如果数据经过离线任务回填,水位线要拉到 10 分钟以上。窗口设 10 分钟、滑动步长 5 分钟,表示每 5 分钟产出一次最近 10 分钟的滚动静态结果,兼顾实时性和窗口重合带来的计算量。
approx_count_distinct用的是 HyperLogLog 近似算法,在 UV 量级很大时比countDistinct省资源,误差大约在 1% 以内。需要精确去重时再改回countDistinct,但要注意它会引入大量 shuffle,看场景取舍。
3.3 把窗口结果输出到 Redis Sorted Set
Structured Streaming 里做外部存储写入最稳的是foreachBatch,它把每个微批当成一个 DataFrame,可以下推过滤条件、复用已有的 Redis 客户端连接,还能处理“最后 5 分钟无数据”的场景。下面是写入方案:
import redis.clients.jedis.Jedis hotStream.writeStream .foreachBatch { (batchDF: DataFrame, batchId: Long) => batchDF .withColumn("hot_score", col("score") * log(col("uv") + 1)) .select("itemId", "hot_score") .collect() .foreach { row => val jedis: Jedis = RedisPool.getJedis() jedis.zadd("hot:rank:10min", row.getDouble(1), row.getString(0)) RedisPool.returnJedis(jedis) } } .outputMode("update") .trigger(Trigger.ProcessingTime("30 seconds")) .option("checkpointLocation", "/data/checkpoint/hot_rank") .start()hot_score = score * log(uv + 1)是一种压制单用户刷量的组合:score 体现行为深度,log 后的 UV 体现覆盖广度。单用户刷 50 次 view 只能把 score 提高 50,但 UV 不涨,log 项不变,总体分数不会爆炸。
outputMode("update")适合这个场景,因为 Redis 是用 ZSet 存储,每次只要更新变化商品的分数即可。ProcessingTime("30 seconds")控制微批节奏,窗口滑动是 5 分钟,输出频率没必要比滑动频率快太多,30 秒到 1 分钟比较合适,太快会频繁写 Redis,拖慢 Executor。
3.4 关注度计算中三个高频问题
第一个是乱序数据导致窗口结果偏低。表现是热门商品分数在每次输出时先涨后跌,跌的部分其实是迟到的浏览记录被水位线丢弃。排查方式是查看 Kafka 消费组的 lag,如果 lag 稳定增长,说明处理速度跟不上,优先调大maxOffsetsPerTrigger或增加 Executor 并行度。
第二个是 Redis 连接被频繁创建。绝不能在每个row.foreach里new Jedis(host)。用连接池统一管理,maxTotal设置成 100 左右,避免把 Redis 压出连接超时。
第三个是把窗口结果直接写 Kafka 再让下游消费。这种做法不是不行,但引入的延迟与 Redis 差不多,反而多维护一个消费端。内部系统走 Redis 直连更快,只有需要给多个团队共享数据时才写 Kafka。
这个模块的实时关注度数据出来之后,下一步就是让推荐系统用起来。如果用户最近的行为集中在某个商品上,能立刻反哺到用户画像,推荐才有“智能”可言。
4. 电商商品智能推荐:ALS 离线训练与实时频道的召回融合
4.1 推荐系统在这个架构里的位置
本系统的推荐链路主要由三个通道组成。第一通道是“实时环境”:用户刚看过某个商品详情页,Redis 里取关联商品列表,快速做同品类、同价格带过滤后展示。第二通道是“历史偏好”:用 ALS 产出的每个用户的 TopN 商品列表,适合冷启动和首页推荐。第三通道是“兜底”:用户行为数据太少时,直接取当前窗口关注度最高的商品。
如果只做其中一条通道,推荐效果都很差。只用协同过滤,用户今天第一次搜索“露营灯”的行为要第二天才能进模型,等推荐出来用户已经买完;只用实时关联,推荐结果会局限在看到商品的相似款,完全丢掉了用户的长期品类偏好。请记住这个三元融合机制,它是一个“大数据实时分析系统”能被称为“智能”的关键。
4.2 ALS 模型训练要点与参数设定
ALS 适合用户行为稀疏的电商场景,它不要求事先准备商品属性特征,只靠“用户—商品—行为”三元组就能训练。实际落地通常按如下方式操作:
import org.apache.spark.ml.recommendation.ALS val als = new ALS() .setMaxIter(10) .setRank(20) .setRegParam(0.1) .setUserCol("userId_index") .setItemCol("itemId_index") .setRatingCol("rating") .setColdStartStrategy("drop") val model = als.fit(trainingData) model.write.save("/models/als_model")下表给出参数与推荐值,每个参数的效果不止影响 auc,也直接影响在线响应速度:
| 参数 | 推荐值 | 说明 |
|---|---|---|
| rank | 20 或 50 | 向量维度,越大表达越细腻,但内存占用和在线相似度耗时上升 |
| maxIter | 10 | 迭代次数,5 到 15 之间,超过 15 一般收益微弱 |
| regParam | 0.1 | 正则化,过小容易过拟合,行为数据稀疏时建议 0.5 |
| implicitPrefs | true | 用隐式反馈时应设为 true,评分是 0/1 行为而不是显式打分 |
注意setColdStartStrategy("drop")是在训练时把冷启动用户直接丢掉,否则会在生成推荐时报 NullPointerException。类似地,userId_index和itemId_index不能用字符串,必须先做 StringIndexer 转换。
用隐式反馈时,rating 值可以用前面关注度计算产出的hot_score代替 0/1,这样同一套评分体系能直接被 ALS 训练使用,用户 5 分钟前加购的商品在模型里产生更高权重,比单纯 1 和 0 的排序更合理。
4.3 把实时行为拼进模型结果
ALS 模型每天或每周重新训练一次,但实时行为必须立刻生效。常见做法是模型推荐结果冷加载到 Redis,实时模块读取用户最近 30 分钟的行为并动态提权。
val recentBehavior = rawStream .filter(col("ts_ms") > currentTime - 30 * 60 * 1000) .groupBy("userId", "itemId") .agg(collect_list("behavior").as("behaviors")) recentBehavior.writeStream .foreachBatch { (df, batchId) => // 写入 Redis Hash, key: user_recent:{userId}, field: itemId, value: behaviors } .outputMode("update") .start()推荐接口拿到 ALS 离线生成的rec:user:{userId}TopN 列表后,把user_recent:{userId}里高频行为涉及的商品提到列表前三位。这就是“阿里那种很多人讲过的粗排提权”的实时化实现。不必为此引入专门的实时推荐引擎,Redis 的读写延迟已经满足接口要求。
4.4 推荐系统最容易踩的坑:回声效应和时间衰减
回声效应是最难排查的坑。用户点击了商品 A,系统实时推荐 A 的相似商品 A1,用户又点 A1,系统继续推 A 的相似商品,用户看到的内容越来越窄。规避方式是在推荐结果里强加品类多样性:从 ALS 结果中抽 20 个候选,按二级类目均匀取样,同一小类目最多占 40%。这个操作要在 Redis 写入阶段做好,不然后端接口拿到什么就推什么,很快就“猜你喜欢”变成“猜你一个品”。
另一个坑是时间衰减缺失。如果模型是上周训练的,它可能持续推季节款。正确做法是候选商品在 Redis 存储时带上时间戳,推荐时过滤掉发布时间超过 30 天的商品,再用实时关注度加权一次。热点商品因为窗口分数高,天然带着时效性,两类信号叠加后,推荐结果就不会显得陈旧。
5. 电商商品关联分析:FP-Growth 在 Spark 上的落地与调参
5.1 为什么用 FP-Growth 而不是 Apriori
标题里的“关联分析”落到工程实现,最标准的落点是商品共现关系。电商场景的订单数据有典型的“长尾”特征:热门商品出现次数多,长尾商品出现频次低。Apriori 需要反复扫描数据集生成候选项集,在千万级订单上基本跑不动;FP-Growth 只需要扫描两遍数据集,用 FP 树压缩频繁项集,两者在内存消耗和速度上有数量级差距。MLlib 自带FPGrowth实现,直接跑在 DataFrame 上,不需要额外引入算法库。
这里要区分一个概念:关联分析基于“同时购买”,基于“点击共现”做关联的置信度非常低。比如用户浏览了 50 个商品页,两两之间未必有强关联;而一个订单里的商品组合才有真正的“意图绑定”。所以下面的实现默认数据源是订单表,不是点击流表。
5.2 Spark 调用 FP-Growth 的完整代码与参数解释
import org.apache.spark.ml.fpm.FPGrowth val orders = spark.sql(""" SELECT order_id, collect_list(item_id) AS items FROM order_detail WHERE dt = '2024-05-20' GROUP BY order_id """) val fpGrowth = new FPGrowth() .setItemsCol("items") .setMinSupport(0.002) .setMinConfidence(0.2) .setNumPartitions(10) val model = fpGrowth.fit(orders) model.setPredictionCol("prediction") val rules = model.associationRules rules.write.mode("overwrite").saveAsTable("rule_item_rel")setMinSupport(0.002)的意思是某个商品组合至少要出现在 0.2% 的订单中。日订单量 100 万时,就是 2000 单。别把 support 设得太小,比如 0.0001,会出现大量噪声组合,比如“手机壳+婴儿湿巾”这种被同一用户凑单但毫无业务含义的组合。setMinConfidence(0.2)表示由 A 推出 B 的条件概率要大于 20%,这个阈值决定规则可靠性,20% 是综合考虑电商场景可接受的下限。
setNumPartitions(10)控制 FP 树构建的并行度。对于 1 亿订单数据量,默认值可能造成单个 Executor 上的 FP 树构建压力过高,建议设成 Executor 总数的 1 到 2 倍。模型生成的associationRules单独落一张 Hive 表,线上推荐服务每天凌晨读一次,缓存到 Redis。
5.3 关联规则如何反哺推荐通道
关联规则与协同过滤的推荐结果必须叠加使用。它们有个关键差异:协同过滤找“像你的人喜欢的”,关联规则找“像这个商品该搭的”,两者合起来才能应对“用户为买帐篷进了店”却需要同时看到“防潮垫”的情况。落地方式如下:
// 伪代码展示线上推荐服务融合逻辑 val relatedItems = redis.zrange("rel:item:" + currentItem, 0, 5) val cfItems = redis.lrange("rec:user:" + userId, 0, 20) val finalList = deduplicate(relatedItems ++ cfItems) .filterNot(blacklist.contains) .sortBy(item => redis.zscore("hot:rank:10min", item).getOrElse(0.0) + itemQualityScore(item))注意rel:item:{itemId}里存的是规则置信度高的商品列表,必须按置信度降序加入推荐流,否则会把“买了 A 的人买了 B”变成“看了 A 就必须看到 B”,对内容生态不一定是好事。
5.4 关联分析在实际项目里的三个坑
第一,别把“点击流共现”当“购买关联”。在电商场景,用户点击和购买的行为意图完全不同,FP-Growth 用点击共现训练,可能有 90% 的规则是“手机→路由器”→“手机→数据线”这种泛化关系,没有增量价值。
第二,严谨处理“同一订单的凑单行为”。比如“满 300-50”活动下用户买了很多不相关商品,它们会被归进一个订单里,产生虚假关联。做法是过滤订单金额过高或商品数大于等于 5 的订单,或者做价格带归一化:只在大促期间跑、但要降低关联结果权重。
第三,规则结果需要业务审核,不能直接上线。每一步关联规则都要输出 support、confidence、lift 三列指标,让运营筛选出可用的规则。只要关联结果和运营直觉严重不符,第一反应应该是检查数据源里是否有测试订单、活动购物车把商品强行凑单的脏数据,而不是先调低 support。
6. Spark 作业调优与数据倾斜排查:把三套逻辑装进一个稳定任务
6.1 先从内存参数说起:统一的 spark-submit 模板
同样的代码,给足 Executor 内存和吝啬地分配内存,运行效率相差可达三倍。流式作业和离线批处理对内存的要求不同,但下面这份模板可以当作起点。
spark-submit \ --class com.example.HotAnalysis \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.kryoserializer.buffer.max=256m \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.streaming.kafka.maxRatePerPartition=3000 \ --conf spark.memory.fraction=0.6 \ --conf spark.memory.storageFraction=0.3 \ app.jarSpark 内存分为执行内存和存储内存,动态占用关系由spark.memory.fraction(默认 0.6)控制:这个参数调大,缓存执行中间结果的可用内存增加,但留给 RDD 缓存和数据块的存储内存减少。spark.memory.storageFraction=0.3表示存储内存至少保留的比例,当任务同时有流计算和 ALS 模型加载时,能避免反复驱逐缓存导致重算。
spark.sql.shuffle.partitions=200控制所有 shuffle 阶段的默认分区数。注意它不随 Executor 数量自动调整,如果 20 个 Executor、每个 4 核,200 个分区算偏少,建议改成executor核数 × executor数量 × 2到×3,也就是 160 到 240 之间。分区太少,单任务处理数据量太大,GC 频繁;分区太多,shuffle 文件的寻址开销上涨。
6.2 数据倾斜会伪装成内存溢出
流式计算关注度时,少数爆款商品会在窗口内收到几十倍于普通商品的行为量,按itemId分组后,个别 Reduce 任务要处理的 key 远超其他任务,表现为某个 Executor OOM,而其他 Executor 只有 20% 的利用率。这种情况在 Spark UI 上非常显眼:某个 Stage 的大部分 Task 秒级完成,但最后一两个 Task 要跑十几分钟。
处理方式分两步。第一步是找到倾斜 key 的证据。用下面这段可以把每个 key 的数据量打印出来:
groupedDF .groupBy("itemId") .count() .orderBy(col("count").desc) .show(20)确认是少数爆款导致的倾斜后,第二步做加盐拆分:把倾斜的商品 ID 拆成多个后缀,分散到不同 Reduce 任务,再在结果层合并。电商场景还可以直接设置爆款隔离:把热门商品单独走一条简单聚合链路,剩下的走通用链路,最后用 union 合并。这个方案比加盐更直观,因为爆款商品在列表页本身就该有独立处理策略。
6.3 用 DataFrame.explain 和 Spark UI 验证调优效果
任务跑完后,不要只依赖吞吐数值做判断。一个稳定的检查顺序是:先在作业代码里对最核心的聚合调用explain("formatted"),观察执行计划是否为HashAggregate或Partial/ Final分区,如果出现SortAggregate,说明分区键没处理好,先优化 key 分布。再打开 Spark UI 看每个 Stage 的输入数据量和 Shuffle Read 大小,一般 Shuffle Read 超过 Executor 内存一半,就该考虑增大分区数或优化 join 策略。最后看 GC 时间占比,如果超过 10%,增大 executor-memory 或减小单 Executor 核数,让每个 Executor 的并发任务数降下来。
这套检查方法可以直接复用到一个验证技巧上:给每个窗口聚合后加一行filter(score > threshold)或limit,把 Top 商品先输出到日志,用tail -f观察是否每 5 分钟稳定更新。这个动作能把“作业在跑但数据没产出”的问题提前暴露,比事后看 Redis 数据要快得多。
本文还有配套的精品资源,点击获取