Java与Scala混编电商推荐系统:从日志清洗到ALS与ItemBasedCF实践
2026/9/13 16:28:18 网站建设 项目流程

简介:这是一份面向大数据与推荐系统学习者的电商项目资源,采用Java与Scala混合开发,聚焦推荐引擎在电商场景中的落地实现。项目覆盖协同过滤、ALS矩阵分解、基于物品的ItemBasedCF以及按区域热门商品统计等典型模块,适合有一定大数据基础、希望掌握推荐算法工程化表达的读者。压缩包共233个文件,包含26个Java源码、14个Scala源码、162个编译后的class文件,以及配置和项目描述文件,整体仅640KB,目录结构清晰,便于快速阅读和二次修改。资源已有118人学习,通过阅读代码可以直观了解协同过滤与交替最小二乘法在分布式环境下的写法,并掌握Java与Scala混编的项目组织方式,实现从数据预处理、模型训练到推荐结果输出的完整链路。对于正在做课程设计或搭建小型推荐系统的开发者,这份紧凑的代码包具有不错的参考价值。

1. 为什么推荐系统要混编 Java 和 Scala:从日志到推荐的完整链路

电商大数据项目里,Java 和 Scala 混编不是炫技,而是工程取舍。Java 在数据清洗、接口对接、离线分析上代码可控、团队易接手;Scala 在 Spark 生态里写机器学习几乎能节省一半代码量。这个项目从原始日志(LogInfo、AdLogInfo)出发,依次做黑名单过滤、ALS 与 ItemBasedCF 协同过滤、HotProductByArea 热门兜底,最终形成一套可运行的离线推荐链路。适合正在做大数据毕业设计、准备 Java/Scala 大数据面试,或者想把推荐系统工程化的读者。它的价值不在算法有多深,而在于把日志清洗、模型训练、推荐产出、部署验证这一整条链路串起来,每一环都能对照代码去复现。

2. 日志清洗与特征提取:LogInfo、AdLogInfo 到 Spark DataFrame

推荐系统的第一道工序不是建模,而是日志清洗。这里的核心类是 LogInfo 和 AdLogInfo,它们分别表示用户行为日志和广告日志。日志格式不统一、时区不一致、字段缺失,直接送进模型会导致离线指标失真。常见做法是先统一定义 schema,再用 Spark 读入并剔除无效记录。

2.1 电商日志的数据模型:LogInfo 和 AdLogInfo 字段设计

推荐链路需要一个统一的行为视图。LogInfo 一般包含用户 ID、物品 ID、行为类型(view/cart/pay)、行为时间、来源渠道;AdLogInfo 除了这些,还多出广告位 ID、素材 ID、曝光与点击标记。表结构可以按下面的模型收敛:

字段类型说明是否必须
user_idString用户唯一标识
item_idString物品/商品 ID
actionStringview/cart/pay
timestampLong行为发生时间(秒级)
area_idString地区 ID否,用于地区热门
ad_idString广告 ID否,仅 AdLogInfo
is_clickInt是否点击,0/1广告行为需要

黑名单用户(BlackUserList)一般单独存一份,在清洗阶段 join 掉。BlackUserList$.class这个类出现在项目里,说明设计者选择在入口处过滤刷单用户,而不是在模型层过滤。这是对的,因为刷单行为一旦进入训练集,会把协同过滤的相似度矩阵拉偏,后续再调权重都很难救回来。

2.2 用 Scala 把原始日志解析成 DataFrame

日志中的原始行通常是 tab 或逗号分隔。下面是一段可以直接跑的 Scala 解析逻辑,输出 Spark SQL 可查询的 DataFrame:

import org.apache.spark.sql.{SparkSession, Row} import org.apache.spark.sql.types._ val spark = SparkSession.builder() .appName("LogClean") .master("local[*]") .getOrCreate() val schema = StructType(Seq( StructField("user_id", StringType, nullable = false), StructField("item_id", StringType, nullable = false), StructField("action", StringType, nullable = false), StructField("timestamp", LongType, nullable = false), StructField("area_id", StringType, nullable = true) )) val rawRDD = spark.sparkContext.textFile("hdfs:///data/user_log/") .map(_.split("\\t")) // 过滤字段数量不够的记录,同时把时间戳转成 Long val cleanRDD = rawRDD.filter(_.length >= 4).mapPartitions { iter => iter.flatMap { arr => try { val ts = arr(3).trim.toLong val area = if (arr.length > 4) arr(4).trim else "unknown" Some(Row(arr(0), arr(1), arr(2), ts, area)) } catch { case _: NumberFormatException => None } } } val df = spark.createDataFrame(cleanRDD, schema) df.createTempView("user_log") // 删除黑名单:假设 blacklist 是一张表 val filtered = spark.sql( """ |SELECT l.* FROM user_log l |LEFT JOIN blacklist b ON l.user_id = b.user_id |WHERE b.user_id IS NULL """.stripMargin) filtered.show(5)

