简介:本资源是一套基于真实电商场景的Python广告推荐系统源码,面向机器学习初学者、推荐算法实践者及数据科学爱好者,聚焦点击率(CTR)预估这一核心广告技术问题。项目完整复现了从离线召回(ALS协同过滤)、特征工程(品类/品牌打分、Redis缓存设计)到在线服务(实时推荐接口、特征存储)的全流程,依托阿里天池淘宝展示广告数据集构建可运行的端到端方案。压缩包共13个文件,含12个Python模块(如ALS_recall.py、brand_scoring.py、onlineRecommend.py等)与1份README.md说明文档,总大小仅21KB,轻量易读、结构清晰,便于理解各模块职责与调用关系。已有916人学习下载,读者可直接获取完整的工程化代码框架、Redis集成实践、多维度特征构造逻辑及线上/线下双路推荐架构设计思路,是深入理解工业级推荐系统落地细节的优质入门范例。
1. 电商广告推荐不是“猜用户喜欢什么”,而是用 Python 把曝光、点击、转化三类行为串成可训练的信号链
很多刚接触电商推荐的同学以为,写个协同过滤或调个 LightGBM 就能上线跑广告推荐——结果在真实业务里,模型 AUC 高但线上 CTR 不升反降,或者离线指标漂亮,AB 实验却显示 ROI 下滑。问题不在算法本身,而在于电商广告场景的特殊性:用户一次会话中可能浏览 20+ 商品,但只点击 1~2 个,最终下单更少;广告位有首焦、搜索底纹、详情页关联等多个异构位置;预算、出价、实时竞价策略又和模型预测强耦合。这个标题里的“Python电商广告推荐系统源码”不是玩具 demo,它必须包含行为日志清洗→特征工程→多目标建模→在线服务封装四层闭环。适合两类人:一是正在搭建电商业务中台的数据工程师,需要可复用的特征 pipeline 和模型 serving 框架;二是算法工程师,想跳过从零搭环境的环节,直接基于真实字段结构(如user_id,ad_id,position_type,expose_time,click_time,pay_amount)调试多任务 loss 权重。下面我们就从数据源头开始,把这套系统真正跑通。
2. 用 Pandas + PySpark 构建电商广告行为日志清洗流水线:从原始埋点到可建模宽表
电商广告日志不是规整 CSV,而是由 Nginx 日志、客户端 SDK 上报、服务端打点三路汇聚的半结构化数据。常见字段包括event_type(expose/click/pay)、user_id(加密 ID)、ad_id(广告创意 ID)、position(广告位编码,如 “search_bottom_3”)、ts(毫秒级时间戳)、page_url(当前页面路径)、device_type(iOS/Android/H5)。直接用 Pandas 加载全量日志会 OOM,必须分层处理。
2.1 原始日志解析与 schema 对齐
假设原始日志为 JSON 行式文件(每行一个事件),需先统一字段命名和类型。注意user_id和ad_id在不同端可能用不同加密方式,需做映射对齐:
# log_parser.py import json import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.functions import col, from_json, to_timestamp, when, lit from pyspark.sql.types import StructType, StructField, StringType, LongType, IntegerType # 定义统一 schema(关键字段必须显式声明,避免 inferSchema 失败) schema = StructType([ StructField("event_type", StringType(), True), StructField("user_id", StringType(), True), StructField("ad_id", StringType(), True), StructField("position", StringType(), True), StructField("ts", LongType(), True), # 毫秒时间戳 StructField("page_url", StringType(), True), StructField("device_type", StringType(), True), StructField("session_id", StringType(), True), ]) def parse_raw_log(spark, input_path): # 读取原始 JSON 日志(支持压缩格式如 .json.gz) df = spark.read.option("multiLine", "false").json(input_path, schema=schema) # 转换时间戳并标准化 event_type df = df.withColumn("event_time", to_timestamp(col("ts") / 1000.0)) \ .withColumn("event_type", when(col("event_type") == "show", "expose") .when(col("event_type") == "clk", "click") .otherwise(col("event_type"))) # 过滤无效事件(空 user_id/ad_id 或非法时间) df = df.filter( (col("user_id").isNotNull()) & (col("ad_id").isNotNull()) & (col("ts") > 1609459200000) & # 2021-01-01 00:00:00 UTC (col("ts") < 2524608000000) # 2050-01-01 00:00:00 UTC ) return df # 使用示例(本地小样本调试) spark = SparkSession.builder.appName("ad-log-parse").getOrCreate() raw_df = parse_raw_log(spark, "hdfs://namenode:8020/logs/ad_raw_202406*.json.gz") raw_df.select("event_type", "user_id", "ad_id", "position", "event_time").show(5)提示:
to_timestamp(col("ts") / 1000.0)是关键,电商日志普遍用毫秒时间戳,Spark 默认按秒解析会错位 3 位。filter中的时间范围硬约束能避免脏数据污染后续特征计算。
2.2 构建用户-广告会话级宽表:曝光→点击→转化的时序归因
单条日志无法建模,必须聚合为“用户在某次会话中对某广告的完整交互链”。核心逻辑是:以user_id + session_id + ad_id为粒度,提取首次曝光时间、是否点击、是否下单、点击延迟(ms)、下单金额等字段。这里用窗口函数比 groupby 更稳:
# session_aggregator.py from pyspark.sql.window import Window from pyspark.sql.functions import ( min, max, count, sum, first, last, row_number, collect_list, struct, datediff, expr, when, lit ) def build_session_wide_table(raw_df): # 步骤1:按 user_id + session_id + ad_id 排序事件,标记序号 window_spec = Window.partitionBy("user_id", "session_id", "ad_id").orderBy("ts") df_with_rank = raw_df.withColumn("event_rank", row_number().over(window_spec)) # 步骤2:提取每个组合的首次曝光时间、最后点击时间、是否发生点击/支付 agg_df = df_with_rank.groupBy("user_id", "session_id", "ad_id").agg( min("ts").alias("first_expose_ts"), max(when(col("event_type") == "click", col("ts"))).alias("last_click_ts"), max(when(col("event_type") == "pay", col("ts"))).alias("last_pay_ts"), sum(when(col("event_type") == "expose", 1).otherwise(0)).alias("expose_cnt"), sum(when(col("event_type") == "click", 1).otherwise(0)).alias("click_cnt"), sum(when(col("event_type") == "pay", 1).otherwise(0)).alias("pay_cnt"), # 提取支付金额(假设 pay 事件带 amount 字段,需提前 join 支付明细表) max(when(col("event_type") == "pay", col("pay_amount"))).alias("pay_amount") ) # 步骤3:计算关键衍生指标 wide_df = agg_df.withColumn("is_clicked", when(col("click_cnt") > 0, 1).otherwise(0)) \ .withColumn("is_converted", when(col("pay_cnt") > 0, 1).otherwise(0)) \ .withColumn("click_delay_ms", when(col("is_clicked") == 1, col("last_click_ts") - col("first_expose_ts")) .otherwise(-1)) \ .withColumn("expose_to_click_ratio", col("click_cnt") / col("expose_cnt")) \ .withColumn("position_type", when(col("position").contains("search"), "search") .when(col("position").contains("detail"), "detail") .otherwise("other")) return wide_df # 执行聚合 session_df = build_session_wide_table(raw_df) session_df.select( "user_id", "ad_id", "position_type", "expose_cnt", "is_clicked", "is_converted", "click_delay_ms", "pay_amount" ).show(5)注意:
click_delay_ms必须用毫秒差值,而非datediff(天数),因为电商广告点击常发生在曝光后几秒内。position_type的规则需根据实际埋点规范调整,例如"home_banner_1"应归为"home",不能简单用split("_")[0]——首页 banner 和首页 tab 是完全不同的流量质量。
2.3 特征工程基础表:用户侧、广告侧、上下文侧三类特征生成
宽表只是起点,推荐模型需要结构化特征。我们按维度拆解:
| 维度 | 字段示例 | 计算方式 | 更新频率 |
|---|---|---|---|
| 用户侧 | user_age_group,user_pv_7d,user_click_rate_30d | 基于用户历史行为聚合 | T+1 离线 |
| 广告侧 | ad_ctr_7d,ad_cvr_30d,ad_category | 基于广告 ID 聚合 | T+1 离线 |
| 上下文侧 | hour_of_day,is_weekend,device_type,position_type | 直接取自宽表 | 实时 |
# feature_generator.py from pyspark.sql.functions import ( hour, dayofweek, when, col, avg, stddev, count, sum, broadcast, concat_ws ) def generate_user_features(spark, session_df): # 用户 7 天曝光数、点击率、平均点击延迟 user_agg = session_df.groupBy("user_id").agg( sum("expose_cnt").alias("user_expose_7d"), sum("click_cnt").alias("user_click_7d"), avg("click_delay_ms").alias("user_avg_click_delay_7d"), stddev("click_delay_ms").alias("user_std_click_delay_7d") ).withColumn("user_ctr_7d", col("user_click_7d") / col("user_expose_7d")) # 用户设备偏好(统计各设备曝光占比) device_dist = session_df.groupBy("user_id", "device_type").agg( count("*").alias("device_cnt") ).withColumn("total_device_cnt", sum("device_cnt").over(Window.partitionBy("user_id"))) device_ratio = device_dist.withColumn("device_ratio", col("device_cnt") / col("total_device_cnt")) \ .filter(col("device_ratio") > 0.5) \ .select("user_id", "device_type") return user_agg.join(device_ratio, on="user_id", how="left") def generate_ad_features(session_df): # 广告 30 天 CTR/CVR、所属类目(需 join 广告主提供的 ad_info 表) ad_agg = session_df.groupBy("ad_id").agg( sum("expose_cnt").alias("ad_expose_30d"), sum("click_cnt").alias("ad_click_30d"), sum("pay_cnt").alias("ad_pay_30d"), avg("pay_amount").alias("ad_avg_pay_30d") ).withColumn("ad_ctr_30d", col("ad_click_30d") / col("ad_expose_30d")) \ .withColumn("ad_cvr_30d", col("ad_pay_30d") / col("ad_click_30d")) # 关联广告元信息(假设已存在 ad_info 表) ad_info_df = spark.table("ad_info") # 包含 ad_id, category_id, advertiser_id, budget_type return ad_agg.join(ad_info_df, on="ad_id", how="left") # 合并三类特征 user_feat = generate_user_features(spark, session_df) ad_feat = generate_ad_features(session_df) context_feat = session_df.select( "user_id", "ad_id", hour("first_expose_ts").alias("hour_of_day"), (dayofweek("first_expose_ts") >= 6).cast("int").alias("is_weekend"), "device_type", "position_type" ) final_feature_df = user_feat.join(ad_feat, on="ad_id", how="inner") \ .join(context_feat, on=["user_id", "ad_id"], how="inner")关键参数说明:
user_ctr_7d分母用sum("expose_cnt")而非count(*),因为单次会话可能多次曝光同一广告;ad_cvr_30d的分母是ad_click_30d而非ad_expose_30d,这是 CVR 的定义本质——点击后的转化率。若某广告 30 天无点击,则ad_cvr_30d为null,后续模型需做fillna(0)处理。
3. 构建多目标深度推荐模型:用 PyTorch 实现曝光预估 + 点击预估 + 转化预估联合训练
电商广告不能只优化点击率(CTR),否则会推高点击但拉低 ROI;也不能只优化转化率(CVR),因为没曝光就没转化。工业界主流方案是ESMM(Entire Space Multi-Task Model):用两个子网络分别建模 pCTR 和 pCVR,再通过 pCTCVR = pCTR × pCVR 实现联合优化。我们用 PyTorch 从零实现,不依赖 DeepCTR 等第三方库,确保可调试、可解释。
3.1 数据准备:将 Spark DataFrame 转为 PyTorch Dataset
特征列需对齐,数值型做 min-max 归一化,类别型做 embedding lookup:
# dataset.py import torch from torch.utils.data import Dataset, DataLoader import numpy as np from sklearn.preprocessing import MinMaxScaler class AdDataset(Dataset): def __init__(self, spark_df, numerical_cols, categorical_cols, label_cols): # 转为 Pandas(仅限中小规模数据,大数据用 Arrow + TorchArrow) pdf = spark_df.toPandas() # 数值特征归一化 self.scaler = MinMaxScaler() self.numerical_data = self.scaler.fit_transform(pdf[numerical_cols].values) # 类别特征转索引(需提前构建 vocab) self.categorical_data = {} for col in categorical_cols: # 假设已构建好 vocab: {"user_id": {id1: 0, id2: 1, ...}, "ad_id": {...}} vocab = load_vocab(col) # 从 HDFS 加载预存 vocab dict self.categorical_data[col] = pdf[col].map(vocab.get).fillna(0).astype(int).values # 标签:expose(是否曝光)、click(是否点击)、pay(是否支付) self.labels = { "expose": pdf[label_cols[0]].values, "click": pdf[label_cols[1]].values, "pay": pdf[label_cols[2]].values } def __len__(self): return len(self.labels["expose"]) def __getitem__(self, idx): return { "numerical": torch.tensor(self.numerical_data[idx], dtype=torch.float32), "categorical": {k: torch.tensor(v[idx], dtype=torch.long) for k, v in self.categorical_data.items()}, "labels": { "expose": torch.tensor(self.labels["expose"][idx], dtype=torch.float32), "click": torch.tensor(self.labels["click"][idx], dtype=torch.float32), "pay": torch.tensor(self.labels["pay"][idx], dtype=torch.float32) } } # 使用示例 numerical_cols = ["user_expose_7d", "user_ctr_7d", "ad_ctr_30d", "ad_cvr_30d", "click_delay_ms"] categorical_cols = ["user_id", "ad_id", "position_type", "device_type"] label_cols = ["expose_cnt", "click_cnt", "pay_cnt"] # 注意:此处用 cnt > 0 判断 bool train_dataset = AdDataset(final_feature_df, numerical_cols, categorical_cols, label_cols) train_loader = DataLoader(train_dataset, batch_size=1024, shuffle=True, num_workers=4)3.2 ESMM 模型定义:共享底层 + 独立塔 + 乘法约束
模型结构必须体现“曝光是点击的前提,点击是转化的前提”这一业务逻辑:
# model.py import torch import torch.nn as nn import torch.nn.functional as F class ESMM(nn.Module): def __init__(self, num_numerical, cat_dims, # [(user_id_vocab_size, 8), (ad_id_vocab_size, 16), ...] hidden_dims=[128, 64, 32], dropout=0.2): super().__init__() # 共享底层:数值特征 + 类别 embedding 拼接 self.numerical_bn = nn.BatchNorm1d(num_numerical) self.embedding_layers = nn.ModuleList([ nn.Embedding(cat_dim, embed_dim) for cat_dim, embed_dim in cat_dims ]) # 共享 DNN layers = [] in_dim = num_numerical + sum([embed_dim for _, embed_dim in cat_dims]) for h_dim in hidden_dims: layers.extend([ nn.Linear(in_dim, h_dim), nn.BatchNorm1d(h_dim), nn.ReLU(), nn.Dropout(dropout) ]) in_dim = h_dim self.shared_mlp = nn.Sequential(*layers) # pCTR 塔(曝光→点击) self.ctr_tower = nn.Sequential( nn.Linear(hidden_dims[-1], 64), nn.ReLU(), nn.Dropout(dropout), nn.Linear(64, 1), nn.Sigmoid() ) # pCVR 塔(点击→转化) self.cvr_tower = nn.Sequential( nn.Linear(hidden_dims[-1], 64), nn.ReLU(), nn.Dropout(dropout), nn.Linear(64, 1), nn.Sigmoid() ) def forward(self, numerical, categorical): # 数值特征归一化(已在 Dataset 中完成,此处仅 BN) x_num = self.numerical_bn(numerical) # 类别特征 embedding x_cat = [] for i, (cat_name, cat_idx) in enumerate(categorical.items()): emb = self.embedding_layers[i](cat_idx) x_cat.append(emb) x_cat = torch.cat(x_cat, dim=1) # 拼接 & 共享 MLP x = torch.cat([x_num, x_cat], dim=1) shared_out = self.shared_mlp(x) # 双塔输出 pctr = self.ctr_tower(shared_out).squeeze(-1) # shape: [B] pcvr = self.cvr_tower(shared_out).squeeze(-1) # shape: [B] pctcvr = pctr * pcvr # ESMM 核心:pCTCVR = pCTR * pCVR return {"pctr": pctr, "pcvr": pcvr, "pctcvr": pctcvr} # 初始化模型 model = ESMM( num_numerical=len(numerical_cols), cat_dims=[ (100000, 16), # user_id vocab size, embed dim (50000, 16), # ad_id (10, 4), # position_type (3, 2) # device_type ], hidden_dims=[128, 64] )为什么用乘法而非加法?因为
pCTCVR = pCTR × pCVR是概率论中的链式法则,物理意义明确:用户看到广告(pCTR)且在看到后点击并支付(pCVR)的联合概率。若用pCTR + pCVR,则当两者都为 0.8 时结果为 1.6,违反概率定义。
3.3 多任务损失函数与训练循环:平衡曝光、点击、转化三目标
损失函数需加权,且pCTCVR的监督信号来自pay标签,pCTR来自expose标签,pCVR无直接标签(需用pay / click伪标签):
# train.py import torch.optim as optim from torch.nn import BCELoss def esmm_loss(preds, labels, alpha=0.2, beta=0.3, gamma=0.5): """ ESMM loss: L = alpha * L_ctr + beta * L_cvr + gamma * L_ctcvr 其中 L_ctr = BCE(pctr, expose_label) L_cvr = BCE(pcvr, cvr_label) —— cvr_label = pay_label / click_label (需平滑) L_ctcvr = BCE(pctcvr, pay_label) """ bce = BCELoss() # pCTR loss: 曝光是点击的前提,所以用 expose_label 监督 pctr l_ctr = bce(preds["pctr"], labels["expose"]) # pCVR loss: 构造 pseudo-CVR label # 当 click=0 时,cvr_label 设为 0(避免除零);当 click=1 且 pay=1 时为 1,否则为 0 cvr_label = torch.where( labels["click"] == 1, labels["pay"], torch.zeros_like(labels["pay"]) ) l_cvr = bce(preds["pcvr"], cvr_label) # pCTCVR loss: 直接监督支付 l_ctcvr = bce(preds["pctcvr"], labels["pay"]) return alpha * l_ctr + beta * l_cvr + gamma * l_ctcvr # 训练循环 optimizer = optim.Adam(model.parameters(), lr=1e-3) model.train() for epoch in range(10): total_loss = 0 for batch in train_loader: optimizer.zero_grad() preds = model(batch["numerical"], batch["categorical"]) loss = esmm_loss(preds, batch["labels"]) loss.backward() optimizer.step() total_loss += loss.item() print(f"Epoch {epoch+1}, Loss: {total_loss/len(train_loader):.4f}")参数调优重点:
alpha,beta,gamma的初始值按业务目标设定。若当前阶段要保曝光量,alpha设为 0.5;若要提升 ROI,gamma应 >beta。实践中发现beta=0.1效果更好——因为pCVR标签稀疏(点击用户中只有 ~5% 会支付),过高的beta会导致模型过度拟合点击用户。
4. 模型部署与在线服务:用 Flask + ONNX 实现低延迟广告打分 API
训练好的模型不能停留在 Jupyter 里。电商广告要求P99 延迟 < 50ms,且需支持 AB 测试分流。我们放弃重量级 Serving 框架(如 Triton),用轻量方案:PyTorch → ONNX → Flask。
4.1 导出 ONNX 模型并验证一致性
ONNX 是跨框架中间表示,能规避 PyTorch 版本兼容问题:
# export_onnx.py import torch.onnx # 创建 dummy input(必须与实际推理 shape 一致) dummy_numerical = torch.randn(1, len(numerical_cols)) dummy_categorical = { "user_id": torch.tensor([12345], dtype=torch.long), "ad_id": torch.tensor([67890], dtype=torch.long), "position_type": torch.tensor([1], dtype=torch.long), "device_type": torch.tensor([0], dtype=torch.long) } # 导出 torch.onnx.export( model, (dummy_numerical, dummy_categorical), "esmm_model.onnx", input_names=["numerical", "user_id", "ad_id", "position_type", "device_type"], output_names=["pctr", "pcvr", "pctcvr"], dynamic_axes={ "numerical": {0: "batch_size"}, "user_id": {0: "batch_size"}, "ad_id": {0: "batch_size"}, "pctr": {0: "batch_size"}, "pctcvr": {0: "batch_size"} }, opset_version=12 ) # 验证 ONNX 输出与 PyTorch 一致 import onnxruntime as ort ort_session = ort.InferenceSession("esmm_model.onnx") def to_numpy(tensor): return tensor.detach().cpu().numpy() if tensor.requires_grad else tensor.cpu().numpy() # PyTorch 推理 with torch.no_grad(): torch_out = model(dummy_numerical, dummy_categorical) # ONNX 推理 ort_inputs = { "numerical": to_numpy(dummy_numerical), "user_id": to_numpy(dummy_categorical["user_id"]), "ad_id": to_numpy(dummy_categorical["ad_id"]), "position_type": to_numpy(dummy_categorical["position_type"]), "device_type": to_numpy(dummy_categorical["device_type"]) } ort_outs = ort_session.run(None, ort_inputs) # 比较 print("PyTorch pctr:", torch_out["pctr"].item()) print("ONNX pctr:", ort_outs[0][0]) assert abs(torch_out["pctr"].item() - ort_outs[0][0]) < 1e-44.2 Flask API 实现:支持批量打分与特征校验
API 需处理特征缺失、ID 未登录等异常,并返回结构化 JSON:
# app.py from flask import Flask, request, jsonify import onnxruntime as ort import numpy as np from typing import List, Dict, Any app = Flask(__name__) ort_session = ort.InferenceSession("esmm_model.onnx") # 加载特征 scaler 和 vocab(与训练时一致) scaler = load_scaler() # MinMaxScaler fitted on training data vocab_dict = load_vocab_dict() # {"user_id": {id_str: idx}, ...} def preprocess_request(data: Dict[str, Any]) -> Dict[str, np.ndarray]: """将 JSON 请求转为 ONNX 输入格式""" numerical = np.array([ data.get("user_expose_7d", 0), data.get("user_ctr_7d", 0), data.get("ad_ctr_30d", 0), data.get("ad_cvr_30d", 0), data.get("click_delay_ms", 0) ]).reshape(1, -1) # 归一化 numerical_scaled = scaler.transform(numerical) # 类别特征转索引(未登录 ID 映射为 0) categorical = { "user_id": np.array([vocab_dict["user_id"].get(str(data.get("user_id", "")), 0)]), "ad_id": np.array([vocab_dict["ad_id"].get(str(data.get("ad_id", "")), 0)]), "position_type": np.array([{"search": 0, "detail": 1, "home": 2}.get(data.get("position_type", "other"), 0)]), "device_type": np.array([{"iOS": 0, "Android": 1, "H5": 2}.get(data.get("device_type", "H5"), 2)]) } return {"numerical": numerical_scaled.astype(np.float32), **categorical} @app.route("/score", methods=["POST"]) def score_ad(): try: req_data = request.get_json() if not req_data: return jsonify({"error": "Empty request body"}), 400 # 批量支持(可选) if isinstance(req_data, list): results = [] for item in req_data: inputs = preprocess_request(item) ort_outs = ort_session.run(None, inputs) results.append({ "pctr": float(ort_outs[0][0]), "pcvr": float(ort_outs[1][0]), "pctcvr": float(ort_outs[2][0]) }) return jsonify(results) # 单条 inputs = preprocess_request(req_data) ort_outs = ort_session.run(None, inputs) return jsonify({ "pctr": float(ort_outs[0][0]), "pcvr": float(ort_outs[1][0]), "pctcvr": float(ort_outs[2][0]), "score": float(ort_outs[2][0]) # 默认用 pCTCVR 作为排序分 }) except Exception as e: return jsonify({"error": str(e)}), 500 if __name__ == "__main__": app.run(host="0.0.0.0", port=5000, threaded=True)性能关键点:
threaded=True启用多线程,避免 Flask 默认单线程阻塞;ort_session全局复用,避免重复加载模型;preprocess_request中所有操作均为 NumPy 向量化,无 Python 循环。
4.3 压测与监控:用 Locust 验证 QPS 与延迟
部署前必须压测。Locust 脚本模拟真实请求:
# locustfile.py from locust import HttpUser, task, between import json import random class AdScorerUser(HttpUser): wait_time = between(0.1, 0.5) # 每次请求间隔 100~500ms @task def score_ad(self): payload = { "user_id": str(random.randint(1000, 999999)), "ad_id": str(random.randint(1000, 99999)), "user_expose_7d": random.randint(0, 500), "user_ctr_7d": round(random.random(), 4), "ad_ctr_30d": round(random.random(), 4), "ad_cvr_30d": round(random.random(), 4), "click_delay_ms": random.randint(0, 10000), "position_type": random.choice(["search", "detail", "home"]), "device_type": random.choice(["iOS", "Android", "H5"]) } with self.client.post("/score", json=payload, catch_response=True) as response: if response.status_code != 200: response.failure(f"HTTP {response.status_code}") # 启动命令:locust -f locustfile.py --host http://localhost:5000压测结果示例(4 核 CPU,16GB 内存):
| 并发用户数 | QPS | P99 延迟 | 错误率 |
|---|---|---|---|
| 100 | 185 | 32ms | 0% |
| 500 | 892 | 41ms | 0% |
| 1000 | 1620 | 48ms | <0.1% |
注意:若 P99 > 50ms,优先检查
preprocess_request中的字典查找(vocab_dict["user_id"].get(...))是否用了 Python dict 而非更高效的defaultdict或numpy.searchsorted。线上环境建议将 vocab 加载到 Redis 缓存。
5. 线上效果验证与 AB 实验设计:用双重差分法归因广告推荐收益
模型上线后,不能只看离线 AUC。电商广告的核心指标是千次曝光收入(RPM)和广告 ROI。必须通过 AB 实验隔离模型效果,且需控制混杂变量。
5.1 实验分组与流量分配:基于用户 Hash 的稳定分流
避免按请求随机分,用user_id的 MD5 哈希后取模,保证同一用户始终进入同组:
# ab_utils.py import hashlib def get_ab_group(user_id: str, group_count: int = 4) -> str: """基于 user_id 稳定哈希分组""" hash_val = int(hashlib.md5(user_id.encode()).hexdigest()[:8], 16) group_id = hash_val % group_count return ["control", "treatment_a", "treatment_b", "treatment_c"][group_id] # 示例:用户 u12345 → md5("u12345") = "d41d8cd98f00b204e9800998ecf8427e" → 0xd41d8cd9 = 3558740185 → 3558740185 % 4 = 1 → "treatment_a"5.2 核心指标定义与双重差分(DID)分析
AB 实验易受时间趋势干扰(如大促期间自然增长)。用 DID 消除外部因素:
| 组别 | 实验期 RPM | 对照期 RPM | 差值 |
|---|---|---|---|
| Control | $R_{c1}$ | $R_{c0}$ | $R_{c1} - R_{c0}$ |
| Treatment | $R_{t1}$ | $R_{t0}$ | $R_{t1} - R_{t0}$ |
DID 估计量:$\Delta = (R_{t1} - R_{t0}) - (R_{c1} - R_{c0})$
# did_analysis.py import pandas as pd from scipy import stats def calculate_did(df: pd.DataFrame): """ df columns: user_id, group (control/treatment), period (before/after), rpm """ # 按组和时期聚合 agg = df.groupby(["group", "period"])["rpm"].agg(["mean", "std", "count"]).reset_index() # 提取四象限均值 control_before = agg[(agg["group"]=="control") & (agg["period"]=="before")]["mean"].iloc[0] control_after = agg[(agg["group"]=="control") & (agg["period"]=="after")]["mean"].iloc[0] treatment_before = agg[(agg["group"]=="treatment <p> <a href="https://download.csdn.net/download/weixin_47367099/85542199" style="color:#ec7500;font-size:14px;"> 本文还有配套的精品资源,点击获取 </a> <img alt="menu-r.4af5f7ec.gif" src="https://csdnimg.cn/release/wenkucmsfe/public/img/menu-r.4af5f7ec.gif" style="width:16px;margin-left:4px;vertical-align:text-bottom;cursor:text;"> </p>