creative_roi_calculator.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407
  1. """
  2. 创意级动态 ROI 计算器 — auto_put_ad_mini
  3. 把 roi_calculator.py 的"动态 ROI 7 日均值"公式按 (ad_id, creative_id) 维度下沉,
  4. 为创意级 pause 决策提供数值基础。
  5. 核心公式(与 roi_calculator.py 完全一致,仅聚合维度从 ad_id 改为 (ad_id, creative_id)):
  6. T0裂变系数 = SUM(fission0_count) / SUM(open_count)
  7. arpu = SUM(total_revenue) / SUM(total_return_count)
  8. 当日裂变收益率 = SUM(fission0_count) * arpu / SUM(cost)
  9. 当日回流倍数 = SUM(total_return_count) / SUM(open_count)
  10. T0裂变系数_7日均值 = mean(T0裂变系数) over 7 天
  11. 回流倍数_7日均值 = mean(当日回流倍数) over 7 天
  12. 裂变效率稳定因子 = 回流倍数_7日均值 / T0裂变系数_7日均值
  13. 创意动态ROI = 当日裂变收益率 × 裂变效率稳定因子
  14. 创意动态ROI_7日均值 = mean(创意动态ROI) over 7 天 ← 决策参考值
  15. 前置条件:
  16. - 单日 (ad_id, creative_id) 消耗 < 100 元的天数不参与计算(NaN)
  17. - min_periods=3:至少 3 天合格数据才计算 7 日滚动均值
  18. ⚠️ 归因语义说明:
  19. ODPS 表 loghubods.ad_put_tencent_creative_data_day 的 fission0_count / total_return_count /
  20. total_revenue 等字段,本期假设按"创意级独立归因"处理(即同一用户被多创意触达时归到首次触达
  21. 的 creative_id)。如果实际是按曝光数加权拆分到所有创意,需要另外修正聚合逻辑。
  22. 本期先按现有口径实现,待阶段 6 端到端测试时通过对比"创意 ROI 加权平均"与"广告级 ROI"做交叉验证。
  23. """
  24. import logging
  25. from datetime import datetime, timedelta
  26. from pathlib import Path
  27. from typing import Optional
  28. import numpy as np
  29. import pandas as pd
  30. from agent.tools import tool
  31. from agent.tools.models import ToolContext, ToolResult
  32. logger = logging.getLogger(__name__)
  33. _MINI_DIR = Path(__file__).resolve().parent.parent
  34. _MERGED_DIR = _MINI_DIR / "outputs" / "merged"
  35. _CREATIVE_ROI_DIR = _MINI_DIR / "outputs" / "creative_roi"
  36. # ===== 内部聚合 =====
  37. def _aggregate_creative_to_creative_day(df: pd.DataFrame) -> pd.DataFrame:
  38. """按 (ad_id, creative_id, date) 聚合(同一天可能多分片,需 SUM)。"""
  39. if df.empty:
  40. return pd.DataFrame()
  41. df = df.copy()
  42. # bizdate → date
  43. if "bizdate" in df.columns:
  44. df["date"] = df["bizdate"].astype(str)
  45. elif "date" not in df.columns:
  46. logger.warning("creative_roi: DataFrame 缺少 bizdate/date 列")
  47. return pd.DataFrame()
  48. # 列名标准化(与 roi_calculator 对齐)
  49. COLUMN_RENAME = {
  50. "首层小程序打开数": "open_count",
  51. "裂变0层回流数": "fission0_count",
  52. "裂变层回流数": "fission_count",
  53. "裂变1层回流数": "fission1_count",
  54. "总回流人数": "total_return_count",
  55. "总收入": "total_revenue",
  56. "ad_status": "configured_status",
  57. }
  58. rename_map = {k: v for k, v in COLUMN_RENAME.items() if k in df.columns}
  59. df = df.rename(columns=rename_map)
  60. # 过滤无 creative_id 的行(广告状态行)
  61. df = df[df["creative_id"].notna() & (df["creative_id"].astype(str).str.strip() != "")]
  62. if df.empty:
  63. return pd.DataFrame()
  64. # 数值字段安全转换
  65. numeric_cols = [
  66. "cost", "view_count", "valid_click_count",
  67. "open_count", "fission0_count", "fission_count", "fission1_count",
  68. "total_return_count", "total_revenue",
  69. ]
  70. for col in numeric_cols:
  71. if col in df.columns:
  72. df[col] = pd.to_numeric(df[col], errors="coerce").fillna(0)
  73. agg_dict = {
  74. "account_id": "first",
  75. "ad_name": "first",
  76. "creative_name": "first",
  77. "create_time": "first",
  78. "configured_status": "first",
  79. "package_name": "first",
  80. }
  81. for col in numeric_cols:
  82. if col in df.columns:
  83. agg_dict[col] = "sum"
  84. agg_dict = {k: v for k, v in agg_dict.items() if k in df.columns}
  85. grouped = df.groupby(["ad_id", "creative_id", "date"], as_index=False).agg(agg_dict)
  86. return grouped
  87. # ===== 创意级动态 ROI =====
  88. def _calculate_creative_dynamic_roi(
  89. cdf: pd.DataFrame,
  90. min_daily_cost: float = 100.0,
  91. ) -> pd.DataFrame:
  92. """
  93. 在 (ad_id, creative_id, date) 粒度上计算动态 ROI。
  94. """
  95. if cdf.empty:
  96. return cdf
  97. cdf = cdf.sort_values(["ad_id", "creative_id", "date"]).reset_index(drop=True)
  98. group_keys = ["ad_id", "creative_id"]
  99. # 当日基础指标(单日消耗不足时设 NaN)
  100. cdf["T0裂变系数"] = np.where(
  101. (cdf.get("open_count", 0) > 0) & (cdf.get("cost", 0) >= min_daily_cost),
  102. cdf["fission0_count"] / cdf["open_count"].replace(0, np.nan),
  103. np.nan,
  104. )
  105. cdf["arpu"] = np.where(
  106. (cdf.get("total_return_count", 0) > 0) & (cdf.get("cost", 0) >= min_daily_cost),
  107. cdf["total_revenue"] / cdf["total_return_count"].replace(0, np.nan),
  108. np.nan,
  109. )
  110. cdf["当日裂变收益率"] = np.where(
  111. (cdf.get("cost", 0) > 0) & (cdf.get("cost", 0) >= min_daily_cost),
  112. cdf["fission0_count"] * cdf["arpu"] / cdf["cost"].replace(0, np.nan),
  113. np.nan,
  114. )
  115. cdf["当日回流倍数"] = np.where(
  116. (cdf.get("open_count", 0) > 0) & (cdf.get("cost", 0) >= min_daily_cost),
  117. cdf["total_return_count"] / cdf["open_count"].replace(0, np.nan),
  118. np.nan,
  119. )
  120. # 7 日滚动均值(按创意分组,min_periods=3)
  121. cdf["T0裂变系数_7日均值"] = (
  122. cdf.groupby(group_keys)["T0裂变系数"]
  123. .transform(lambda x: x.rolling(window=7, min_periods=3).mean())
  124. )
  125. cdf["回流倍数_7日均值"] = (
  126. cdf.groupby(group_keys)["当日回流倍数"]
  127. .transform(lambda x: x.rolling(window=7, min_periods=3).mean())
  128. )
  129. cdf["裂变效率稳定因子"] = np.where(
  130. cdf["T0裂变系数_7日均值"] > 0,
  131. cdf["回流倍数_7日均值"] / cdf["T0裂变系数_7日均值"],
  132. np.nan,
  133. )
  134. cdf["创意动态ROI"] = cdf["当日裂变收益率"] * cdf["裂变效率稳定因子"]
  135. cdf["创意动态ROI_7日均值"] = (
  136. cdf.groupby(group_keys)["创意动态ROI"]
  137. .transform(lambda x: x.rolling(window=7, min_periods=3).mean())
  138. )
  139. cdf["roi_valid_days"] = (
  140. cdf.groupby(group_keys)["创意动态ROI"]
  141. .transform(lambda x: x.notna().sum())
  142. )
  143. return cdf
  144. def _build_creative_summary(
  145. cdf: pd.DataFrame,
  146. end_date: str,
  147. ) -> pd.DataFrame:
  148. """
  149. 按 (ad_id, creative_id) 汇总最新一天指标 + 7 日累计 + 创意年龄 + 占比。
  150. """
  151. if cdf.empty:
  152. return pd.DataFrame()
  153. end_dt = datetime.strptime(end_date, "%Y%m%d")
  154. start_dt_7d = end_dt - timedelta(days=6)
  155. start_date_7d = start_dt_7d.strftime("%Y%m%d")
  156. # 最近 7 天累计消耗
  157. df_7d = cdf[(cdf["date"] >= start_date_7d) & (cdf["date"] <= end_date)].copy()
  158. cost_7d = df_7d.groupby(["ad_id", "creative_id"], as_index=False)["cost"].sum()
  159. cost_7d.rename(columns={"cost": "cost_7d"}, inplace=True)
  160. # 广告级 7 日累计(用来算占比)
  161. ad_cost_7d = df_7d.groupby("ad_id", as_index=False)["cost"].sum()
  162. ad_cost_7d.rename(columns={"cost": "ad_cost_7d"}, inplace=True)
  163. summary = cost_7d.merge(ad_cost_7d, on="ad_id", how="left")
  164. summary["cost_share_7d"] = np.where(
  165. summary["ad_cost_7d"] > 0,
  166. (summary["cost_7d"] / summary["ad_cost_7d"]).round(4),
  167. 0.0,
  168. )
  169. # 创意年龄:以创意首次出现日期(min bizdate)为锚点,相对 end_date 计算
  170. first_date = (
  171. cdf.groupby(["ad_id", "creative_id"], as_index=False)["date"]
  172. .min()
  173. .rename(columns={"date": "first_date"})
  174. )
  175. def _age(row):
  176. """计算创意年龄(首次出现日期到 end_date 的天数),解析失败返回 None。"""
  177. try:
  178. first_dt = datetime.strptime(str(row["first_date"]), "%Y%m%d")
  179. return max((end_dt - first_dt).days, 0)
  180. except Exception:
  181. return None
  182. first_date["creative_age_days"] = first_date.apply(_age, axis=1)
  183. summary = summary.merge(first_date[["ad_id", "creative_id", "creative_age_days"]],
  184. on=["ad_id", "creative_id"], how="left")
  185. # 最新一天的创意动态 ROI + 创意属性
  186. latest = cdf[cdf["date"] == end_date][[
  187. c for c in [
  188. "ad_id", "creative_id", "creative_name", "ad_name", "account_id",
  189. "configured_status", "创意动态ROI", "创意动态ROI_7日均值", "roi_valid_days",
  190. ] if c in cdf.columns
  191. ]].copy()
  192. summary = summary.merge(latest, on=["ad_id", "creative_id"], how="left")
  193. # 兜底:创意如果在 end_date 当天没数据,从最近一天回填属性
  194. missing_mask = summary["creative_name"].isna() if "creative_name" in summary.columns else None
  195. if missing_mask is not None and missing_mask.any():
  196. last_seen = (
  197. cdf.sort_values("date")
  198. .groupby(["ad_id", "creative_id"], as_index=False)
  199. .last()[[c for c in [
  200. "ad_id", "creative_id", "creative_name", "ad_name", "account_id",
  201. "configured_status", "创意动态ROI_7日均值", "roi_valid_days",
  202. ] if c in cdf.columns]]
  203. .rename(columns={
  204. "creative_name": "_creative_name_fb",
  205. "ad_name": "_ad_name_fb",
  206. "account_id": "_account_id_fb",
  207. "configured_status": "_configured_status_fb",
  208. "创意动态ROI_7日均值": "_roi_7d_fb",
  209. "roi_valid_days": "_roi_valid_days_fb",
  210. })
  211. )
  212. summary = summary.merge(last_seen, on=["ad_id", "creative_id"], how="left")
  213. for col, fb in [
  214. ("creative_name", "_creative_name_fb"),
  215. ("ad_name", "_ad_name_fb"),
  216. ("account_id", "_account_id_fb"),
  217. ("configured_status", "_configured_status_fb"),
  218. ("创意动态ROI_7日均值", "_roi_7d_fb"),
  219. ("roi_valid_days", "_roi_valid_days_fb"),
  220. ]:
  221. if col in summary.columns and fb in summary.columns:
  222. summary[col] = summary[col].where(summary[col].notna(), summary[fb])
  223. summary.drop(columns=[fb], inplace=True)
  224. summary["roi_valid_days"] = summary["roi_valid_days"].fillna(0).astype(int)
  225. return summary
  226. # ===== 工具入口 =====
  227. @tool(description="计算创意级动态 ROI(7 日均值),用于创意级 pause 决策")
  228. async def calculate_creative_roi(
  229. ctx: ToolContext = None,
  230. end_date: str = "yesterday",
  231. min_daily_cost: float = 100.0,
  232. window_days: int = 30,
  233. ) -> ToolResult:
  234. """
  235. 创意级动态 ROI 计算工具。
  236. 工作流:
  237. 1. 加载最近 window_days 天的 merged_*.csv
  238. 2. 按 (ad_id, creative_id, date) 聚合
  239. 3. 计算每天的 T0 裂变系数 / arpu / 当日裂变收益率 / 当日回流倍数
  240. 4. 计算 7 日滚动均值 + 裂变效率稳定因子
  241. 5. 计算创意动态 ROI / 创意动态 ROI 7 日均值
  242. 6. 输出 outputs/creative_roi/creative_roi_{end_date}.csv
  243. Args:
  244. end_date: 结束日期(YYYYMMDD 或 "yesterday")
  245. min_daily_cost: 单日消耗门槛(默认 100 元),低于此值的天数不参与
  246. window_days: 加载历史窗口天数(默认 30 天,与 roi_calculator 一致)
  247. Returns:
  248. ToolResult,包含 csv_path / 创意总数 / eligible 创意数 / 全体均值
  249. """
  250. try:
  251. # 解析日期
  252. if end_date == "yesterday":
  253. end_dt = datetime.now() - timedelta(days=1)
  254. else:
  255. end_dt = datetime.strptime(end_date.replace("-", ""), "%Y%m%d")
  256. end_date_str = end_dt.strftime("%Y%m%d")
  257. # 加载 merged
  258. start_dt = end_dt - timedelta(days=window_days - 1)
  259. merged_dfs = []
  260. for i in range(window_days):
  261. d = (start_dt + timedelta(days=i)).strftime("%Y%m%d")
  262. csv = _MERGED_DIR / f"merged_{d}.csv"
  263. if not csv.exists():
  264. continue
  265. df = pd.read_csv(csv, dtype={"ad_id": str, "creative_id": str, "account_id": str})
  266. merged_dfs.append(df)
  267. if not merged_dfs:
  268. return ToolResult(
  269. title="创意级 ROI 计算失败",
  270. output=f"未找到任何 merged 数据({_MERGED_DIR})",
  271. )
  272. creative_df = pd.concat(merged_dfs, ignore_index=True)
  273. logger.info("创意级 ROI: 加载 merged 数据 %d 行(%d 天)", len(creative_df), len(merged_dfs))
  274. # 同步 roi_calculator 的"近 7 天累计消耗 = 0 视为已关闭"前置过滤
  275. last_7_start = (end_dt - timedelta(days=6)).strftime("%Y%m%d")
  276. bz = creative_df["bizdate"].astype(str)
  277. recent7 = creative_df[(bz >= last_7_start) & (bz <= end_date_str)]
  278. zero_ads = (
  279. recent7.groupby("ad_id")["cost"].sum()
  280. .pipe(lambda s: s[s.fillna(0) <= 0].index.tolist())
  281. )
  282. if zero_ads:
  283. before = len(creative_df)
  284. creative_df = creative_df[~creative_df["ad_id"].isin(zero_ads)].reset_index(drop=True)
  285. logger.info(
  286. "创意级 ROI: 前置过滤 %d 条近 7 天 0 消耗广告,creative_df %d → %d 行",
  287. len(zero_ads), before, len(creative_df),
  288. )
  289. # 聚合到 (ad_id, creative_id, date)
  290. cdf = _aggregate_creative_to_creative_day(creative_df)
  291. if cdf.empty:
  292. return ToolResult(
  293. title="创意级 ROI 计算失败",
  294. output="creative-day 聚合结果为空(可能 creative_id 全为空)",
  295. )
  296. logger.info("创意级 ROI: 聚合到 (ad_id, creative_id, date) %d 行", len(cdf))
  297. # 计算动态 ROI
  298. cdf = _calculate_creative_dynamic_roi(cdf, min_daily_cost)
  299. # 汇总到 (ad_id, creative_id)
  300. summary = _build_creative_summary(cdf, end_date_str)
  301. if summary.empty:
  302. return ToolResult(
  303. title="创意级 ROI 计算失败",
  304. output="creative summary 为空",
  305. )
  306. # 输出列排序
  307. out_cols = [
  308. "ad_id", "account_id", "ad_name", "creative_id", "creative_name",
  309. "configured_status", "creative_age_days",
  310. "cost_7d", "ad_cost_7d", "cost_share_7d",
  311. "创意动态ROI", "创意动态ROI_7日均值", "roi_valid_days",
  312. ]
  313. out_cols = [c for c in out_cols if c in summary.columns]
  314. summary = summary[out_cols].copy()
  315. # 保存
  316. _CREATIVE_ROI_DIR.mkdir(parents=True, exist_ok=True)
  317. out_path = _CREATIVE_ROI_DIR / f"creative_roi_{end_date_str}.csv"
  318. summary.to_csv(out_path, index=False, encoding="utf-8-sig")
  319. logger.info("创意级 ROI CSV 已保存: %s", out_path)
  320. # 统计
  321. total = len(summary)
  322. roi_col = "创意动态ROI_7日均值"
  323. valid = int(summary[roi_col].notna().sum()) if roi_col in summary.columns else 0
  324. roi_mean = float(summary[roi_col].mean()) if valid > 0 else float("nan")
  325. ad_count = int(summary["ad_id"].nunique())
  326. lines = [
  327. f"✅ 创意级动态 ROI 计算完成(截至 {end_date_str})",
  328. f" 输出文件:{out_path}",
  329. f" 创意总数:{total}(覆盖 {ad_count} 条广告)",
  330. f" 有 ROI 值的创意数(7 日均值非 NaN):{valid}",
  331. f" 创意动态 ROI_7 日均值 全体均值:{roi_mean:.4f}",
  332. ]
  333. return ToolResult(
  334. title=f"创意级 ROI 计算完成({total} 个创意)",
  335. output="\n".join(lines),
  336. metadata={
  337. "csv_path": str(out_path),
  338. "total_creatives": total,
  339. "valid_creatives": valid,
  340. "ad_count": ad_count,
  341. "roi_mean": roi_mean,
  342. "end_date": end_date_str,
  343. },
  344. )
  345. except Exception as e:
  346. logger.exception("calculate_creative_roi 失败")
  347. return ToolResult(title="创意级 ROI 计算异常", output=f"错误:{e}")