这段代码的逻辑分三步:按 tab 切开一行日志,逐字段检查格式,再与黑名单表做左连接过滤。使用mapPartitions而不是map,可以减少连接 Driver 和 Executor 的次数,适合每条记录都做独立解析的场景。flatMap内返回Option,时间戳解析失败时直接None,把脏数据静默丢弃,避免整个 Stage 报错。项目里同时出现LogInfo.classAdLogInfo.class,说明有两种日志流,建议在入口处打上来源标签,合并成统一 schema 后再进入后续计算。

2.3 黑名单过滤的工程要点

黑名单过滤不能只做一次,因为一次推荐任务里可能涉及多个数据源。推荐的做法是将黑名单表广播到各 Executor,在内存里直接查,而不是每次 join 都触发 Shuffle:

val blackUsers = spark.sparkContext.broadcast( blacklistDF.select("user_id").rdd.map(_.getString(0)).collect().toSet ) val filteredRDD = cleanRDD.filter(row => !blackUsers.value.contains(row.getString(0))) spark.createDataFrame(filteredRDD, schema).createTempView("user_log_clean")

核心是 Spark Broadcast 变量。黑名单通常只有几万到几十万个 userId,相对于亿级行为日志是明显的小表,做 map-side join 比 sort-merge join 快一个数量级。广播变量是只读的,不能在里面做累加操作。黑名单每天更新时,就在每天跑批前重新构建广播变量。

3. 协同过滤的实现:ALS 与 ItemBasedCF 对比

这个项目里同时出现了ALSDemo$.classItemBasedCF$.class,正好覆盖协同过滤的两种实现思路。ALS(交替最小二乘)适合在 Spark 上做分布式矩阵分解,ItemBasedCF 更适合逻辑简单、可解释性要求高的场景。

3.1 两种算法的工作原理

ALS 把用户对物品的评分矩阵分解成两个低维矩阵:用户特征矩阵 U 和物品特征矩阵 V,让 U·V^T 逼近原始评分矩阵。由于 Spark MLlib 里的 ALS 是分布式交替优化,它能处理百万级用户和千万级物品。ItemBasedCF 不依赖矩阵分解,它先计算物品之间的相似度,再根据用户已经交互过的物品,去加权求和未交互物品的预测分。两者核心区别:ALS 训练成本高但离线预测快,ItemBasedCF 计算相似度矩阵快,但物品量一大,内存占用会明显上升。项目里 Scala 写 ALS、Java 写 ItemBasedCF,正好对应两个语言生态的长处:Scala 写机器学习 API 简洁,Java 做线上工程维护更稳。

3.2 用 Scala 实现 ALS 模型训练与评估

下面是可运行的 ALS 训练代码。先加载清洗后的行为数据,把 view/cart/pay 映射成不同权重,再训练和评估:

import org.apache.spark.ml.recommendation.ALS import org.apache.spark.sql.functions._ val ratings = spark.table("user_log_clean") .filter("action in ('view','cart','pay')") .withColumn("rating", when(col("action") === "view", 1.0) .when(col("action") === "cart", 3.0) .when(col("action") === "pay", 5.0)) .select("user_id", "item_id", "rating") .where("rating is not null") val Array(train, test) = ratings.randomSplit(Array(0.8, 0.2), seed = 42) val als = new ALS() .setMaxIter(10) .setRank(12) .setRegParam(0.01) .setUserCol("user_id") .setItemCol("item_id") .setRatingCol("rating") .setColdStartStrategy("drop") val model = als.fit(train) val predictions = model.transform(test) val evaluator = new org.apache.spark.ml.evaluation.RegressionEvaluator() .setMetricName("rmse") .setLabelCol("rating") .setPredictionCol("prediction") val rmse = evaluator.evaluate(predictions) println(s"RMSE = $rmse")

ALS 参数里最值得调的是setRanksetRegParam。rank 控制隐因子维度,调大能拟合更复杂的用户兴趣,但也会增加过拟合风险;regParam 是正则化系数,越大模型越平滑。生产环境一般把 rank 调到 20~50,maxIter 调到 20 以上。setColdStartStrategy("drop")必须保留,否则测试集里出现新物品时会得到 NaN 预测,导致评估失败。行为权重 1/3/5 是经验值,如果有真实转化率数据,应该替换成对应商品的实际价值。

