尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

移动推荐项目实战:从Pandas特征工程到Spark分布式训练

移动推荐项目实战:从Pandas特征工程到Spark分布式训练 简介这是阿里天池新人赛“移动推荐”赛题的试玩项目压缩包面向人工智能、计算机科学与技术等相关专业的学生尤其适合用来完成课程作业或毕业设计中的推荐系统实践。项目代码经过验证可运行围绕移动端用户商品推荐场景完整覆盖数据预分析、特征构造、模型训练与预测等流程并提供了基于逻辑回归、GBDT与XGBoost的对比实现也包含OnSpark并行处理相关的脚本方便学习者理解推荐建模的核心步骤。压缩包内共有25个文件以19个Python源码文件为主另有说明文档、特征列表及图示等辅助材料整体仅138KB结构清晰。目前已有134人学习下载可直接参考README快速上手在本地环境中复现竞赛数据处理与模型调优的思路作为课题设计或入门天池竞赛的参考资料都很合适。1. 为什么这个移动推荐项目值得拆开看阿里天池新人赛的“移动推荐算法”常被当成入门练习但这套源码里藏着的不是“跑通一个baseline”的作业而是从单机特征工程到Spark分布式训练的一整套工程链路。压缩包里同时保留了model_lr_and_gdbt_and_xgboost和OnSpark_model两条路径前者用LR、GBDT、XGBoost做模型融合后者把同样的特征在Spark上重写了一遍解决的是“特征构造速度跟不上数据量”的真实痛点。适合两类人准备打推荐类比赛的新手以及想看“同一套逻辑如何从Pandas平滑迁移到Spark DataFrame”的工程向读者。我拆完这份代码后最大的感受是新手赛提交的分数差距往往不是模型差异而是训练集和预测集在特征统计口径上不一致。2. 移动推荐的数据结构与特征构造核心2.1 赛题数据到底长什么样这份资源对应的天池移动推荐赛题核心是用户、商品、品牌、行为日志四个维度。行为日志包含点击、收藏、加购、购买四种行为时间跨度为若干天要求预测某一天用户会购买哪些商品。官方给出的评价指标是F1值也就是预测出的“用户-商品”购买对中准确率和召回率的调和平均。这个指标决定了我们不需要精确排序只需要尽可能多地命中真实购买对同时控制误报。拿到压缩包后先看feature_construct目录里面的脚本把原始日志整理成用户维度和商品维度两大块。我一般在做这类赛题时会把数据先做一次透视看看每种行为的占比这能避免后期特征偏向高频行为。下表整理了资源里最常见的几类特征构造维度维度典型特征构造方式用户维度用户总点击数、总购买数、活跃天数按user_id分组聚合行为日志商品维度商品被点击数、被购买数、转化率按item_id分组聚合用户-商品交叉该用户对该商品点击次数、是否收藏过按user_iditem_id分组聚合时间衰减近3天、近7天行为窗口对时间戳做滑窗统计这些特征在feature_list.txt里都有完整的字段清单我强烈建议先把这份特征清单打印出来对着看而不是直接跑代码。尤其是里面的时间窗特征标注了“统计截止时间”这个参数它直接决定了训练样本和预测样本的特征分布是否一致。2.2 单机特征构造从Pandas到DataFrame的统计特征单机版特征构造的核心逻辑是把行为日志按不同key分组然后做聚合统计。下面是这份资源里典型的用户维度特征构造代码我做了简化但保留了关键口径import pandas as pd def build_user_features(action_df, end_time, windows[3, 7]): action_df: 原始行为日志含user_id, item_id, behavior_type, time end_time: 统计截止时间保证训练集和预测集用同一口径 # 限定截止时间之前的数据避免未来信息泄露 df action_df[action_df[time] end_time].copy() for w in windows: # 滑动窗口切分 mask df[time] end_time - w * 86400 wdf df[mask] # 用户在某窗口内的行为统计 tmp wdf.groupby(user_id)[behavior_type].agg( click_cntlambda x: (x 1).sum(), buy_cntlambda x: (x 4).sum() ).reset_index() tmp.columns [user_id, fuser_click_{w}d, fuser_buy_{w}d] if w windows[0]: result tmp else: result result.merge(tmp, onuser_id, howleft) return result这段代码的关键是end_time参数。很多新手直接对整个数据集做统计然后切出训练集和预测集这会导致预测集的特征里混入“未来”信息线下评估虚高线上分数崩盘。正确做法是先固定一个截止时间统计特征只用截止时间之前的行为然后滑动到下一个时间点构造训练样本。behavior_type在原始数据里用1到4表示点击、收藏、加购、购买上面的聚合直接把行为类型映射成不同计数列。需要说明的是内存不足时不要用Pandas硬扛可以把action_df换成Spark DataFrame或者按用户ID分桶后循环处理。我一般会先用小样本跑通流程确认特征列数量和结果没有空值再全量跑。3. LR GBDT XGBoost 三模型融合的玩法3.1 为什么是“线性 树模型”的组合天池赛里用LR做基线很常见因为特征维度高、稀疏性强线性模型在这个场景下足够稳定。但单纯的LR学不到特征之间的非线性交互所以一般会用GBDT或XGBoost来做高阶组合。这份资源里的model_lr_and_gdbt_and_xgboost目录把三个模型串成了一条流水线先用LR跑一版再用GBDT和XGBoost分别训练最后对预测概率做加权融合。我实测下来融合后的F1通常比单模型高出1到2个点代价只是训练时间翻倍。从原理上看GBDT对特征单调变换不敏感能自动捕捉阈值XGBoost在GBDT基础上加了二阶导和正则化训练速度更快LR则对特征做了线性建模容易校准概率。这三个模型的输出的分布不一样融合时不要直接平均而是先看各自的线下F1给模型权重。3.2 训练脚本与参数设置下面是资源里训练XGBoost模型的核心代码存储为model/xgb_train.pyimport xgboost as xgb def train_xgb(train_x, train_y, valid_x, valid_y): d_train xgb.DMatrix(train_x, labeltrain_y) d_valid xgb.DMatrix(valid_x, labelvalid_y) params { objective: binary:logistic, eta: 0.05, max_depth: 6, subsample: 0.8, colsample_bytree: 0.7, eval_metric: auc, seed: 42 } watchlist [(d_train, train), (d_valid, valid)] model xgb.train(params, d_train, num_boost_round800, evalswatchlist, early_stopping_rounds50, verbose_eval20) return modeleta0.05是学习率调低一点能让模型收敛更稳但需要增加num_boost_roundmax_depth6对高维稀疏特征来说不会过深subsample和colsample_bytree控制随机性防止过拟合。early_stopping_rounds50是这份资源里很实用的设定——树模型在800轮内如果验证集AUC连续50轮不提升就停止省时间也防止后期过拟合。注意eval_metric用的是auc因为F1与阈值有关而AUC能平滑地反映排序能力训练结束后我会额外搜索最优阈值。GBDT和LR的训练脚本在同一个目录下LR直接用sklearn.linear_model.LogisticRegression关键点是开启C正则且把solver设为liblinear因为特征维度高用默认的lbfgs收敛慢。训练完成后三个模型各自输出预测概率融合公式如下def blend_pred(lr_prob, gbdt_prob, xgb_prob, weights(0.2, 0.3, 0.5)): final_prob (weights[0] * lr_prob weights[1] * gbdt_prob weights[2] * xgb_prob) return final_prob权重组合怎么定不要在测试集上调应该基于线上榜或者再切一个验证集用网格搜索试几组比如(0.2,0.3,0.5)和(0.1,0.4,0.5)看哪组在验证集上F1最高。注意最终提交用的是预测截止时间后一天的真实购买行为所以验证集的时间切片也要模拟这个结构。4. OnSpark把特征工程搬上分布式计算4.1 什么时候需要Spark版本当用户行为日志从几百万膨胀到几千万行Pandas的groupby.agg会非常慢甚至直接OOM。这份资源里的OnSpark_model目录就是把前面Pandas特征构造改成Spark DataFrame算子核心脚本有onspark_data_preprocssing.py、onspark_generate_feature_user.py、onspark_generate_feature_product.py、onspark_merge_feature.py和onspark_prediction.py。它的价值不只是提速而是把数据处理流程拆成了可重跑的任务链预处理、用户特征、商品特征、合并、预测每个环节的产物都落盘任何一步失败不需要从头跑。4.2 Spark特征生成的典型写法以下代码是onspark_generate_feature_user.py的核心片段用于在Spark DataFrame上构造用户行为特征from pyspark.sql import functions as F def generate_user_features(spark, action_path, end_time): df spark.read.csv(action_path, headerTrue, inferSchemaTrue) # 统一时间口径过滤截止时间之后的数据 df df.filter(F.col(time) end_time) # 用户维度聚合统计 user_feat df.groupBy(user_id).agg( F.sum(F.when(F.col(behavior_type) 1, 1).otherwise(0)).alias(click_cnt), F.sum(F.when(F.col(behavior_type) 4, 1).otherwise(0)).alias(buy_cnt), F.countDistinct(item_id).alias(item_cnt), F.countDistinct(brand_id).alias(brand_cnt) ) # 时间窗口特征利用case when做近3天统计 recent_df df.filter(F.col(time) end_time - 3 * 86400) recent_feat recent_df.groupBy(user_id).agg( F.sum(F.when(F.col(behavior_type) 4, 1).otherwise(0)).alias(buy_cnt_3d) ) result user_feat.join(recent_feat, user_id, left) result.write.parquet(user_features.parquet) return result与Pandas版本最大的差别是Spark的groupBy.agg返回的是分布式的DataFrame不会把中间结果拉回DriverF.sum(F.when(...))替代了Pandas的lambda聚合避免在UDF里写Python循环。inferSchemaTrue在数据量大时容易全表扫描建议在建表或写CSV时显式指定schema。另一个实用点是所有任务用end_time参数控制统计截止时间这样Spark任务可以和单机版共用同一份特征配置避免两边统计口径漂移。4.3 从单机到Spark的迁移踩坑迁移过程中最常见的坑是空值处理。Pandas的groupby.agg对没有出现过的组会给出NaN但Spark的join(..., howleft)不会自动填充空值需要在合并后统一fillna(0)。我在这份资源里看到onspark_merge_feature.py里专门有一段fillna(0)的操作这就是一个典型的分布式与单机差异点。另一个坑是排序稳定性。Spark的groupBy之后结果顺序不固定如果后续要把特征向量拼接成定长数组建议在ID列排序后再collect_list否则同一行样本在不同次运行里特征顺序可能不同直接影响模型预测。资源里的onspark_generate_feature_user_product.py应该也是类似逻辑它生成的是用户-商品组合特征这一块在分布式环境下最容易产生shuffle倾斜如果某个热门商品被大量用户交互需要做repartition或加盐。5. 复现链路与线上分数提升技巧5.1 一步步跑通压缩包拿到压缩包后不要急着运行全部脚本。先解压确认目录结构没有损坏。如果解压报错提示invalid zip archive用以下命令检查unzip -t 阿里天池竞赛新人赛试玩_移动推荐.zip-t参数会测试压缩包完整性如果有失败的条目说明下载不完整需要重新下载。确认完整后按下面顺序执行# 1. 数据预处理生成统一格式的行为日志 python onspark_data_preprocssing.py --input raw_data --output processed_data # 2. 生成用户特征和商品特征 python onspark_generate_feature_user.py --input processed_data --end_time 2014-12-17 python onspark_generate_feature_product.py --input processed_data --end_time 2014-12-17 python onspark_generate_feature_user_product.py --input processed_data --end_time 2014-12-17 # 3. 合并所有特征 python onspark_merge_feature.py --feature_dir feature_out --output all_features.parquet # 4. 跑单机版三模型融合 cd model_lr_and_gdbt_and_xgboost python train_lr.py python train_gbdt.py python train_xgboost.py # 5. 生成提交结果 python onspark_prediction.py --model xgb_model --threshold 0.18 --output submission.csv注意end_time参数在每次运行时都要保持一致否则特征统计的时间窗口错开。训练模型的三条命令用串行执行是为了保证每个模型都完成训练后再融合。threshold参数是购买概率的判定阈值它在最后一步对F1影响非常大。5.2 剪裁阈值被忽略的提分点很多人训练完模型直接取概率大于0.5作为购买预测但在这种极度不平衡的数据里购买行为占比往往不足1%0.5这个阈值会召回极低。我拿到这份资源后第一件事就是单独写了一个脚本用验证集扫描0.05到0.5之间的阈值寻找F1最大值。from sklearn.metrics import f1_score def find_best_threshold(y_true, y_prob): best_score, best_thr 0, 0.1 for thr in range(50, 500, 5): thr / 1000.0 pred (y_prob thr).astype(int) score f1_score(y_true, pred) if score best_score: best_score, best_thr score, thr return best_thr, best_score这个脚本的收益立竿见影通常能把F1从0.12拉到0.2以上。原理是F1对阈值敏感而模型输出的概率分布本身就偏向低值区所以你应该在验证集上搜索阈值而不是用默认的0.5。最后提交前我还会做一次“购买对去重”因为同一个用户可能对同一商品有多个行为记录生成提交结果时要按user_id item_id去重保留概率最大的一条这一步能避免重复计数带来的召回率虚高。这份资源里还有一个细节值得学onspark_prediction.py在处理预测结果时会先把预测概率和验证集里的真实购买标签join在一起再计算F1。这比单独输出CSV后再评测要高效得多因为它绕过了文件读写的开销也避免了键值类型不一致的坑。如果你准备把这份代码迁移到自己的项目建议保留这个join验证的做法。本文还有配套的精品资源点击获取
返回列表