FEATURED · 精选文章

Java与Scala混编电商推荐系统:从日志清洗到ALS与ItemBasedCF实践

发布时间 / 2026/9/13 16:28:36
来源 / 创域科博编辑部
栏目 / 资讯中心
Java与Scala混编电商推荐系统:从日志清洗到ALS与ItemBasedCF实践 简介这是一份面向大数据与推荐系统学习者的电商项目资源采用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否仅 AdLogInfois_clickInt是否点击0/1广告行为需要黑名单用户BlackUserList一般单独存一份在清洗阶段 join 掉。BlackUserList$.class这个类出现在项目里说明设计者选择在入口处过滤刷单用户而不是在模型层过滤。这是对的因为刷单行为一旦进入训练集会把协同过滤的相似度矩阵拉偏后续再调权重都很难救回来。2.2 用 Scala 把原始日志解析成 DataFrame日志中的原始行通常是 tab 或逗号分隔。下面是一段可以直接跑的 Scala 解析逻辑输出 Spark SQL 可查询的 DataFrameimport 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.class和AdLogInfo.class说明有两种日志流建议在入口处打上来源标签合并成统一 schema 后再进入后续计算。2.3 黑名单过滤的工程要点黑名单过滤不能只做一次因为一次推荐任务里可能涉及多个数据源。推荐的做法是将黑名单表广播到各 Executor在内存里直接查而不是每次 join 都触发 Shuffleval 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$.class和ItemBasedCF$.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(sRMSE $rmse)ALS 参数里最值得调的是setRank和setRegParam。rank 控制隐因子维度调大能拟合更复杂的用户兴趣但也会增加过拟合风险regParam 是正则化系数越大模型越平滑。生产环境一般把 rank 调到 20~50maxIter 调到 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 MapString, SetString 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) { SetString usersA itemUsers.getOrDefault(itemA, new HashSet()); SetString usersB itemUsers.getOrDefault(itemB, new HashSet()); if (usersA.isEmpty() || usersB.isEmpty()) return 0.0; SetString 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 左右。否则纯显式模型在稀疏行为数据上收敛很慢。ALSDemo和ItemBasedCF两个类名没有明确说明反馈类型但电商日志里 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_id和item_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.partitions400 \ --class com.ecommerce.recommend.ALSDemo \ recommend-1.0.jar \ --trainPath /data/user_log \ --modelPath /model/als这里的--class要替换成你打包后的完整类名。executor-memory 8g需要结合数据量调整太大反而引起 GC 停顿spark.sql.shuffle.partitions400用于避免 Shuffle 聚合集中在少数 task 上。ALS 训练结束后建议显式保存model.save后续离线推荐直接ALSModel.load不用每天重训。如果频繁出现 Executor OOM先调大spark.memory.offHeap.enabled或减小 rank不要无脑加内存因为 OOM 往往来自 Shuffle 溢出和广播变量膨胀。5.2 离线评测同时看 RMSE 和排序命中率RMSE 只能反映评分拟合的误差推荐任务更关心排序质量。评估时建议同时看 RMSE、PrecisionK 和 RecallK指标计算公式目标RMSEsqrt(sum((真实评分-预测评分)^2)/N)越小越好PrecisionK推荐前K中用户真实交互的物品数 / K0.02~0.3RecallK推荐前K中命中物品数 / 用户真实交互总数0.1~0.5覆盖率被推荐的物品数 / 总物品数越高越好跑评测时候选集要限制在用户确实有过行为的物品范围内否则覆盖率虚高、Recall 失真。一般在离线阶段把评测集切到最近 7 天行为训练集用之前 30 天这样更贴近线上时间分布。5.3 响应时效与实时推荐兜底如果业务要求推荐结果每小时更新可以采用批流分离ALS 每天夜里训练产出 userFactors 和 itemFactors 向量放进 RedisSpark Streaming 每小时消费用户行为日志更新hot_product_area热门榜线上 Java 服务直接读 Redis 里的向量做内积排序再结合热门榜结果做混合。这里有个技巧userFactors 和 itemFactors 可以序列化成 Float 数组存入 Redis用内积替代 transform省掉每次推荐都启动 Spark Context 的开销。广告场景下AdLogInfo也可以走同样的热门统计逻辑按广告位维度聚合曝光和点击率再参与最终排序。本文还有配套的精品资源点击获取
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