"""Versioned north-star ROI metric computation. This module only transforms pandas DataFrames. It does not know about policy, approval, Tencent writes, storage, scheduling, or Feishu. """ from __future__ import annotations from typing import Dict, Iterable, List, Mapping, Sequence import numpy as np import pandas as pd from .fission_multiplier import ( DISPLAY_MULTIPLIER_COLUMN, DISPLAY_TOTAL_TO_FIRST_COLUMN, FissionMultiplierParameters, apply_fission_multiplier, ) METRIC_VERSION = "north_star_roi_t15_v8" METRIC_RUN_SUFFIX = "m8" SELF_CHANNEL = "小程序投流-稳定" GZH_CHANNEL = "公众号合作-即转-稳定" QIWEI_CHANNEL = "群/企微合作-稳定" ENTITY_SELF = "self" ENTITY_SELF_AD = "self_ad" ENTITY_GZH = "gzh" ENTITY_QIWEI = "qiwei" DIMENSION_COLUMNS = [ "channel", "代理名称", "账号id", "账号名称", "广告id", "广告名称", "包名", "广告优化目标", "创意id", "合作方名", "公众号名", ] NUMERIC_COLUMNS = [ "首层UV", "T0裂变数", "成本", "效率收入", "裂变效率收入", ] ENTITY_KEYS: Mapping[str, Sequence[str]] = { ENTITY_SELF: ( "channel", "账号id", "广告id", "包名", "广告优化目标", "创意id", ), ENTITY_SELF_AD: ( "channel", "账号id", "广告id", "包名", "广告优化目标", ), ENTITY_GZH: ("channel", "合作方名", "公众号名"), ENTITY_QIWEI: ("channel", "合作方名"), } def prepare_daily_metrics( raw: pd.DataFrame, fission_parameters: FissionMultiplierParameters, ) -> pd.DataFrame: """清洗日聚合结果并用当日裂变效率收入计算 T15 预测 ROI。""" raw = raw.rename( columns={ "首层uv": "首层UV", "t0裂变数": "T0裂变数", } ) required = {"dt", "entity_type", *DIMENSION_COLUMNS, *NUMERIC_COLUMNS} missing = sorted(required - set(raw.columns)) if missing: raise ValueError(f"日聚合数据缺少字段: {', '.join(missing)}") daily = raw.copy() daily["dt"] = daily["dt"].astype(str) daily["entity_type"] = daily["entity_type"].astype(str) for column in DIMENSION_COLUMNS: daily[column] = daily[column].fillna("").astype(str) for column in NUMERIC_COLUMNS: daily[column] = pd.to_numeric(daily[column], errors="coerce").fillna(0.0) daily = apply_fission_multiplier(daily, fission_parameters) daily["T0实际裂变收入"] = daily["裂变效率收入"] daily["实际全链路效率收入"] = ( daily["效率收入"] + daily["T0实际裂变收入"] ) daily["预测T1-T15裂变收入"] = daily["T0实际裂变收入"] * ( daily[DISPLAY_MULTIPLIER_COLUMN] - 1 ) daily["预测T0-T15裂变收入"] = ( daily["T0实际裂变收入"] + daily["预测T1-T15裂变收入"] ) daily["预测全链路效率收入"] = ( daily["效率收入"] + daily["预测T0-T15裂变收入"] ) daily["全链路效率收入"] = daily["预测全链路效率收入"] daily["实际ROI"] = np.where( daily["成本"] > 0, daily["实际全链路效率收入"] / daily["成本"], np.nan, ) daily["ROI"] = np.where( daily["成本"] > 0, daily["预测全链路效率收入"] / daily["成本"], np.nan, ) daily["T0裂变率"] = np.where( daily["首层UV"] > 0, daily["T0裂变数"] / daily["首层UV"], np.nan, ) return daily def _age_map(ad_age: pd.DataFrame | None) -> Dict[str, int]: if ad_age is None or ad_age.empty: return {} required = {"广告id", "广告age"} missing = required - set(ad_age.columns) if missing: raise ValueError(f"广告age数据缺少字段: {', '.join(sorted(missing))}") ages = ad_age.copy() ages["广告id"] = ages["广告id"].fillna("").astype(str) ages["广告age"] = pd.to_numeric(ages["广告age"], errors="coerce").fillna(0).astype(int) return dict(zip(ages["广告id"], ages["广告age"])) def _expected_three_days(expected_dates: Iterable[str]) -> List[str]: dates = sorted({str(value) for value in expected_dates}) if len(dates) != 3: raise ValueError(f"必须提供连续三个日期,实际为: {dates}") parsed = pd.to_datetime(dates, format="%Y%m%d") gaps = parsed.to_series().diff().dropna().dt.days.tolist() if gaps != [1, 1]: raise ValueError(f"日期必须连续,实际为: {dates}") return dates def _summarize_entity( entity_daily: pd.DataFrame, entity_type: str, expected_dates: Sequence[str], ages: Mapping[str, int], ) -> pd.DataFrame: keys = list(ENTITY_KEYS[entity_type]) records: List[dict] = [] for key_values, group in entity_daily.groupby(keys, dropna=False, sort=False): if not isinstance(key_values, tuple): key_values = (key_values,) parameter_columns = [ DISPLAY_MULTIPLIER_COLUMN, DISPLAY_TOTAL_TO_FIRST_COLUMN, "传播裂变系数匹配层级", "传播裂变系数来源", "传播裂变参数版本", "传播裂变参数cohort日期", "传播裂变参数状态", "调控参与状态", "ROI计算口径", ] parameter_values = { column: group[column].dropna().unique() for column in parameter_columns } inconsistent = [ column for column, values in parameter_values.items() if len(values) > 1 or (len(values) == 0 and column != DISPLAY_TOTAL_TO_FIRST_COLUMN) ] if inconsistent: raise ValueError( f"同一实体三日内传播裂变参数不一致: {key_values}; fields={inconsistent}" ) observed_dates = set(group["dt"].astype(str)) by_date = group.groupby("dt", as_index=False)[ [ "首层UV", "T0裂变数", "成本", "效率收入", "裂变效率收入", "T0实际裂变收入", "实际全链路效率收入", "预测T1-T15裂变收入", "预测T0-T15裂变收入", "预测全链路效率收入", ] ].sum() by_date = ( by_date.set_index("dt") .reindex(list(expected_dates), fill_value=0) ) daily_actual_roi = np.where( by_date["成本"] > 0, by_date["实际全链路效率收入"] / by_date["成本"], np.nan, ) daily_predicted_roi = np.where( by_date["成本"] > 0, by_date["预测全链路效率收入"] / by_date["成本"], np.nan, ) record = dict(zip(keys, key_values)) for column in DIMENSION_COLUMNS: if column in record: continue values = group.sort_values("dt")[column].dropna().astype(str) record[column] = next( (value for value in reversed(values.tolist()) if value), "", ) record["entity_type"] = entity_type record["覆盖天数"] = len(observed_dates.intersection(expected_dates)) record["窗口首层UV"] = float(by_date["首层UV"].sum()) record["日均首层UV"] = float(by_date["首层UV"].mean()) record["窗口最小单日首层UV"] = float(by_date["首层UV"].min()) record["最新日首层UV"] = float(by_date.iloc[-1]["首层UV"]) record["窗口T0裂变数"] = float(by_date["T0裂变数"].sum()) record["T0裂变率"] = ( record["窗口T0裂变数"] / record["窗口首层UV"] if record["窗口首层UV"] > 0 else np.nan ) record["成本"] = float(by_date["成本"].sum()) record["窗口最小单日成本"] = float(by_date["成本"].min()) record["效率收入"] = float(by_date["效率收入"].sum()) record["裂变效率收入"] = float( by_date["裂变效率收入"].sum() ) record["T0实际裂变收入"] = float( by_date["T0实际裂变收入"].sum() ) for column, values in parameter_values.items(): record[column] = values[0] if len(values) else np.nan record[DISPLAY_MULTIPLIER_COLUMN] = float( record[DISPLAY_MULTIPLIER_COLUMN] ) if pd.notna(record[DISPLAY_TOTAL_TO_FIRST_COLUMN]): record[DISPLAY_TOTAL_TO_FIRST_COLUMN] = float( record[DISPLAY_TOTAL_TO_FIRST_COLUMN] ) record["实际全链路效率收入"] = float( by_date["实际全链路效率收入"].sum() ) record["预测T1-T15裂变收入"] = float( by_date["预测T1-T15裂变收入"].sum() ) record["预测T0-T15裂变收入"] = float( by_date["预测T0-T15裂变收入"].sum() ) record["预测全链路效率收入"] = float( by_date["预测全链路效率收入"].sum() ) record["全链路效率收入"] = record["预测全链路效率收入"] record["实际ROI"] = ( record["实际全链路效率收入"] / record["成本"] if record["成本"] > 0 else np.nan ) record["ROI"] = ( record["预测全链路效率收入"] / record["成本"] if record["成本"] > 0 else np.nan ) window_days = float(len(expected_dates)) record["日均T0裂变数"] = record["窗口T0裂变数"] / window_days record["日均成本"] = record["成本"] / window_days record["日均效率收入"] = record["效率收入"] / window_days record["日均T0裂变效率收入"] = ( record["T0实际裂变收入"] / window_days ) record["日均总预估效率收入"] = ( record["预测全链路效率收入"] / window_days ) record["三日加权平均T0裂变率"] = record["T0裂变率"] record["三日加权平均实际ROI"] = record["实际ROI"] record["三日加权平均效率ROI"] = record["ROI"] record["首日首层UV"] = float(by_date.iloc[0]["首层UV"]) record["首日效率ROI"] = ( float(daily_predicted_roi[0]) if np.isfinite(daily_predicted_roi[0]) else np.nan ) record["最新日效率ROI"] = ( float(daily_predicted_roi[-1]) if np.isfinite(daily_predicted_roi[-1]) else np.nan ) for index, dt in enumerate(expected_dates): actual_roi = ( float(daily_actual_roi[index]) if np.isfinite(daily_actual_roi[index]) else np.nan ) predicted_roi = ( float(daily_predicted_roi[index]) if np.isfinite(daily_predicted_roi[index]) else np.nan ) record[f"首层UV_{dt}"] = float(by_date.iloc[index]["首层UV"]) record[f"T0裂变数_{dt}"] = float(by_date.iloc[index]["T0裂变数"]) record[f"成本_{dt}"] = float(by_date.iloc[index]["成本"]) record[f"效率收入_{dt}"] = float(by_date.iloc[index]["效率收入"]) record[f"裂变效率收入_{dt}"] = float( by_date.iloc[index]["裂变效率收入"] ) record[f"预测全链路效率收入_{dt}"] = float( by_date.iloc[index]["预测全链路效率收入"] ) record[f"实际ROI_{dt}"] = actual_roi record[f"ROI_{dt}"] = predicted_roi if entity_type in (ENTITY_SELF, ENTITY_SELF_AD): record["广告age"] = int(ages.get(str(record.get("广告id", "")), 0)) records.append(record) return pd.DataFrame(records) def compute_roi_summary( raw_daily: pd.DataFrame, expected_dates: Iterable[str], ad_age: pd.DataFrame | None = None, *, fission_parameters: FissionMultiplierParameters, ) -> tuple[pd.DataFrame, list[str]]: """Return three-day entity ROI snapshots and normalized dates.""" dates = _expected_three_days(expected_dates) daily = prepare_daily_metrics(raw_daily, fission_parameters) daily = daily[daily["dt"].isin(dates)].copy() ages = _age_map(ad_age) summaries = [] for entity_type in (ENTITY_SELF, ENTITY_SELF_AD, ENTITY_GZH): entity_daily = daily[daily["entity_type"].eq(entity_type)] if entity_daily.empty: continue entity_summary = _summarize_entity( entity_daily, entity_type, dates, ages, ) if not entity_summary.empty: summaries.append(entity_summary) if not summaries: raise ValueError("三个目标实体均没有可用的最近三日数据") summary = pd.concat(summaries, ignore_index=True, sort=False) return summary, dates