run_impact_report.py 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160
  1. #!/usr/bin/env python3
  2. """Run the DAU percentile, fission-impact, and commercial-impact report."""
  3. from __future__ import annotations
  4. import argparse
  5. import json
  6. import sys
  7. from datetime import datetime, timedelta
  8. from pathlib import Path
  9. from string import Template
  10. import pandas as pd
  11. PROJECT_DIR = Path(__file__).resolve().parents[1]
  12. WORKSPACE_ROOT = Path(__file__).resolve().parents[2]
  13. TEMPLATE_PATH = PROJECT_DIR / "sql" / "dau_commercial_impact_template.sql"
  14. DISPLAY_COLUMNS = {
  15. "dt": "dt",
  16. "安全分分位点": "安全分分位点",
  17. "score": "Score",
  18. "p20_score": "P20得分",
  19. "p40_score": "P40得分",
  20. "p60_score": "P60得分",
  21. "p80_score": "P80得分",
  22. "dau(当天日活总数)": "DAU(当天日活总数)",
  23. "访问uv(每个分位点对应人数)": "访问UV(每个分位点对应人数)",
  24. "访问uv占比": "访问UV占比",
  25. "当日总分享次数": "当日总分享次数",
  26. "分享次数(每个分位点对应次数)": "分享次数(每个分位点对应次数)",
  27. "当日总回流": "当日总回流",
  28. "分享次数占比": "分享次数占比",
  29. "当日分享当日回流总人数": "当日分享当日回流总人数",
  30. "当日分享当日回流(各分位点人数)": "当日分享当日回流(各分位点人数)",
  31. "当日分享当日回流人数占比": "当日分享当日回流人数占比",
  32. "总广告无曝光人数": "总广告无曝光人数",
  33. "无广告曝光人数": "无广告曝光人数",
  34. "无广告曝光人数占比": "无广告曝光人数占比",
  35. "总点击人数": "总点击人数",
  36. "点击人数": "点击人数",
  37. "点击人数占比": "点击人数占比",
  38. "总已转化人数": "总已转化人数",
  39. "已转化人数": "已转化人数",
  40. "已转化人数占比": "已转化人数占比",
  41. }
  42. PERCENT_COLUMNS = [
  43. "访问UV占比",
  44. "分享次数占比",
  45. "当日分享当日回流人数占比",
  46. "无广告曝光人数占比",
  47. "点击人数占比",
  48. "已转化人数占比",
  49. ]
  50. def render_sql(bizdate: str, lookback_days: int, active_weight: str) -> str:
  51. anchor = datetime.strptime(bizdate, "%Y%m%d")
  52. history_start = (anchor - timedelta(days=lookback_days)).strftime("%Y%m%d")
  53. history_end = (anchor - timedelta(days=1)).strftime("%Y%m%d")
  54. return Template(TEMPLATE_PATH.read_text(encoding="utf-8")).substitute(
  55. bizdate=bizdate,
  56. history_start_dt=history_start,
  57. history_end_dt=history_end,
  58. active_weight=active_weight,
  59. )
  60. def format_and_validate(result: pd.DataFrame) -> pd.DataFrame:
  61. result.columns = [str(column).lower() for column in result.columns]
  62. missing = sorted(set(DISPLAY_COLUMNS) - set(result.columns))
  63. if missing:
  64. raise ValueError("结果缺少字段: " + ", ".join(missing))
  65. result = result[list(DISPLAY_COLUMNS)].rename(columns=DISPLAY_COLUMNS)
  66. if len(result) != 5:
  67. raise AssertionError(f"结果应为5档,实际为{len(result)}档")
  68. total_dau = int(result["DAU(当天日活总数)"].iloc[0])
  69. if int(result["访问UV(每个分位点对应人数)"].sum()) != total_dau:
  70. raise AssertionError("五档访问UV之和不等于DAU")
  71. additive_metrics = [
  72. ("分享次数(每个分位点对应次数)", "当日总分享次数"),
  73. ("无广告曝光人数", "总广告无曝光人数"),
  74. ("点击人数", "总点击人数"),
  75. ("已转化人数", "总已转化人数"),
  76. ]
  77. for band_column, total_column in additive_metrics:
  78. if int(result[band_column].sum()) != int(result[total_column].iloc[0]):
  79. raise AssertionError(f"五档{band_column}之和不等于{total_column}")
  80. thresholds = result[["P20得分", "P40得分", "P60得分", "P80得分"]].iloc[0]
  81. if not thresholds.is_monotonic_increasing:
  82. raise AssertionError("P20/P40/P60/P80得分不是单调递增")
  83. display = result.copy()
  84. for column in PERCENT_COLUMNS:
  85. display[column] = pd.to_numeric(display[column], errors="coerce").map(
  86. lambda value: "" if pd.isna(value) else f"{value:.2%}"
  87. )
  88. return display
  89. def main() -> None:
  90. parser = argparse.ArgumentParser(description=__doc__)
  91. parser.add_argument("--bizdate", required=True, help="统计日,格式yyyyMMdd")
  92. parser.add_argument("--lookback-days", type=int, default=180)
  93. parser.add_argument("--active-weight", choices=["0.5", "1.0"], default="0.5")
  94. parser.add_argument("--run-name", required=True)
  95. parser.add_argument("--sql-only", action="store_true")
  96. parser.add_argument("--format-only", action="store_true")
  97. args = parser.parse_args()
  98. if args.lookback_days < 1:
  99. raise ValueError("lookback-days必须大于0")
  100. output_dir = PROJECT_DIR / "output" / args.run_name / "impact"
  101. sql_dir = PROJECT_DIR / "sql" / args.run_name / "impact"
  102. output_dir.mkdir(parents=True, exist_ok=True)
  103. sql_dir.mkdir(parents=True, exist_ok=True)
  104. sql = render_sql(args.bizdate, args.lookback_days, args.active_weight)
  105. rendered_sql_path = sql_dir / "query.sql"
  106. rendered_sql_path.write_text(sql + "\n", encoding="utf-8")
  107. raw_path = output_dir / "result_raw.csv"
  108. output_path = output_dir / "dau_commercial_impact.csv"
  109. if args.sql_only:
  110. print(f"rendered_sql={rendered_sql_path}")
  111. return
  112. if args.format_only:
  113. if not raw_path.exists():
  114. raise ValueError(f"缺少原始结果: {raw_path}")
  115. display = format_and_validate(pd.read_csv(raw_path))
  116. display.to_csv(output_path, index=False)
  117. print(f"output={output_path}")
  118. return
  119. if str(WORKSPACE_ROOT) not in sys.path:
  120. sys.path.insert(0, str(WORKSPACE_ROOT))
  121. from odps_module import ODPSClient
  122. client = ODPSClient(project="loghubods")
  123. instance = client.odps.run_sql(sql, hints={"odps.sql.allow.cartesian": "true"})
  124. run = {"instance_id": instance.id, "logview": instance.get_logview_address()}
  125. (output_dir / "odps_runs.json").write_text(
  126. json.dumps([run], ensure_ascii=False, indent=2) + "\n",
  127. encoding="utf-8",
  128. )
  129. print(f"instance={run['instance_id']}", flush=True)
  130. print(f"logview={run['logview']}", flush=True)
  131. instance.wait_for_success()
  132. with instance.open_reader(tunnel=True) as reader:
  133. result = reader.to_pandas()
  134. result.to_csv(raw_path, index=False)
  135. display = format_and_validate(result)
  136. display.to_csv(output_path, index=False)
  137. print(f"output={output_path}")
  138. if __name__ == "__main__":
  139. main()