十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

移动推荐竞赛基线工程:从数据预处理到模型融合与OnSpark迁移

移动推荐竞赛基线工程:从数据预处理到模型融合与OnSpark迁移 简介压缩包整理自阿里天池新人赛“移动推荐”赛题的实战代码与方案非常适合数据挖掘入门者作为第一份完整参考。资源共26个文件压缩包大小约138KB其中19个Python脚本覆盖了数据预处理、特征工程、模型训练与预测等核心环节并提供了逻辑回归、梯度提升、XGBoost三种常用算法的实现同时包含基于Spark的OnSpark版本便于处理更大规模数据。另附README说明、特征列表等文本文件目录按数据分析、特征构造、模型等模块拆分阅读路径清晰。目前已有114人学习下载。通过研读这套代码可以学习到用户、商品、品牌三维特征的构造思路以及多模型融合和分布式调优的工程技巧帮助读者建立从原始数据到最终预测的完整认知为后续参加更复杂的竞赛打好基础。1. 天池“移动推荐”试玩工程先看懂这份完整基线再改代码从 zip 包里解压出 ali_freshman_compatiton-master 之后建议先打开 README.md 和 model 目录而不是急着找运行入口。这份工程对应的是阿里天池新人赛里的“移动推荐”赛题输入一段用户对商品的浏览、收藏、加购、购买行为目标是预测用户在指定日期对某款商品是否会产生购买行为。工程把数据预处理、特征构造、LR/GBDT/XGBoost 训练、基于 Spark 的 OnSpark 预测脚本全部串了起来是一套能跑通全流程的基线代码。对刚开始做推荐竞赛的人这套代码的价值在于把常见的特征套路整理成了可复用的脚本对做推荐开发的工程师OnSpark_model 目录演示了如何把 pandas 特征逻辑改写为分布式 Spark 任务适合用来对照学习。整份代码不是只能提交答案的片段而是从原始日志到最终预测结果的一条完整链路。注意工程里有几处文件名拼写不一致比如 onspark_data_preprocssing.py 少了字母 eREADME 里描述的命令和实际文件名也有出入。解压后建议先做一次改名避免后续引用时路径找不到。后面几章我会按“预处理 → 特征 → 模型 → 分布式预测”的顺序拆解并给出可以直接落地的代码片段。2. 数据预处理行为日志别直接用先做时间窗口切分2.1 移动推荐赛题的数据字段与行为类型原始行为日志通常包含四类字段用户 ID、商品 ID、品牌 ID、行为类型外加时间戳和地理位置。工程里 data_preanalysis 目录下的脚本会输出每个字段缺失率与取值分布这是第一个必须执行的步骤。推荐赛题的原始数据在离线导出时经常带重复点击、异常时间戳和空地理信息不先做统计后续特征构造会把脏数据带进去。字段含义典型取值预处理建议user_id用户标识长整型无需清洗但要做 ID 连续性检查item_id商品标识长整型存在商品被下架后行为清零的情况注意缺失brand_id品牌标识长整型会出现品牌 ID 相同但商品不同的情况behavior_type行为类型浏览/收藏/加购/购买要统一编码常见映射为 1/2/3/4timestamp行为时间unix 秒或日期字符串必须转成 datetime 并按天绑定geo地理位置城市编码缺失值建议单独建模不直接填充data_preanalysis 下的脚本一般会输出每个字段的 min、max、缺失比例以及行为类型计数。这里有一个常见问题行为类型编码在不同版本赛题里不一样工程里如果直接写死if behavior 4换一份数据就可能失效。我一般会把行为类型映射抽成一个常量字典统一放在脚本头部。import pandas as pd import numpy as np BEHAVIOR_ENCODING { view: 1, # 浏览 cart: 2, # 加购 favor: 3, # 收藏 buy: 4 # 购买 }这段代码把原始行为名映射成稳定数字编码。实际工程里原始日志可能是英文行为名也可能是数字编码统一做一次映射可以避免后面特征构造脚本里每个函数各自处理一次也能减少改参数时遗漏的坑。映射表的语义顺序建议沿用“浏览-加购-收藏-购买”后续计算转化率时可以直接按表名引用不用记住数字含义。2.2 时间窗口切分固定窗口比滑窗更可控预测目标决定了时间窗口的划分方式。天池移动推荐官方的做法是给出一段连续时间的行为数据要求预测某一天每个用户对某个商品是否购买。因此训练集的特征窗口不能用全部历史至少要留出最后的预测日作为验证。常见的做法是用前 N 天行为构造特征用第 N1 天的购买记录构造标签然后整体切分成训练和验证两部分。我在处理这类赛题时会先固定一个验证截止时间而不是随机切分。因为用户行为在时间上是强相关的随机切分会让特征窗口和标签窗口互相污染线上预测会严重失准。一个简单可用的切分逻辑如下# 假设行为日志已经读入 df_log df_log[dt] pd.to_datetime(df_log[timestamp], units) train_end pd.Timestamp(2014-12-15) label_date pd.Timestamp(2014-12-16) # 特征窗口截止到 train_end 当天 feature_data df_log[df_log[dt] train_end] # 标签窗口train_end 到 label_date 之间的购买行为 label_data df_log[(df_log[dt] train_end) (df_log[dt] label_date) (df_log[behavior] 4)]这里feature_data用来生成用户和商品的统计特征label_data用来构建正样本标签。注意train_end使用的是严格小于防止当天数据同时进入特征和标签。label_date是预测日推荐赛题里通常只预测一天所以标签窗口不能跨多天否则模型学到的购买概率是累计购买概率线上预测时阈值会漂移。时间窗口的参数一般写在配置文件中例如用FEATURE_DAYS7、LABEL_DAYS1。如果验证集效果AUC不稳定优先调整FEATURE_DAYS不要一上来就调模型参数。特征窗口太短会丢失长期偏好太长会把冷启动用户的特征噪声放大。2.3 数据清洗的两次去重行为日志里最常见的脏数据是重复点击。用户在同一分钟内反复刷新商品页会产生多条浏览记录如果不做去重商品特征里的浏览次数会翻倍而购买次数没有变化导致转化率被低估。去重逻辑一般分两层第一层是同一用户、同一商品、同一行为、同一秒内的重复记录直接去掉第二层是同一用户对同一商品在同一个小时内多次浏览只保留第一条。df_log df_log.drop_duplicates( subset[user_id, item_id, behavior, timestamp] ) # 构造小时标记用于第二层去重 df_log[behavior_hour] df_log[dt].dt.floor(h) df_log df_log.drop_duplicates( subset[user_id, item_id, behavior, behavior_hour], keepfirst )第二层去重先构造behavior_hour再用四个字段组合去重这样可以过滤掉短时间内刷屏导致的行为次数虚高但不会误删跨小时的真实重复行为。需要特别注意的是第二层去重只对浏览行为生效不要把购买行为按小时折叠用户在货架页连续购买多个商品是最典型的真实购买场景小时级去重会把这类购买误删。所以更稳妥的做法是在第二层去重前加条件df_log[df_log[behavior] ! buy]。提示第二层去重只对浏览行为生效。购买行为即使发生在同一秒也可能是不同订单必须保留。第2章覆盖了从字段理解、窗口切分到去重三个核心动作。预处理脚本输出一张清洗后的宽表后就可以进入特征构造阶段。3. 特征工程用户、商品、品牌三层特征的构造顺序特征工程的第一步是先决定特征挂在哪个层级上。工程里feature_construct目录下的脚本按聚合粒度分成四类我把它整理成下面的对应关系。特征层级聚合键代表特征侧重点用户层user_id行为占比、活跃天数用户活跃度商品层item_id近7天热度、转化率商品生命周期品牌层brand_id品牌下商品数、品牌热度品牌影响力用户-商品层user_id, item_id交互次数、间隔个性化意图3.1 用户维度特征行为次数占比比绝对值更稳定用户特征的核心是回答“这个用户平时喜欢怎么逛”。只看浏览次数控制不了用户活跃度的差异。比如一个用户有 1000 次浏览但只有 1 次购买另一个用户有 10 次浏览但买了 5 次前者的购买意愿未必更低。特征工程里更常用的是行为占比加购率、收藏率、购买转化率以及最近 7 天和最近 1 天的行为分布差异。工程里feature_construct目录下应该能看到按 user_id 分组的特征生成脚本。这个脚本输出的特征通常包括用户行为总数、各行为类型数量、行为时间跨度、活跃天数、最后一次行为距离预测日的小时数。其中“最后一次行为距离预测日的时间差”是点击率预估里非常有效的特征需要单独处理因为它不服从正态分布后续树模型可以直接用但线性模型需要做分桶或 log 变换。def build_user_features(log, label_dt, b_map): log log.copy() log[behavior] log[behavior].map(b_map) grouped log.groupby(user_id) u_feat pd.DataFrame() u_feat[u_cnt] grouped.size() for b, name in [(1, u_view), (2, u_cart), (3, u_favor), (4, u_buy)]: u_feat[name] grouped[behavior].apply(lambda s: (s b).sum()) u_feat[u_buy_rate] u_feat[u_buy] / (u_feat[u_view] 1) u_feat[u_last_dt] grouped[dt].max() u_feat[u_last_gap_hour] ( pd.Timestamp(label_dt) - u_feat[u_last_dt]).dt.total_seconds() / 3600 return u_feat这段代码里u_buy_rate的分母加了 1是为了防止除零。u_last_gap_hour表示用户最后一次活跃距离标签日的小时数这个值越小说明用户在预测日附近活跃重要性越高。特征列名加了u_前缀是为了后续模型融合时拼接多个特征文件能直接区分来源避免只有字段名没有来源导致调试困难。3.2 商品与品牌维度特征关注商品生命周期商品维度特征要回答“这个商品最近是不是在被大量加购”。商品本身存在生命周期刚上架的商品浏览多、购买少临下架的商品可能收藏多但购买集中在最后一两天。如果只用全量天数的平均转化率会把新品的真实热度压低。常见做法是分别计算最近 1 天、最近 3 天、最近 7 天的浏览量和购买量再组合成长短期热度比。品牌维度的处理思路不同。品牌 ID 本身是类别特征直接 one-hot 会导致维度过高。常见做法是把品牌聚合到商品上计算品牌下所有商品的购买率、品牌在预测日之前最后一条行为记录的时间。工程里单独建了onspark_generate_feature_user_brand.py说明品牌特征在作者看来值得单独生成一份而不是放进商品特征里一起聚合。def build_brand_features(log, b_map): log log.copy() log[behavior] log[behavior].map(b_map) b log.groupby(brand_id).agg( b_cnt(behavior, count), b_buy(behavior, lambda s: (s 4).sum()), b_item_cnt(item_id, nunique), b_last_dt(dt, max) ) b[b_buy_rate] b[b_buy] / (b[b_cnt] 1) return b这里把品牌下唯一商品数b_item_cnt也作为特征。品牌下商品过少可能意味着这是一个新品牌或小众品牌它的行为总量即使很大也可能来自单个爆款不能代表品牌整体热度。树模型能看到这个字段后会在特征分裂时把“单爆款品牌”和“多商品品牌”区分开。b_last_dt在后续样本拼接时要转成距离预测日的时长否则不同日期预测时这个字段没有可比性。3.3 用户-商品交叉特征推荐模型的上限所在用户特征和商品特征单独做出来之后模型只能学到“喜欢购物的用户”和“热门的商品”真正能抓住点击率预估本质的是用户和商品组合的行为。比如用户之前是否加购过这个商品、是否收藏过、浏览后多久再次出现这些信息比单纯的用户活跃度和商品热度要直接得多。交叉特征通常用groupby([user_id, item_id])生成核心字段是用户-商品交互次数、最后一次交互时间、中间隔了多少天。构造交叉特征时要特别注意样本量。用户-商品组合的正样本非常稀疏直接统计购买次数很容易让模型在训练集上记住 ID但对没有历史交互的组合失效。工程里的做法是生成两类特征一类是具体组合的历史统计另一类是用户最近交互过的商品列表的向量化表示。后一种在工程脚本里可能体现在feature_list.txt中里面会列出每个特征所属维度。def build_ui_features(log, b_map): log log.copy() log[behavior] log[behavior].map(b_map) grouped log.groupby([user_id, item_id]) ui grouped.agg( ui_view(behavior, lambda s: (s 1).sum()), ui_cart(behavior, lambda s: (s 2).sum()), ui_favor(behavior, lambda s: (s 3).sum()), ui_buy(behavior, lambda s: (s 4).sum()), ui_first_dt(dt, min), ui_last_dt(dt, max), ) ui[ui_days] (ui[ui_last_dt] - ui[ui_first_dt]).dt.days return ui.reset_index()ui_days表示用户从第一次到最后一次交互这个商品跨越的天数天数越长说明用户对商品持续关注可能是反复比价。这个特征对“收藏很久没有买”和“第一次浏览就买”两种行为有区分作用。特征构造完成之后需要统一把类别字段转换为整数索引再分别保存训练集特征、验证集特征和预测集特征模型读取时直接按feature_list.txt对齐列名即可。4. 模型训练LR、GBDT、XGBoost 的参数与融合策略4.1 三个模型在推荐场景里的分工工程目录名model_lr_and_gdbt_and_xgboost说明作者把三个模型都放进了同一套特征里。为什么不是只选一个效果最好的因为在移动推荐赛题这种高维稀疏场景里LR 的作用是保证模型能在线性维度上兜底GBDT 用来捕获非线性组合XGBoost 则负责通过正则化防止过拟合。三个模型的预测结果做简单融合后AUC 通常能比单模型高 0.5 到 1 个点这部分提升来自模型结构差异而不是调参。LR 需要自己构造交叉特征因为线性模型不会自动完成特征组合GBDT 和 XGBoost 对原始统计特征更友好可以直接吃连续值。但 GBDT 在训练时会把特征切分成阶梯函数对稀疏类别特征很敏感如果 one-hot 后维度太高训练速度会大幅下降。所以这套工程里特征工程生成的统计特征大多是数值型刚好适合 GBDT。下面表格整理了三种模型在工程中常用到的参数配置具体数值来自我对这个赛题的复现经验可以作为调节起点。模型关键参数设置值调参方向LRsolverliblinear数据量大时换 sagaLRC1.0过拟合时减小到 0.1GBDTn_estimators100根据验证集早停GBDTmax_depth5增大到 7 可能过拟合XGBoosteta0.05从 0.1 减小到 0.03XGBoostmax_depth6与 subsample 配合调整XGBoostsubsample0.8防止样本抖动XGBoostcolsample_bytree0.8特征多时降低到 0.6参数表里的值不是固定的关键在于理解每个参数影响的是偏差还是方差。eta调小之后必须增加num_boost_round否则欠拟合max_depth增大后会覆盖更多高阶特征组合但训练集表现会快速上升验证集表现反而下降。4.2 训练脚本的骨架特征文件、标签、早停工程里 model 目录下的脚本一般会先读取feature_list.txt指定特征列再读取特征生成阶段输出的训练和验证特征文件。这一步有一点容易踩坑训练集和验证集的特征数量必须完全一致feature_list.txt里多一个或少一个列都会导致 XGBoost 报特征维度不一致。所以特征构造和模型训练之间要有一个对齐脚本把两份特征文件统一按同一份特征列名排序。import pandas as pd import xgboost as xgb train pd.read_csv(train_features.csv, dtype{user_id: int32}) val pd.read_csv(val_features.csv, dtype{user_id: int32}) features [c for c in train.columns if c not in [user_id, item_id, label]] y_train train[label].values dtrain xgb.DMatrix(train[features], labely_train) dval xgb.DMatrix(val[features], labelval[label].values) params { objective: binary:logistic, eval_metric: auc, eta: 0.05, max_depth: 6, subsample: 0.8, colsample_bytree: 0.8, } bst xgb.train( params, dtrain, num_boost_round500, evals[(dtrain, train), (dval, val)], early_stopping_rounds50, verbose_eval50 )代码里early_stopping_rounds50是最关键的参数它让模型在验证集 AUC 连续 50 轮不增长时停止训练避免对训练集过拟合。evals同时输出训练集和验证集指标如果看到训练集 AUC 很高但验证集 AUC 在 100 轮后开始下降说明特征噪声太大需要回特征工程层而不是继续调参。verbose_eval50表示每 50 轮打印一次指标方便在长时间训练时观察趋势。4.3 简单加权融合与提交结果三个模型各自输出[user_id, item_id, prob]之后融合不一定非要上 Stacking。天池新人赛的分数差距往往在小数点后第三位Stacking 需要再训练一层模型容易过拟合。更稳妥的做法是先在验证集上用网格搜索寻找权重比如遍历w1, w2, w3使线性加权后验证集 AUC 最大。权重的搜索步长可以是 0.1一般两轮就能收敛。from sklearn.metrics import roc_auc_score best_weight, best_auc None, 0 for w1 in [0.2, 0.3, 0.4]: for w2 in [0.2, 0.3, 0.4]: w3 1 - w1 - w2 if w3 0: continue blend w1 * pred_lr w2 * pred_gbdt w3 * pred_xgb auc roc_auc_score(val[label], blend) if auc best_auc: best_auc, best_weight auc, (w1, w2, w3) print(best_weight, best_auc)融合权重搜出来之后可以对最终概率做一次 rank 归一化再搜索权重。因为不同模型输出的概率分布差异很大LR 的输出普遍接近 0XGBoost 的输出会在 0.01 到 0.9 之间波动直接相加相当于给了 XGBoost 更大的权重。把每个模型的预测值按行排序转换成百分位后再加权权重搜索结果会更稳定。5. OnSpark 迁移大规模特征生成的分布式改造与排错5.1 为什么要单独维护一套 OnSpark 脚本单机 pandas 在 10 万用户、1000 万行为样本的规模下勉强能跑但特征维度一旦超过几百列groupby([user_id, item_id])的内存占用会迅速膨胀。工程里 OnSpark_model 目录专门存放了基于 Spark 的特征生成和预测脚本目的是在分布式环境下处理更大数据集。onspark_data_preprocssing.py对应单机版的预处理onspark_generate_feature_user_product.py、onspark_generate_feature_user_brand.py分别生成不同维度的特征最终的onspark_prediction.py输出预测结果。这里要注意拼写目录里是onspark_data_preprocssing.py少了一个 e如果是在集群上跑官方 README 里的命令需要改成和实际文件名一致。更推荐的做法是先把这几个脚本统一加一个spark_前缀至少可以避免调用时手滑打错。5.2 从 pandas groupby 到 Spark 聚合的改写思路pandas 里groupby([user_id, item_id]).agg(...)到了 Spark 上不能直接平移因为 Spark DataFrame API 对多聚合的计算图有更严格的类型约束。工程里的onspark_generate_feature_user_product.py大概率用了groupBy加agg的组合关键点是把每个行为类型转换成独立的求和列而不是在lambda里逐行判断。from pyspark.sql.functions import when, lit, max as spark_max from pyspark.sql.functions import sum as spark_sum ui log.groupBy(user_id, item_id).agg( spark_sum(when(log[behavior] view, lit(1)).otherwise(lit(0))).alias(ui_view), spark_sum(when(log[behavior] buy, lit(1)).otherwise(lit(0))).alias(ui_buy), spark_max(dt).alias(ui_last_dt) )上面这段用pyspark.sql.functions实现。when和otherwise的组合相当于把单机版的布尔判断替换成 Spark SQL 的条件表达式这样sum才能作用于整列。这里最大的差异在于pandas 的groupby.apply会把整个分组的数据拉回到 driver 端而agg表达式会先在每个 executor 上做部分聚合再把结果合并数据传输量完全不同。5.3 分布式环境下的两个典型坑第一个坑是特征对齐。单机版feature_list.txt里列的顺序就是训练脚本读取的顺序但在 Spark 上如果使用feature_list.txt选取列很容易因为分区顺序不同导致列顺序不一致。我一般会在onspark_prediction.py中显式按照feature_list.txt重新排序 DataFrame 列并在生成预测结果前打印一次前 5 行列名确认没有丢失列。第二个坑是时间戳类型。Spark 读取 parquet 时会把时间戳转成整型秒而 pandas 读 CSV 得到的是 datetime64两者的时间差会导致last_gap_hour出现几百倍的偏差。解决办法是统一在数据预处理阶段把timestamp转成 long 类型的秒数所有特征计算都基于这个整数最后预测日再转换成可读时间用于验证。这个改动看着小但对树模型的 AUC 影响可以达到 1 到 2 个点是隐藏的高价值排错点。OnSpark 脚本在集群上跑通后可以把模型导出成 PMML 或者保存在 HDFS 上。工程里的onspark_prediction.py输出结果通常以user_id, item_id, prob的格式落盘到指定目录提交前再按概率降序排序。由于只预测一个固定日期可以先按用户分组用窗口函数取每个用户概率最高的前 N 个商品避免把全量预测结果写到单机后再做二次排序。本文还有配套的精品资源点击获取
返回列表