前阵子有学员拿一套Spark电商推荐系统的毕业设计源码找我,说代码在原环境怎么都跑不起来,也不知道怎么跟导师讲清楚。我帮他从头到尾梳理了一遍数据链路之后发现,这套东西代码量不大,但涉及的点非常杂:数据清洗、特征加工、协同过滤召回、排序模型、Spark参数调优、集群部署,哪一环没对上都会出问题。尤其很多人买来设计源文件和万字报告,结果只是能打开,真到要复现、要改业务场景、要扩展成自己的方案时,就完全卡住了。
这篇内容我就把一套从0到1可落地的Spark大数据电商推荐系统完整拆开讲一遍,包含整体架构设计、数据预处理、召回和排序实现、集群调优、冷启动和评估,以及我实际调试时踩过的坑。不管你是正在做大数据毕业设计,还是刚接推荐类项目想快速看懂Spark这套玩法,都可以直接拿这篇文章当“导读”来用。我会尽量按实操顺序走,每个环节都给出可复用的代码逻辑和参数参考,而不是空谈架构。
1. 项目设计与架构思路拆解
1.1 推荐系统到底在解决什么问题
做电商推荐之前,先要搞清楚一个核心问题:用户为什么需要推荐?
因为电商平台商品数量远超用户浏览能力。一个中型电商平台可能有几十万、上百万个SKU(库存量单位),用户不可能逐页翻完。传统的搜索是“人找货”,用户带着明确目标输入关键词;推荐则是“货找人”,在用户没有明确表达意图时主动预测他可能感兴趣的商品。推荐系统的目标不是单纯提升点击率,而是在用户、商品、场景三者的关系中,找到当前时刻最合适的匹配结果。比如用户之前浏览过某类数码产品,系统在首页推荐位给他展示同类新品,就是在降低他发现商品的成本。
从业务指标来看,推荐系统最直接的贡献是提升点击率、转化率和客单价,同时还能承担一部分“长尾商品分发”的职责。很多人做项目时容易陷入技术细节,一上来就调模型、调参数,忽略了业务定义。我建议你先明确推荐位在哪个页面、面向什么用户、希望优化什么指标。这个决定后面所有特征和样本设计的方向。普遍情况是:首页推荐位看点击率,购物车和结算页看转化率,详情页推荐看关联购买率。不同的业务目标,样本标注和模型评估方式是完全不一样的。
1.2 整体数据链路与模块划分
一套完整的推荐系统,从数据产生到最终推荐结果呈现在用户面前,大致经历这样一条链路:
数据源(用户行为日志、商品明细、订单数据)→ 数据清洗与标准化(Spark任务)→ 特征加工(用户特征、物品特征、交叉特征)→ 召回阶段(多路召回,包括协同过滤、热度召回、规则召回)→ 排序阶段(特征拼接、模型打分)→ 结果存储(Redis或数据库)→ 推荐服务接口 → 客户端展示。
我按模块划分成四层来理解:
- 数据层:负责收集和存储原始数据,包括用户行为日志、商品信息表、订单表。离线场景下通常落HDFS,业务库数据通过同步工具抽到数仓。
- 计算层:用Spark做离线批处理,完成清洗、聚合、特征计算、模型训练和推理。这一层是整个系统的核心,也是文章接下来重点展开的部分。
- 算法层:召回算法(ALS协同过滤、Item-CF、热度兜底)和排序算法(逻辑回归、GBDT等)都在这一层。召回的目标是从全量商品中快速缩小到几百个候选,排序的目标是把候选集按用户兴趣精准排序。
- 服务层:提供查询接口,接收用户ID后快速返回推荐列表。一般用Redis做缓存,保证毫秒级响应。
这套分层思路不仅适用于毕业设计,生产环境也基本是这个结构,只不过生产环境多了实时流计算(如Kafka+Flink)和AB实验平台。做项目时我建议先按离线链路做通,再考虑实时化,否则复杂度会翻倍。
1.3 为什么选Spark,而不是纯Python或SQL
很多初学者会问:数据量也不大,用Pandas直接算不行吗?为什么一定要上Spark?
这里要分清场景。Pandas单机处理几百万行数据其实挺流畅,但推荐系统在真实业务里要处理的是用户行为日志的宽表、商品维表、订单维表,动辄几亿行。一旦数据超过单机内存,Pandas就会频繁触发Swap,任务直接卡死。SQL能解决部分聚合问题,但复杂特征工程、矩阵分解、模型训练这些算法逻辑在SQL里写非常痛苦,维护成本高。
Spark的价值在于它提供了统一的分布式计算框架,既能做SQL式的结构化数据处理,又能写自定义算法逻辑,还自带MLlib机器学习库,ALS、逻辑回归等算法直接调用即可。用生活类比来说,Pandas相当于你自己在小厨房里炒菜,适合小分量;Spark是一条中央厨房流水线,多个灶台同时开火,菜量再大也能按流程出餐。推荐系统的数据链路天然是流水线,每一步都可以用Spark算子完成,工程上衔接最顺畅。
1.4 技术栈清单与角色说明
下面的表格整理了这套系统常用的技术组件,以及每个组件承担的角色。做项目的时候不用全部上,按自己机器资源量力而行。
| 组件 | 角色 | 使用说明 |
|---|---|---|
| Spark Core / Spark SQL | 数据清洗、聚合、特征加工 | 核心计算引擎,建议用Spark 3.x以上版本 |
| Spark MLlib | 模型训练和推理 | ALS协同过滤、逻辑回归、GBT等算法库 |
| HDFS | 分布式存储 | 存放原始数据、中间结果、模型文件 |
| YARN | 资源调度 | 集群模式下管理CPU和内存资源 |
| Redis | 线上缓存 | 存用户最终推荐列表和Item相似结果 |
| MySQL / MongoDB | 业务元数据 | 存商品表、用户表,以及跑批后的结果表 |
| Kafka | 消息队列 | 实时行为日志采集,进阶扩展时使用 |
| Zookeeper | 集群协调 | Hadoop和Kafka集群的协调服务 |
还有一点要注意:代码和报告里可以把架构画得很高大上,但实际落地时一定要控制规模。我见过太多人把组件堆到七八个,结果跑在个人电脑上光集群维护就消耗了大部分精力,核心算法反而没有时间搞。稳妥的做法是,先保证Spark+HDFS+MySQL/Redis能把链路跑通,其他组件作为扩展点写在报告里即可。
2. 数据准备与预处理实操
2.1 数据来源与表结构设计
做推荐系统,第一步不是写算法,而是先把数据准备好。最常用的电商数据来自三张表:用户行为日志表、商品信息表、用户画像/订单表。
用户行为日志表是推荐建模最重要的数据,字段一般包含:
CREATE TABLE user_behavior ( user_id BIGINT COMMENT '用户ID', item_id BIGINT COMMENT '商品ID', behavior_type STRING COMMENT '行为类型:pv/cart/fav/buy', category_id BIGINT COMMENT '商品类目ID', session_id STRING COMMENT '会话ID', timestamp BIGINT COMMENT '行为时间戳', event_date STRING COMMENT '日期分区,格式yyyyMMdd' ) PARTITIONED BY (event_date STRING);商品表主要用来做特征工程和召回后的规则过滤,字段包括商品ID、类目ID、标题、价格、销量、评分、上下架时间等。用户画像表则包括年龄、性别、城市、注册时间等,数据稀疏时用来做冷启动和人口统计学特征。
如果是做毕业设计,原始数据可以用公开数据集,比如经典的MovieLens评分数据、淘宝用户行为数据集或一些开源的电商模拟数据。用公开数据有一个好处是字段规范、文档齐全,但要注意把评分数据转换成电商行为数据时,需要自己定义行为权重的映射关系,这一点后面会详细说。
2.2 数据清洗与格式统一
拿到原始日志后,第一件事是清洗。日志数据通常存在重复、空值、异常时间、非法ID等问题,不处理直接用于训练,会让结果偏差很大。下面是我常用的PySpark清洗逻辑:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, to_date, unix_timestamp spark = SparkSession.builder \ .appName("ECommerceRec-ETL") \ .enableHiveSupport() \ .getOrCreate() df = spark.read.parquet("/data/raw/user_behavior") # 1. 去重:同一用户同一商品同一行为同一会话只保留一次 df = df.dropDuplicates(["user_id", "item_id", "behavior_type", "session_id"]) # 2. 过滤异常数据 df = df.filter(col("user_id").isNotNull()) df = df.filter(col("item_id").isNotNull()) df = df.filter(col("item_id") > 0) # 3. 时间清洗:过滤未来时间和太老的数据 df = df.filter((col("timestamp") >= unix_timestamp("2023-01-01")) & (col("timestamp") <= unix_timestamp("2023-12-31"))) # 4. 行为类型规范化,统一小写 df = df.withColumn("behavior_type", when(col("behavior_type").isin("PV", "Pv"), "pv") .when(col("behavior_type").isin("CART", "Cart"), "cart") .otherwise(col("behavior_type")))三个容易忽略的细节:
第一,行为类型统一大小写,别小看这个,很多线上日志因为来源不同,同一行为有各种写法,后面join的时候会莫名丢数据。
第二,异常时间过滤一定要做,真实日志里经常出现1970年、2038年这样的脏时间戳,不处理会导致时间衰减特征计算出来的权重完全错乱。
第三,同用户同商品短时间内反复刷新产生的“无效曝光”要不要保留,取决于业务。做点击率预估时,重复曝光点击通常只保留第一次,否则会把一个样本重复放大N倍,造成模型过拟合。
2.3 行为权重与标签映射
清洗完之后,需要给不同类型的行为定义一个用于训练的“评分”。这里有两种思路,很多人会混为一谈。
第一种思路是用于协同过滤的隐式反馈评分,ALS算法里可以指定confidence权重,比如点击算1,加购算3,收藏算4,购买算5。这个比例不是拍脑袋定的,要结合业务转化漏斗。正常情况下,从点击到加购的转化率大概在5%左右,从加购到购买的转化率可能在30%左右,所以加购行为的信息量明显高于点击。我常用的映射如下:
| 行为类型 | 权重值 | 说明 |
|---|---|---|
| pv | 1.0 | 点击/曝光,最弱信号 |
| fav | 4.0 | 收藏,表达明确兴趣 |
| cart | 5.0 | 加购,购买意图强 |
| buy | 10.0 | 购买,最强正反馈 |
第二种思路是构造排序模型的标签,二分类问题中,曝光未点击是负样本,点击或购买是正样本。如果拿不到曝光数据,就用召回结果里用户未点击的商品作为负样本。这个细节后面章节专门讲,这里先记住:协同过滤的rating和排序模型的label是两套东西,不要混用。
2.4 特征加工的常用方式
特征工程决定了模型效果的上限。实际项目中,特征往往比模型算法更影响结果。我把这套系统里最核心的特征分成三类,分别说一下加工方式。
用户侧特征:包括用户近7天浏览/收藏/加购/购买的商品数,用户最常购买的类目,用户的活跃度(按行为天数或行为总数分桶),用户价格带偏好(浏览商品的均价、最高价、最低价)。用Spark实现就是在一个时间窗口内按user_id做groupBy聚合。
物品侧特征:包括商品近7天曝光量、点击率、转化率、收藏率、加购率,商品所属类目以下单量计算的热度分,以及商品上架天数。商品热度有一个偏移问题,新上架商品天然数据少,直接按原值排序会被老爆品压住,所以一般会做贝叶斯平滑或直接加时间衰减。
交叉特征:最常用的是“用户在某类目下的行为次数”,比如用户过去30天在“手机数码”类目下点击了多少次、购买了多少次。这种特征能反映用户对某个品类的偏好强度,逻辑回归这类线性模型很吃这种交叉信息。
下面是一段典型的特征聚合代码:
from pyspark.sql import functions as F user_cat_feature = df.filter(col("behavior_type") == "buy") \ .groupBy("user_id", "category_id") \ .agg( F.count("*").alias("buy_cnt"), F.sum(when(col("behavior_type") == "buy", 1).otherwise(0)).alias("cnt_30d") )特征加工完成后,建议统一输出成Parquet列式存储,按时间分区保存到HDFS。Parquet比CSV省空间,读起来也快,后面跑模型不用每次重新加工一遍。
3. 召回层实现:多路召回策略
3.1 ALS协同过滤实现个性化召回
协同过滤是推荐系统最经典的召回算法。核心思想很简单:找到和我相似的用户,把相似用户喜欢的商品推荐给我;或者找到我喜欢的商品的相似商品,推荐给我。Spark MLlib里的ALS(交替最小二乘法)做的是矩阵分解,把“用户-商品”评分矩阵分解成两个低秩矩阵的乘积,用低维向量表示用户和商品的隐含特征。
ALS的代码在PySpark里很简洁:
from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator als = ALS( userCol="user_id", itemCol="item_id", ratingCol="rating", coldStartStrategy="drop", implicitPrefs=True, alpha=40.0, rank=20, maxIter=15, regParam=0.1, numUserBlocks=10, numItemBlocks=10 ) model = als.fit(train_df) user_recs = model.recommendForAllUsers(50)几个关键参数我重点解释一下:
rank是隐含特征维度,决定了模型容量。rank太小,表达能力不够;rank太大,训练慢而且容易过拟合。常规项目从10到50之间调,小数据集可以先用20作为起点。alpha是隐式反馈的置信度系数,只在implicitPrefs=True时生效,alpha越大,行为次数的差异化影响越强。regParam是正则化系数,防止过拟合,常用范围是0.01到1之间。coldStartStrategy="drop"表示预测时遇到新用户或新商品就跳过,不产出NaN,这个必须设置,否则结果表里会飘着一堆NaN,写入Redis时会报错。
ALS训练完成后,recommendForAllUsers会为每个用户返回TopN商品。这一步得到的是个性化候选集,但实际工程里不会只用ALS一路召回,因为ALS对冷启动用户无能为力,对头部热门商品容易过度集中,所以需要多路召回组合。
3.2 Item-CF基于物品的协同过滤
ALS之外,我习惯再加一路Item-CF召回,它的逻辑是:用户对商品A感兴趣,那么和A相似的物品B也应该推荐给该用户。“相似”的定义来源于用户行为共现——两个商品被同一批用户点击或购买过,就认为它们有相似性。
Spark实现Item-CF的核心是“物品对共现计数”,代码如下:
from pyspark.sql import functions as F # 输入:用户-商品-行为 behavior_df = df.filter(col("behavior_type").isin(["cart", "buy"])) # 自连接,生成同一用户下的商品两两组合 pair_df = behavior_df.alias("a") \ .join(behavior_df.alias("b"), (F.col("a.user_id") == F.col("b.user_id")) & (F.col("a.item_id") < F.col("b.item_id")), "inner") \ .select( F.col("a.item_id").alias("item_a"), F.col("b.item_id").alias("item_b") ) # 统计共现次数 item_sim = pair_df.groupBy("item_a", "item_b") \ .agg(F.count("*").alias("co_cnt")) \ .filter(F.col("co_cnt") >= 3) # 过滤噪声共现这里有个筛选条件co_cnt >= 3,意思是两个商品至少被3个不同用户共同购买过,才认为它们相似。具体阈值根据数据稀疏程度调,如果数据量小可以放宽到2,否则相似关系里全是噪声。
Item-CF相对ALS的优势是结果可解释性强,推荐理由可以说“买了A的用户也买了B”,产品上更容易展示。同时Item-CF在线下计算好后,可以存成“商品-相似商品列表”,线上接口根据用户最近点击过的商品查这个表,实时拼出候选集,逻辑简单响应快。
3.3 热度兜底与规则召回
不管模型建得多好,总有用户没有任何历史行为,或者行为量太少,ALS和Item-CF都拿不到有效结果。这时候需要一路“兜底召回”,最简单的就是热销榜和新品榜。
热度分不能直接用销量排序,否则排行榜常年不变,新品永远没有出头机会。我常用的热度分公式是:
[ score = \frac{\log(1 + click + 3 \times cart + 5 \times buy)}{1 + \log(1 + days_since_on_shelf)} ]
分子表示商品综合热度,分母是时间衰减因子,上架时间越长,热度分被稀释得越多,这样新品只要短期内表现不错,就有机会冲到榜单前列。用Spark实现就是按商品聚合行为次数,再代入公式计算Score,取TopN。
另外还有一种规则召回:基于类目的偏好召回。比如用户最近买过手机,那就把同价位的手机配件、耳机等周边商品作为候选。这种召回简单但非常实用,尤其对用户行为稀疏的场景,比纯模型更稳。
4. 排序层实现:从候选集到最终列表
4.1 样本组织与特征拼接
召回层让每个用户得到几十到几百个候选商品,接下来要做的就是把这批候选商品精确排序。排序模型的本质是一个学习排序问题,需要一个监督信号来训练模型。
训练样本的组织方式是:把用户、候选商品、特征、标签拼在一起。标签的定义通常为:
- 正样本:用户点击过、加购过、购买过的商品。如果想区分行为等级,可以做多目标,但入门阶段建议先做二分类,点击或购买为正样本,未交互为负样本。
- 负样本:曝光未点击的商品,或随机从召回候选里抽取用户未行为过的商品。
负样本的比例很关键。一般正负样本比控制在1:2到1:10之间,极端不平衡会导致模型把所有样本都预测为正或都预测为负。真实业务里负样本量远大于正样本,项目里如果发现模型准确率很高但AUC很低,多半是负样本量过少或者采样方式不对。
特征拼接的时候要特别注意特征泄漏问题。比如用“用户是否购买过该商品”作为特征再预测用户是否购买该商品,这在离线评估时AUC会虚高,线上完全无效。正确做法是,构造特征时只能使用预测时间点之前的数据,比如预测7月1日会不会点击,只能用6月30日及以前的行为特征。
4.2 用Spark ML实现排序模型
入门阶段我推荐先用逻辑回归(LR),原因有三个:第一,模型可解释性强,每个特征的权重直接告诉你哪个因素影响最大;第二,Spark MLlib对LR支持成熟稳定;第三,LR可以作为后续复杂模型的上线基线。
实现代码如下:
from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression from pyspark.ml.evaluation import BinaryClassificationEvaluator feature_cols = ["user_click_cnt", "user_buy_cnt", "item_click_rate", "item_cart_rate", "user_cat_buy_cnt", "item_price", "item_rank"] assembler = VectorAssembler(inputCols=feature_cols, outputCol="raw_features") scaler = StandardScaler(inputCol="raw_features", outputCol="features", withStd=True, withMean=False) lr = LogisticRegression(featuresCol="features", labelCol="label", maxIter=50, regParam=0.01) # Pipeline串联执行 from pyspark.ml import Pipeline pipeline = Pipeline(stages=[assembler, scaler, lr]) model = pipeline.fit(train_df) pred_df = model.transform(test_df) evaluator = BinaryClassificationEvaluator(labelCol="label", metricName="areaUnderROC") auc = evaluator.evaluate(pred_df)StandardScaler这一步很多人会省,但逻辑回归对特征尺度敏感。物品点击率是0到1的小数,用户行为计数可能是几千的大数,不归一化的话,大数值特征会主导梯度更新,严重影响收敛速度。先做标准差归一化,训练速度会有肉眼可见的提升。
如果之后想提升效果,可以换Spark的GBTClassifier(梯度提升树),树模型对特征尺度不敏感,且能自动捕捉非线性关系,但训练时间更长,调参也更复杂。我一般以LR作为baseline,GBDT效果如果明显提升再上线。注意不要一上来就上深度模型,数据量不够时深度模型很容易被传统模型吊打。
4.3 TopN生成与业务规则合并
排序模型输出每个用户-商品对的点击概率,接着要按概率倒序截取TopN。这里还要叠加一些业务规则,比如:
- 过滤掉用户近30天已经购买过5件以上的同类商品,避免重复推荐;
- 过滤掉当前已下架或库存为0的商品;
- 同一商品品牌或同一卖家,单次推荐列表里不能出现过多,防止推荐结果太集中;
- 每隔一定位置插入运营指定的广告位或活动商品。
用Spark Window函数可以很方便地取TopN:
from pyspark.sql import Window from pyspark.sql import functions as F window = Window.partitionBy("user_id").orderBy(F.col("score").desc()) top_df = scored_df.withColumn("rank", F.row_number().over(window)) \ .filter(F.col("rank") <= 50) \ .drop("rank")到这里,离线结果基本就成型了。每天凌晨跑一次批量任务,把TopN结果写入线上的存储组件。
5. 集群部署与Spark性能调优
5.1 集群部署策略
系统能不能跑起来,很多时候不取决于代码,而取决于集群部署。本地开发用local模式,SparkSession不指定master,跑起来方便调试。但要上线或者做演示,需要部署集群。
这里推荐最通用的方案:三台机器搭建Hadoop YARN集群,Spark运行在YARN上。机器配置不用太高,学习环境16核64G内存就能跑得动中型数据。部署步骤大致是:
- 三台机器都装JDK 8,配置SSH免密登录;
- 安装Zookeeper并启动,保证HDFS高可用;
- 安装Hadoop,配置hdfs-site.xml、yarn-site.xml、core-site.xml,启动NameNode和DataNode;
- 安装Spark,配置spark-env.sh中JAVA_HOME、HADOOP_CONF_DIR,复制spark-defaults.conf;
- 验证:启动后运行spark-submit提交任务,检查YARN Web UI上的执行状态。
关于是选Standalone还是YARN,我个人的经验是:如果没有Hadoop环境,只是为了单机跑代码,Standalone模式更快;但如果要做多租户资源管理、和其他任务共享集群,一定要用YARN。YARN的好处是资源隔离和队列管理,Spark任务跑挂了不会拖垮整个HDFS。
集群部署有个关键点要提醒:内存和磁盘要提前规划。Spark在shuffle阶段会在本地磁盘写大量临时文件,如果/tmp空间不足,任务会报“No space left on device”。我习惯在core-site.xml里把hadoop.tmp.dir指向空间最大的数据盘,并单独分出200G以上给Spark的local.dir。
5.2 内存模型与资源分配实例
Spark调优最核心的是内存。很多人直接照着网上抄参数,结果任务频繁OOM,甚至YARN直接把container杀掉了。先看懂Spark内存模型再配参数,才能少踩坑。
Spark Executor内存由三部分组成:执行内存(Execution Memory)、存储内存(Storage Memory)、预留内存(Reserved Memory)。默认配置下:
- spark.memory.fraction=0.6:表示JVM堆内可用内存中60%用于执行和存储共享;
- spark.memory.storageFraction=0.5:表示共享区内存储内存初始占一半,之后如果执行内存不足,可以抢占存储内存;
- 剩下40%留给用户代码、内部元数据和防止OOM的安全余量。
举个例子说明:假如executor申请的堆内存是8GB,则:
- 可用内存约为8GB * 0.6 = 4.8GB;
- 初始存储内存为4.8GB * 0.5 = 2.4GB;
- 初始执行内存为4.8GB * 0.5 = 2.4GB;
- 实际执行内存不足时可以抢占存储空闲部分。
所以申请8GB堆内存,不代表shuffle能用到8GB。如果你的任务经常做大规模groupBy或join,那么executor内存要适当调大,或者增加分区数减少单分区数据量。
典型的生产环境提交命令模板:
spark-submit \ --master yarn \ --deploy-mode cluster \ --name rec-offline \ --num-executors 10 \ --executor-cores 4 \ --executor-memory 8G \ --driver-memory 4G \ --conf spark.default.parallelism=200 \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.memory.fraction=0.7 \ --conf spark.memory.storageFraction=0.4 \ --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \ --py-files deps.zip \ main.pynum-executors乘以executor-cores就是集群总并发数,要和YARN队列的最大资源匹配。有人会遇到“spark on yarn cpu只能用1个”的问题,其实不是只能用1个,而是YARN的vcore分配和Spark任务并发度之间有换算关系。如果yarn.nodemanager.resource.cpu-vcores配的是8,但每个container默认只申请1个vcore,那么executor-cores=4时,一个executor会占用4个vcore配额。要查清楚yarn-site.xml里的最大分配配置,同时确认spark.executor.cores是否真的生效。单纯从YARN界面看到每个container 1 vcore,大概率是没设spark.executor.cores,或者设了但被YARN的调度器上限限制住了。
5.3 数据倾斜的3种解决思路
跑推荐任务时最容易遇到的问题是数据倾斜:某个热门商品或热门用户的数据量远大于其他key,导致大部分Task很快执行完了,一两个Task卡在99%不动。这种情况不是因为计算量大,而是数据分布不均匀,单个Task要处理几亿条记录。
我常用的解决思路有三类:
第一类,提高并行度。最直接的方式是把spark.sql.shuffle.partitions调大,比如从200调到400或800。如果倾斜不严重,这样做就能把单个Task的压力分摊掉。这是最省事的方法,但也只是缓解,不是根治。
第二类,对热点Key加盐。先找出数据量最大的几个Key,比如热门商品的item_id,然后给它们加上随机前缀,再做join或者聚合。这样热点数据会被打散到多个Task里。处理完后需要去掉前缀再合并结果。这种方法适合做大Key的join,但代码逻辑会复杂一些。
第三类,广播小表。如果两张表join时,一张表很小(比如商品维表只有几万行),就不应该做Shuffle Join,而是用广播变量让每个Executor都存一份小表副本。这样能避免一边数据倾斜,一边没有Shuffle压力。设置spark.sql.autoBroadcastJoinThreshold,默认是10MB,16核64G的机器可以调到20MB甚至50MB,不过要监控driver端内存,广播太大driver会成为瓶颈。
注意,数据倾斜问题没有一个万能解法,要先去Spark UI看哪个Stage卡住了,点开详情看某个Task的Shuffle Read大小,再对症下药。盲目加内存或者加cores,很多时候治标不治本。
6. 冷启动、评估与线上接入
6.1 冷启动问题处理策略
推荐系统逃不开冷启动问题。冷启动分用户冷启动和商品冷启动。
用户冷启动指新注册用户没有任何行为记录,协同过滤模型拿不到他的偏好。我的处理方案分三层:
- 第一层:默认给热度榜,至少保证首页不空;
- 第二层:让用户在注册时选择感兴趣的类目,或者接入第三方数据(微信授权、位置信息)粗粒度判断偏好;
- 第三层:用户产生第一次点击后,立即根据点击行为把去重回溯刷新推荐列表,这需要实时或准实时计算。
商品冷启动指新上架商品没有行为数据,模型不会推荐它。方案是提取商品标题、类目、标签等信息,计算它和现有热销品的类目相似度或文本相似度,给一个基础曝光值。系统里我加了一个规则:所有新商品在24小时内随机出现在部分用户的“新品推荐”坑位,保证有少量曝光,进而产生行为数据,尽快进入正常推荐体系。
6.2 离线评估指标怎么算
推荐模型效果的评估,不能只看模型训练时的损失函数,要站在业务视角看离线指标。
召回阶段的指标主要是召回率、精确率和覆盖率:
- 精确率Precision@K:推荐列表TopK中,用户实际交互过的商品比例;
- 召回率Recall@K:用户实际交互过的商品中,被推荐出来的比例;
- 覆盖率Coverage:推荐出来的商品占全量商品的比例,覆盖率太低说明模型只推荐头部热门,长尾分发效果差。
排序阶段的核心指标是AUC。AUC表示模型把正样本排在负样本前面的概率,0.5代表随机,0.7以上在推荐排序里算比较可用。训练集AUC和测试集AUC相差过大,说明模型过拟合,需要调大正则参数或减少特征维度。
计算TopK指标时,建议按用户维度分开计算再取平均,而不是把所有用户预测结果混在一起算。用户行为量差异很大,混在一起算会被高频用户带偏。
6.3 推荐结果如何接入线上服务
离线任务跑完后,最终结果需要为线上服务可用。最常用的方式是把用户TopN结果写入Redis,Key按固定格式设计,比如rec:user:{userId},Value存一个有序的JSON数组。
PySpark写Redis有几个做法。简单场景下,直接循环分区数据用Redis客户端写就好:
import redis def write_to_redis(rows): r = redis.Redis(host="10.0.0.8", port=6379, db=0) pipe = r.pipeline(transaction=False) for row in rows: key = f"rec:user:{row['user_id']}" value = json.dumps([r["item_id"] for r in row["recs"]]) pipe.setex(key, 86400 * 7, value) pipe.execute() result_df.foreachPartition(write_to_redis)需要注意两点:一是设置过期时间,推荐结果每天更新,旧结果过期后不再返回;二是用pipeline批量写入,不要逐条set,否则几千个用户的写入会非常慢。
如果线上查询需求是实时的、候选集动态变化,离线写好TopN的方式就不够灵活了。进阶方案是写入Item-CF的相似商品表,线上根据用户实时点击去查相似表,再合并热度和规则召回的结果做实时排序。这是业界标准的“离线计算候选 + 线上实时拼接”思路。
7. 常见问题与排查实录
下面整理几个我在实际调试这套系统时遇到的典型问题,附带排查思路和解决建议。
| 问题现象 | 原因分析 | 解决方案 |
|---|---|---|
| Spark on YARN提交后每个Executor的CPU只显示1 vCore | spark.executor.cores未配置,或yarn.scheduler.maximum-allocation-vcores限制过小 | 在spark-submit中显式设置--executor-cores 4,并检查yarn-site.xml中的vcore配额 |
| ALS预测结果出现大量NaN | 测试数据里包含训练集中没出现过的用户或商品,且coldStartStrategy未设置 | 设置coldStartStrategy="drop",或先过滤掉无历史行为的数据 |
| 跑批任务每天越跑越慢,HDFS小文件过多 | 每次写入都生成大量小文件,导致NameNode内存压力大、任务调度慢 | 输出前用coalesce或repartition控制到合理分区数,尽量写Parquet格式 |
| 排序模型训练集AUC很高但线上点击率反而下降 | 特征泄漏或离线在线特征不一致 | 检查特征是否用了未来数据,线上特征加工逻辑必须与离线完全一致 |
| 推荐列表中90%都是同一类目商品 | 特征工程里缺少多样性约束,排序层没有控制类目比例 | 在业务规则层限制同一类目商品在TopN中不超过一定数量 |
| Redis里Key太多导致内存持续增长 | 用户量太大,且过期时间设置过长 | 设置合理的TTL,比如7天;同时单个Key的Value限制长度,只保留TopN≤50 |
| 新用户上线没有推荐兜底 | ALS和Item-CF都没有处理冷启动用户 | 增加热度榜兜底召回,或接入用户注册时的偏好选择 |
再补充一个问题:很多人会发现Spark UI上某个Stage的Shuffle Read特别大,持续数小时。这种一般是groupBy或join操作出现数据倾斜,优先用加盐方案处理,不要急着扩内存。扩内存从1T涨到2T效果很有限,反而会增加GC开销。
调试的另一个技巧是,任务跑到一半失败时,优先查看YARN日志中的Executor日志,特别是stderr和stdout,Caused by那一行通常就是真正的报错原因。我看过很多人在群里贴一堆堆栈,前面全是WARN日志,真正的异常埋在最后面。
最后分享一个调试小技巧
文章结尾我想分享一个自己一直坚持的调试习惯。做这种离线推荐系统,千万不要一开始就拿全量数据跑。我一般会先用1%或者更小的采样数据把链路跑通,确认每个Stage都能出结果、格式都是对的,再切换到全量数据。这样做的好处非常明显:小数据量下任务几分钟就能跑完,报错日志定位起来快;全量跑的时候再把并行度、executor资源调上去,兼容数据规模变化产生的问题。
另外一个经验是,任何线上推荐系统都要有最基础的监控。哪怕只是一个定时脚本,每天检查推荐结果的曝光点击率是否正常,一旦指标明显回落,立刻去查是离线任务没按时产出,还是某个头部品类商品全部下架导致推荐空窗。推荐系统不是模型上线就结束了,它更像一个需要持续观察和调整的运营系统。希望这套从设计到落地的记录,能帮你少踩一些我给上面总结过的这些坑。