#!/usr/bin/env python3 """Run the DAU percentile, fission-impact, and commercial-impact report.""" from __future__ import annotations import argparse import json import sys from datetime import datetime, timedelta from pathlib import Path from string import Template import pandas as pd PROJECT_DIR = Path(__file__).resolve().parents[1] WORKSPACE_ROOT = Path(__file__).resolve().parents[2] TEMPLATE_PATH = PROJECT_DIR / "sql" / "dau_commercial_impact_template.sql" DISPLAY_COLUMNS = { "dt": "dt", "安全分分位点": "安全分分位点", "score": "Score", "p20_score": "P20得分", "p40_score": "P40得分", "p60_score": "P60得分", "p80_score": "P80得分", "dau(当天日活总数)": "DAU(当天日活总数)", "访问uv(每个分位点对应人数)": "访问UV(每个分位点对应人数)", "访问uv占比": "访问UV占比", "当日总分享次数": "当日总分享次数", "分享次数(每个分位点对应次数)": "分享次数(每个分位点对应次数)", "当日总回流": "当日总回流", "分享次数占比": "分享次数占比", "当日分享当日回流总人数": "当日分享当日回流总人数", "当日分享当日回流(各分位点人数)": "当日分享当日回流(各分位点人数)", "当日分享当日回流人数占比": "当日分享当日回流人数占比", "总广告无曝光人数": "总广告无曝光人数", "无广告曝光人数": "无广告曝光人数", "无广告曝光人数占比": "无广告曝光人数占比", "总点击人数": "总点击人数", "点击人数": "点击人数", "点击人数占比": "点击人数占比", "总已转化人数": "总已转化人数", "已转化人数": "已转化人数", "已转化人数占比": "已转化人数占比", } PERCENT_COLUMNS = [ "访问UV占比", "分享次数占比", "当日分享当日回流人数占比", "无广告曝光人数占比", "点击人数占比", "已转化人数占比", ] def render_sql(bizdate: str, lookback_days: int, active_weight: str) -> str: anchor = datetime.strptime(bizdate, "%Y%m%d") history_start = (anchor - timedelta(days=lookback_days)).strftime("%Y%m%d") history_end = (anchor - timedelta(days=1)).strftime("%Y%m%d") return Template(TEMPLATE_PATH.read_text(encoding="utf-8")).substitute( bizdate=bizdate, history_start_dt=history_start, history_end_dt=history_end, active_weight=active_weight, ) def format_and_validate(result: pd.DataFrame) -> pd.DataFrame: result.columns = [str(column).lower() for column in result.columns] missing = sorted(set(DISPLAY_COLUMNS) - set(result.columns)) if missing: raise ValueError("结果缺少字段: " + ", ".join(missing)) result = result[list(DISPLAY_COLUMNS)].rename(columns=DISPLAY_COLUMNS) if len(result) != 5: raise AssertionError(f"结果应为5档,实际为{len(result)}档") total_dau = int(result["DAU(当天日活总数)"].iloc[0]) if int(result["访问UV(每个分位点对应人数)"].sum()) != total_dau: raise AssertionError("五档访问UV之和不等于DAU") additive_metrics = [ ("分享次数(每个分位点对应次数)", "当日总分享次数"), ("无广告曝光人数", "总广告无曝光人数"), ("点击人数", "总点击人数"), ("已转化人数", "总已转化人数"), ] for band_column, total_column in additive_metrics: if int(result[band_column].sum()) != int(result[total_column].iloc[0]): raise AssertionError(f"五档{band_column}之和不等于{total_column}") thresholds = result[["P20得分", "P40得分", "P60得分", "P80得分"]].iloc[0] if not thresholds.is_monotonic_increasing: raise AssertionError("P20/P40/P60/P80得分不是单调递增") display = result.copy() for column in PERCENT_COLUMNS: display[column] = pd.to_numeric(display[column], errors="coerce").map( lambda value: "" if pd.isna(value) else f"{value:.2%}" ) return display def main() -> None: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--bizdate", required=True, help="统计日,格式yyyyMMdd") parser.add_argument("--lookback-days", type=int, default=180) parser.add_argument("--active-weight", choices=["0.5", "1.0"], default="0.5") parser.add_argument("--run-name", required=True) parser.add_argument("--sql-only", action="store_true") parser.add_argument("--format-only", action="store_true") args = parser.parse_args() if args.lookback_days < 1: raise ValueError("lookback-days必须大于0") output_dir = PROJECT_DIR / "output" / args.run_name / "impact" sql_dir = PROJECT_DIR / "sql" / args.run_name / "impact" output_dir.mkdir(parents=True, exist_ok=True) sql_dir.mkdir(parents=True, exist_ok=True) sql = render_sql(args.bizdate, args.lookback_days, args.active_weight) rendered_sql_path = sql_dir / "query.sql" rendered_sql_path.write_text(sql + "\n", encoding="utf-8") raw_path = output_dir / "result_raw.csv" output_path = output_dir / "dau_commercial_impact.csv" if args.sql_only: print(f"rendered_sql={rendered_sql_path}") return if args.format_only: if not raw_path.exists(): raise ValueError(f"缺少原始结果: {raw_path}") display = format_and_validate(pd.read_csv(raw_path)) display.to_csv(output_path, index=False) print(f"output={output_path}") return if str(WORKSPACE_ROOT) not in sys.path: sys.path.insert(0, str(WORKSPACE_ROOT)) from odps_module import ODPSClient client = ODPSClient(project="loghubods") instance = client.odps.run_sql(sql, hints={"odps.sql.allow.cartesian": "true"}) run = {"instance_id": instance.id, "logview": instance.get_logview_address()} (output_dir / "odps_runs.json").write_text( json.dumps([run], ensure_ascii=False, indent=2) + "\n", encoding="utf-8", ) print(f"instance={run['instance_id']}", flush=True) print(f"logview={run['logview']}", flush=True) instance.wait_for_success() with instance.open_reader(tunnel=True) as reader: result = reader.to_pandas() result.to_csv(raw_path, index=False) display = format_and_validate(result) display.to_csv(output_path, index=False) print(f"output={output_path}") if __name__ == "__main__": main()