3.3 用 Java 实现 ItemBasedCF 核心逻辑

如果团队以 Java 为主,可以用轻量级自研实现代替 MLlib 里的 CF。下面是一个简化版的物品相似度计算核心:

import java.util.HashMap; import java.util.HashSet; import java.util.Map; import java.util.Set; public class ItemBasedCF { // item_id -> user_id set,表示哪些用户与该物品产生过交互 private Map<String, Set<String>> itemUsers = new HashMap<>(); public void addItemUser(String itemId, String userId) { itemUsers.computeIfAbsent(itemId, k -> new HashSet<>()).add(userId); } // 余弦相似度:交集 / 各自用户数乘积的平方根 public double cosineSimilarity(String itemA, String itemB) { Set<String> usersA = itemUsers.getOrDefault(itemA, new HashSet<>()); Set<String> usersB = itemUsers.getOrDefault(itemB, new HashSet<>()); if (usersA.isEmpty() || usersB.isEmpty()) return 0.0; Set<String> intersect = new HashSet<>(usersA); intersect.retainAll(usersB); return intersect.size() / Math.sqrt(usersA.size() * (double) usersB.size()); } }

这段代码把物品-用户倒排表放在 HashMap 里,两个物品之间的相似度通过“共同交互用户数”归一化得到。它没有外部依赖,方便本地调试;缺点是所有数据必须放进单机内存,物品数超过 20 万后建议改用 Spark SQL 的crosstab或 GraphX 计算共现矩阵。ItemBasedCF 的实际推荐逻辑是:用户对某物品产生行为后,取 Top N 个最相似物品,按相似度加权汇总生成候选评分。实现时要注意把用户已经交互过的物品过滤掉,否则推荐位会被老内容占据。

3.4 两种模型的选型边界与参数设置

在电商项目里,ALS 适合做“猜你喜欢”这类个性化排序场景,ItemBasedCF 更适合“看了又看”“买了又买”这类强关联场景。两者还可以做混合:用 ALS 生成候选集,再用 ItemBasedCF 的相似度得分对候选排序做二次修正。

场景推荐算法数据量训练频率
首页猜你喜欢ALS亿级行为每天一次
商品详情页相关推荐ItemBasedCF千万级每小时一次
新用户冷启动热门榜无历史行为实时计算

另一个容易踩的坑是setImplicitPrefs的设置。rating 列如果是通过点击/加购/支付映射出来的,本质上是隐式反馈,建议setImplicitPrefs(true),并把setAlpha调在 40 左右。否则纯显式模型在稀疏行为数据上收敛很慢。ALSDemoItemBasedCF两个类名没有明确说明反馈类型,但电商日志里 90% 以上是浏览行为,更适合按隐式反馈处理。

4. 冷启动与热门榜:HotProductByArea 如何完成兜底推荐

协同过滤模型最大的短板是冷启动。新用户没有历史行为,新商品没有交互记录,ALS 和 ItemBasedCF 都无法直接给出个性化推荐。项目里的HotProductByArea$.class正是为这种情况兜底:基于日志中的 area_id 做聚合统计,把点击、加购、支付热度高的商品优先推给新用户。

4.1 冷启动问题与热门榜定位

冷启动发生在三个节点:新用户第一次登录、新品刚上架、老用户跨地区访问。前两种用热门榜顶上是电商平台的常规操作;第三种通常也要结合地区热门避免推荐太“偏”。热门榜不依赖训练,计算快,结果稳定,缺点是只反映流量热点而不反映个体偏好。所以它不能替代协同过滤,而是作为混合推荐的最底层。推荐结果融合时,默认策略是:个性化候选为空时用热门榜填补;个性化候选非空时,热门榜商品占据推荐位 20%~30% 的比例,用于探索新兴趣。

4.2 基于 Spark SQL 的地区热门商品统计

过滤掉无效日志和黑名单用户后,热度统计可以直接用 SQL 完成。以地区分组,按行为权重算热分:

SELECT area_id, item_id, SUM(CASE WHEN action = 'view' THEN 1 WHEN action = 'cart' THEN 3 WHEN action = 'pay' THEN 5 END) AS hot_score FROM user_log_clean WHERE dt >= '2024-01-01' AND dt < '2024-01-08' GROUP BY area_id, item_id ORDER BY hot_score DESC

执行计划会先按时间分区过滤,再做area_iditem_id两级聚合。SUM(CASE...WHEN...)把行为换算成加权分数,比单纯数点击数更能体现支付和加购的贡献。最后的ORDER BY hot_score DESC做全局排序。如果日志量每天达到十亿条,全局排序会变成瓶颈,可以改为每个分区内排序后截取 Top K,减少 Shuffle 输出量。

同样的逻辑在 Scala DataFrame API 里可以写成:

import org.apache.spark.sql.expressions.Window val hot = spark.table("user_log_clean") .filter(col("dt").between("2024-01-01", "2024-01-07")) .groupBy("area_id", "item_id") .agg( sum(when(col("action") === "view", 1) .when(col("action") === "cart", 3) .when(col("action") === "pay", 5)).alias("hot_score") ) .withColumn("rn", row_number().over(Window.partitionBy("area_id").orderBy(col("hot_score").desc))) .filter(col("rn") <= 100)

这里用row_number()开窗函数在每个地区内独立排序,避免把所有明细拉到单节点再排序。生产环境还需要把item_id与商品表 join 一次,过滤掉下架商品。离线作业可以每天凌晨跑,输出一张hot_product_area表,供推荐服务直接读取。

4.3 热门榜与协同过滤结果混合推荐

混合推荐不要简单地把两组结果并排输出,而是要有分层逻辑。常见做法:ALS 先给每个用户生成 200 个候选物品,热门榜给出每个地区 Top 100,然后合并时按权重打分,最后用随机扰动保证多样性。

val personalTopK = alsModel.recommendForUserSubset(usersDF, 200) val areaHot = spark.table("hot_product_area") val merged = personalTopK .join(areaHot, Seq("item_id"), "full_outer") .withColumn("final_score", when(col("als_score").isNull, col("hot_score") * 0.3) .when(col("hot_score").isNull, col("als_score") * 0.7) .otherwise(col("als_score") * 0.7 + col("hot_score") * 0.3))

这个方案给个性化打分 0.7 的权重,热门分 0.3。个性化缺失时,热门分直接做底;个性化存在时,双方按比例混合。权重 0.7/0.3 是经验值,实际生产中要通过小流量实验测试,不能直接复制。混合时还要做商品去重和类目打散,避免连续出现多个相似类目的商品。

5. 运行验证:Spark 提交参数、离线评测与全链路调优

模型写完后,最花时间的往往是把作业跑稳定。这里给出常见的部署命令和验证手段。

5.1 Spark 提交命令与资源配置

以 YARN 集群模式运行完整作业时,我会用一个脚本管理参数:

spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --driver-memory 4g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.sql.shuffle.partitions=400 \ --class com.ecommerce.recommend.ALSDemo \ recommend-1.0.jar \ --trainPath /data/user_log \ --modelPath /model/als

这里的--class要替换成你打包后的完整类名。executor-memory 8g需要结合数据量调整,太大反而引起 GC 停顿;spark.sql.shuffle.partitions=400用于避免 Shuffle 聚合集中在少数 task 上。ALS 训练结束后建议显式保存model.save,后续离线推荐直接ALSModel.load,不用每天重训。如果频繁出现 Executor OOM,先调大spark.memory.offHeap.enabled或减小 rank,不要无脑加内存,因为 OOM 往往来自 Shuffle 溢出和广播变量膨胀。

5.2 离线评测:同时看 RMSE 和排序命中率

RMSE 只能反映评分拟合的误差,推荐任务更关心排序质量。评估时建议同时看 RMSE、Precision@K 和 Recall@K:

指标计算公式目标
RMSEsqrt(sum((真实评分-预测评分)^2)/N)越小越好
Precision@K推荐前K中用户真实交互的物品数 / K0.02~0.3
Recall@K推荐前K中命中物品数 / 用户真实交互总数0.1~0.5
覆盖率被推荐的物品数 / 总物品数越高越好

跑评测时,候选集要限制在用户确实有过行为的物品范围内,否则覆盖率虚高、Recall 失真。一般在离线阶段把评测集切到最近 7 天行为,训练集用之前 30 天,这样更贴近线上时间分布。

5.3 响应时效与实时推荐兜底

如果业务要求推荐结果每小时更新,可以采用批流分离:ALS 每天夜里训练,产出 userFactors 和 itemFactors 向量,放进 Redis;Spark Streaming 每小时消费用户行为日志,更新hot_product_area热门榜;线上 Java 服务直接读 Redis 里的向量做内积排序,再结合热门榜结果做混合。这里有个技巧:userFactors 和 itemFactors 可以序列化成 Float 数组存入 Redis,用内积替代 transform,省掉每次推荐都启动 Spark Context 的开销。广告场景下,AdLogInfo也可以走同样的热门统计逻辑,按广告位维度聚合曝光和点击率,再参与最终排序。

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询