metrics.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374
  1. """Versioned north-star ROI metric computation.
  2. This module only transforms pandas DataFrames. It does not know about policy,
  3. approval, Tencent writes, storage, scheduling, or Feishu.
  4. """
  5. from __future__ import annotations
  6. from typing import Dict, Iterable, List, Mapping, Sequence
  7. import numpy as np
  8. import pandas as pd
  9. from .fission_multiplier import (
  10. DISPLAY_MULTIPLIER_COLUMN,
  11. DISPLAY_TOTAL_TO_FIRST_COLUMN,
  12. FissionMultiplierParameters,
  13. apply_fission_multiplier,
  14. )
  15. METRIC_VERSION = "north_star_roi_t15_v8"
  16. METRIC_RUN_SUFFIX = "m8"
  17. SELF_CHANNEL = "小程序投流-稳定"
  18. GZH_CHANNEL = "公众号合作-即转-稳定"
  19. QIWEI_CHANNEL = "群/企微合作-稳定"
  20. ENTITY_SELF = "self"
  21. ENTITY_SELF_AD = "self_ad"
  22. ENTITY_GZH = "gzh"
  23. ENTITY_QIWEI = "qiwei"
  24. DIMENSION_COLUMNS = [
  25. "channel",
  26. "代理名称",
  27. "账号id",
  28. "账号名称",
  29. "广告id",
  30. "广告名称",
  31. "包名",
  32. "广告优化目标",
  33. "创意id",
  34. "合作方名",
  35. "公众号名",
  36. ]
  37. NUMERIC_COLUMNS = [
  38. "首层UV",
  39. "T0裂变数",
  40. "成本",
  41. "效率收入",
  42. "裂变效率收入",
  43. ]
  44. ENTITY_KEYS: Mapping[str, Sequence[str]] = {
  45. ENTITY_SELF: (
  46. "channel",
  47. "账号id",
  48. "广告id",
  49. "包名",
  50. "广告优化目标",
  51. "创意id",
  52. ),
  53. ENTITY_SELF_AD: (
  54. "channel",
  55. "账号id",
  56. "广告id",
  57. "包名",
  58. "广告优化目标",
  59. ),
  60. ENTITY_GZH: ("channel", "合作方名", "公众号名"),
  61. ENTITY_QIWEI: ("channel", "合作方名"),
  62. }
  63. def prepare_daily_metrics(
  64. raw: pd.DataFrame,
  65. fission_parameters: FissionMultiplierParameters,
  66. ) -> pd.DataFrame:
  67. """清洗日聚合结果并用当日裂变效率收入计算 T15 预测 ROI。"""
  68. raw = raw.rename(
  69. columns={
  70. "首层uv": "首层UV",
  71. "t0裂变数": "T0裂变数",
  72. }
  73. )
  74. required = {"dt", "entity_type", *DIMENSION_COLUMNS, *NUMERIC_COLUMNS}
  75. missing = sorted(required - set(raw.columns))
  76. if missing:
  77. raise ValueError(f"日聚合数据缺少字段: {', '.join(missing)}")
  78. daily = raw.copy()
  79. daily["dt"] = daily["dt"].astype(str)
  80. daily["entity_type"] = daily["entity_type"].astype(str)
  81. for column in DIMENSION_COLUMNS:
  82. daily[column] = daily[column].fillna("").astype(str)
  83. for column in NUMERIC_COLUMNS:
  84. daily[column] = pd.to_numeric(daily[column], errors="coerce").fillna(0.0)
  85. daily = apply_fission_multiplier(daily, fission_parameters)
  86. daily["T0实际裂变收入"] = daily["裂变效率收入"]
  87. daily["实际全链路效率收入"] = (
  88. daily["效率收入"] + daily["T0实际裂变收入"]
  89. )
  90. daily["预测T1-T15裂变收入"] = daily["T0实际裂变收入"] * (
  91. daily[DISPLAY_MULTIPLIER_COLUMN] - 1
  92. )
  93. daily["预测T0-T15裂变收入"] = (
  94. daily["T0实际裂变收入"] + daily["预测T1-T15裂变收入"]
  95. )
  96. daily["预测全链路效率收入"] = (
  97. daily["效率收入"] + daily["预测T0-T15裂变收入"]
  98. )
  99. daily["全链路效率收入"] = daily["预测全链路效率收入"]
  100. daily["实际ROI"] = np.where(
  101. daily["成本"] > 0,
  102. daily["实际全链路效率收入"] / daily["成本"],
  103. np.nan,
  104. )
  105. daily["ROI"] = np.where(
  106. daily["成本"] > 0,
  107. daily["预测全链路效率收入"] / daily["成本"],
  108. np.nan,
  109. )
  110. daily["T0裂变率"] = np.where(
  111. daily["首层UV"] > 0,
  112. daily["T0裂变数"] / daily["首层UV"],
  113. np.nan,
  114. )
  115. return daily
  116. def _age_map(ad_age: pd.DataFrame | None) -> Dict[str, int]:
  117. if ad_age is None or ad_age.empty:
  118. return {}
  119. required = {"广告id", "广告age"}
  120. missing = required - set(ad_age.columns)
  121. if missing:
  122. raise ValueError(f"广告age数据缺少字段: {', '.join(sorted(missing))}")
  123. ages = ad_age.copy()
  124. ages["广告id"] = ages["广告id"].fillna("").astype(str)
  125. ages["广告age"] = pd.to_numeric(ages["广告age"], errors="coerce").fillna(0).astype(int)
  126. return dict(zip(ages["广告id"], ages["广告age"]))
  127. def _expected_three_days(expected_dates: Iterable[str]) -> List[str]:
  128. dates = sorted({str(value) for value in expected_dates})
  129. if len(dates) != 3:
  130. raise ValueError(f"必须提供连续三个日期,实际为: {dates}")
  131. parsed = pd.to_datetime(dates, format="%Y%m%d")
  132. gaps = parsed.to_series().diff().dropna().dt.days.tolist()
  133. if gaps != [1, 1]:
  134. raise ValueError(f"日期必须连续,实际为: {dates}")
  135. return dates
  136. def _summarize_entity(
  137. entity_daily: pd.DataFrame,
  138. entity_type: str,
  139. expected_dates: Sequence[str],
  140. ages: Mapping[str, int],
  141. ) -> pd.DataFrame:
  142. keys = list(ENTITY_KEYS[entity_type])
  143. records: List[dict] = []
  144. for key_values, group in entity_daily.groupby(keys, dropna=False, sort=False):
  145. if not isinstance(key_values, tuple):
  146. key_values = (key_values,)
  147. parameter_columns = [
  148. DISPLAY_MULTIPLIER_COLUMN,
  149. DISPLAY_TOTAL_TO_FIRST_COLUMN,
  150. "传播裂变系数匹配层级",
  151. "传播裂变系数来源",
  152. "传播裂变参数版本",
  153. "传播裂变参数cohort日期",
  154. "传播裂变参数状态",
  155. "调控参与状态",
  156. "ROI计算口径",
  157. ]
  158. parameter_values = {
  159. column: group[column].dropna().unique()
  160. for column in parameter_columns
  161. }
  162. inconsistent = [
  163. column
  164. for column, values in parameter_values.items()
  165. if len(values) > 1
  166. or (len(values) == 0 and column != DISPLAY_TOTAL_TO_FIRST_COLUMN)
  167. ]
  168. if inconsistent:
  169. raise ValueError(
  170. f"同一实体三日内传播裂变参数不一致: {key_values}; fields={inconsistent}"
  171. )
  172. observed_dates = set(group["dt"].astype(str))
  173. by_date = group.groupby("dt", as_index=False)[
  174. [
  175. "首层UV",
  176. "T0裂变数",
  177. "成本",
  178. "效率收入",
  179. "裂变效率收入",
  180. "T0实际裂变收入",
  181. "实际全链路效率收入",
  182. "预测T1-T15裂变收入",
  183. "预测T0-T15裂变收入",
  184. "预测全链路效率收入",
  185. ]
  186. ].sum()
  187. by_date = (
  188. by_date.set_index("dt")
  189. .reindex(list(expected_dates), fill_value=0)
  190. )
  191. daily_actual_roi = np.where(
  192. by_date["成本"] > 0,
  193. by_date["实际全链路效率收入"] / by_date["成本"],
  194. np.nan,
  195. )
  196. daily_predicted_roi = np.where(
  197. by_date["成本"] > 0,
  198. by_date["预测全链路效率收入"] / by_date["成本"],
  199. np.nan,
  200. )
  201. record = dict(zip(keys, key_values))
  202. for column in DIMENSION_COLUMNS:
  203. if column in record:
  204. continue
  205. values = group.sort_values("dt")[column].dropna().astype(str)
  206. record[column] = next(
  207. (value for value in reversed(values.tolist()) if value),
  208. "",
  209. )
  210. record["entity_type"] = entity_type
  211. record["覆盖天数"] = len(observed_dates.intersection(expected_dates))
  212. record["窗口首层UV"] = float(by_date["首层UV"].sum())
  213. record["日均首层UV"] = float(by_date["首层UV"].mean())
  214. record["窗口最小单日首层UV"] = float(by_date["首层UV"].min())
  215. record["最新日首层UV"] = float(by_date.iloc[-1]["首层UV"])
  216. record["窗口T0裂变数"] = float(by_date["T0裂变数"].sum())
  217. record["T0裂变率"] = (
  218. record["窗口T0裂变数"] / record["窗口首层UV"]
  219. if record["窗口首层UV"] > 0
  220. else np.nan
  221. )
  222. record["成本"] = float(by_date["成本"].sum())
  223. record["窗口最小单日成本"] = float(by_date["成本"].min())
  224. record["效率收入"] = float(by_date["效率收入"].sum())
  225. record["裂变效率收入"] = float(
  226. by_date["裂变效率收入"].sum()
  227. )
  228. record["T0实际裂变收入"] = float(
  229. by_date["T0实际裂变收入"].sum()
  230. )
  231. for column, values in parameter_values.items():
  232. record[column] = values[0] if len(values) else np.nan
  233. record[DISPLAY_MULTIPLIER_COLUMN] = float(
  234. record[DISPLAY_MULTIPLIER_COLUMN]
  235. )
  236. if pd.notna(record[DISPLAY_TOTAL_TO_FIRST_COLUMN]):
  237. record[DISPLAY_TOTAL_TO_FIRST_COLUMN] = float(
  238. record[DISPLAY_TOTAL_TO_FIRST_COLUMN]
  239. )
  240. record["实际全链路效率收入"] = float(
  241. by_date["实际全链路效率收入"].sum()
  242. )
  243. record["预测T1-T15裂变收入"] = float(
  244. by_date["预测T1-T15裂变收入"].sum()
  245. )
  246. record["预测T0-T15裂变收入"] = float(
  247. by_date["预测T0-T15裂变收入"].sum()
  248. )
  249. record["预测全链路效率收入"] = float(
  250. by_date["预测全链路效率收入"].sum()
  251. )
  252. record["全链路效率收入"] = record["预测全链路效率收入"]
  253. record["实际ROI"] = (
  254. record["实际全链路效率收入"] / record["成本"]
  255. if record["成本"] > 0
  256. else np.nan
  257. )
  258. record["ROI"] = (
  259. record["预测全链路效率收入"] / record["成本"]
  260. if record["成本"] > 0
  261. else np.nan
  262. )
  263. window_days = float(len(expected_dates))
  264. record["日均T0裂变数"] = record["窗口T0裂变数"] / window_days
  265. record["日均成本"] = record["成本"] / window_days
  266. record["日均效率收入"] = record["效率收入"] / window_days
  267. record["日均T0裂变效率收入"] = (
  268. record["T0实际裂变收入"] / window_days
  269. )
  270. record["日均总预估效率收入"] = (
  271. record["预测全链路效率收入"] / window_days
  272. )
  273. record["三日加权平均T0裂变率"] = record["T0裂变率"]
  274. record["三日加权平均实际ROI"] = record["实际ROI"]
  275. record["三日加权平均效率ROI"] = record["ROI"]
  276. record["首日首层UV"] = float(by_date.iloc[0]["首层UV"])
  277. record["首日效率ROI"] = (
  278. float(daily_predicted_roi[0])
  279. if np.isfinite(daily_predicted_roi[0])
  280. else np.nan
  281. )
  282. record["最新日效率ROI"] = (
  283. float(daily_predicted_roi[-1])
  284. if np.isfinite(daily_predicted_roi[-1])
  285. else np.nan
  286. )
  287. for index, dt in enumerate(expected_dates):
  288. actual_roi = (
  289. float(daily_actual_roi[index])
  290. if np.isfinite(daily_actual_roi[index])
  291. else np.nan
  292. )
  293. predicted_roi = (
  294. float(daily_predicted_roi[index])
  295. if np.isfinite(daily_predicted_roi[index])
  296. else np.nan
  297. )
  298. record[f"首层UV_{dt}"] = float(by_date.iloc[index]["首层UV"])
  299. record[f"T0裂变数_{dt}"] = float(by_date.iloc[index]["T0裂变数"])
  300. record[f"成本_{dt}"] = float(by_date.iloc[index]["成本"])
  301. record[f"效率收入_{dt}"] = float(by_date.iloc[index]["效率收入"])
  302. record[f"裂变效率收入_{dt}"] = float(
  303. by_date.iloc[index]["裂变效率收入"]
  304. )
  305. record[f"预测全链路效率收入_{dt}"] = float(
  306. by_date.iloc[index]["预测全链路效率收入"]
  307. )
  308. record[f"实际ROI_{dt}"] = actual_roi
  309. record[f"ROI_{dt}"] = predicted_roi
  310. if entity_type in (ENTITY_SELF, ENTITY_SELF_AD):
  311. record["广告age"] = int(ages.get(str(record.get("广告id", "")), 0))
  312. records.append(record)
  313. return pd.DataFrame(records)
  314. def compute_roi_summary(
  315. raw_daily: pd.DataFrame,
  316. expected_dates: Iterable[str],
  317. ad_age: pd.DataFrame | None = None,
  318. *,
  319. fission_parameters: FissionMultiplierParameters,
  320. ) -> tuple[pd.DataFrame, list[str]]:
  321. """Return three-day entity ROI snapshots and normalized dates."""
  322. dates = _expected_three_days(expected_dates)
  323. daily = prepare_daily_metrics(raw_daily, fission_parameters)
  324. daily = daily[daily["dt"].isin(dates)].copy()
  325. ages = _age_map(ad_age)
  326. summaries = []
  327. for entity_type in (ENTITY_SELF, ENTITY_SELF_AD, ENTITY_GZH):
  328. entity_daily = daily[daily["entity_type"].eq(entity_type)]
  329. if entity_daily.empty:
  330. continue
  331. entity_summary = _summarize_entity(
  332. entity_daily,
  333. entity_type,
  334. dates,
  335. ages,
  336. )
  337. if not entity_summary.empty:
  338. summaries.append(entity_summary)
  339. if not summaries:
  340. raise ValueError("三个目标实体均没有可用的最近三日数据")
  341. summary = pd.concat(summaries, ignore_index=True, sort=False)
  342. return summary, dates