run_bucket_source_offline_5day_report.py 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197
  1. #!/usr/bin/env python3
  2. from pathlib import Path
  3. import sys
  4. import pandas as pd
  5. from odps_module import ODPSClient
  6. from run_bucket_source_pv_uv_return_realtime import (
  7. FACT_COLUMNS,
  8. add_rates,
  9. add_relative_changes,
  10. column_mapping,
  11. format_fact_columns_as_integers,
  12. format_rate_columns,
  13. reorder_columns,
  14. )
  15. BASE = Path(__file__).resolve().parent
  16. CONFIGS = {
  17. "offline": {
  18. "sql": BASE / "bucket_source_pv_uv_return_offline_20260716_20260720_apptype4_all_versions.sql",
  19. "output": BASE / "bucket_source_full_report_offline_20260716_20260720_apptype4_all_versions.csv",
  20. "experiment_buckets": set("01"),
  21. "experiment_label": "实验组(0、1)",
  22. "control_label": "对照组(其余14桶)",
  23. "experiment_tail_label": "0、1",
  24. "control_tail_label": "其余14桶",
  25. },
  26. "offline_20260713_20260716": {
  27. "sql": BASE / "bucket_source_pv_uv_return_offline_20260713_20260716_apptype4_all_versions.sql",
  28. "output": BASE / "bucket_source_full_report_offline_20260713_20260716_apptype4_all_versions.csv",
  29. "experiment_buckets": set("01"),
  30. "experiment_label": "实验组(0、1)",
  31. "control_label": "对照组(其余14桶)",
  32. "experiment_tail_label": "0、1",
  33. "control_tail_label": "其余14桶",
  34. },
  35. "offline_20260612_1508": {
  36. "sql": BASE / "bucket_source_pv_uv_return_offline_20260612_apptype0_version1508.sql",
  37. "output": BASE / "bucket_source_full_report_offline_20260612_apptype0_version1508.csv",
  38. "experiment_buckets": set("0123456789ab"),
  39. "experiment_label": "实验组(0-b)",
  40. "control_label": "对照组(c-f)",
  41. "experiment_tail_label": "0-b",
  42. "control_tail_label": "c-f",
  43. },
  44. "offline_20260617_all": {
  45. "sql": BASE / "bucket_source_pv_uv_return_offline_20260617_apptype0_all_versions.sql",
  46. "output": BASE / "bucket_source_full_report_offline_20260617_apptype0_all_versions.csv",
  47. "experiment_buckets": set("01"),
  48. "experiment_label": "实验组(0、1)",
  49. "control_label": "对照组(其余14桶)",
  50. "experiment_tail_label": "0、1",
  51. "control_tail_label": "其余14桶",
  52. },
  53. "offline_20260617_1512": {
  54. "sql": BASE / "bucket_source_pv_uv_return_offline_20260617_apptype0_version1512.sql",
  55. "output": BASE / "bucket_source_full_report_offline_20260617_apptype0_version1512.csv",
  56. "experiment_buckets": set("01"),
  57. "experiment_label": "实验组(0、1)",
  58. "control_label": "对照组(其余14桶)",
  59. "experiment_tail_label": "0、1",
  60. "control_tail_label": "其余14桶",
  61. },
  62. "offline_20260721_1578": {
  63. "sql": BASE / "bucket_source_pv_uv_return_offline_20260721_apptype4_version1578.sql",
  64. "output": BASE / "bucket_source_full_report_offline_20260721_apptype4_version1578.csv",
  65. "experiment_buckets": set("01"),
  66. "experiment_label": "实验组(0、1)",
  67. "control_label": "对照组(其余14桶)",
  68. "experiment_tail_label": "0、1",
  69. "control_tail_label": "其余14桶",
  70. },
  71. "offline_20260722_all": {
  72. "sql": BASE / "bucket_source_pv_uv_return_offline_20260722_apptype4_all_versions.sql",
  73. "output": BASE / "bucket_source_full_report_offline_20260722_apptype4_all_versions.csv",
  74. "experiment_buckets": set("01"),
  75. "experiment_label": "实验组(0、1)",
  76. "control_label": "对照组(其余14桶)",
  77. "experiment_tail_label": "0、1",
  78. "control_tail_label": "其余14桶",
  79. },
  80. "offline_20260717_20260721_exclude_qywx": {
  81. "sql": BASE / "bucket_source_pv_uv_return_offline_20260717_20260721_apptype4_all_versions_exclude_qywx.sql",
  82. "output": BASE / "bucket_source_full_report_offline_20260717_20260721_apptype4_all_versions_exclude_qywx.csv",
  83. "experiment_buckets": set("01"),
  84. "experiment_label": "实验组(0、1)",
  85. "control_label": "对照组(其余14桶)",
  86. "experiment_tail_label": "0、1",
  87. "control_tail_label": "其余14桶",
  88. },
  89. "offline_20260721_1578_exclude_qywx": {
  90. "sql": BASE / "bucket_source_pv_uv_return_offline_20260721_apptype4_version1578_exclude_qywx.sql",
  91. "output": BASE / "bucket_source_full_report_offline_20260721_apptype4_version1578_exclude_qywx.csv",
  92. "experiment_buckets": set("01"),
  93. "experiment_label": "实验组(0、1)",
  94. "control_label": "对照组(其余14桶)",
  95. "experiment_tail_label": "0、1",
  96. "control_tail_label": "其余14桶",
  97. },
  98. }
  99. MODE = next((arg for arg in sys.argv[1:] if not arg.startswith("--")), "offline")
  100. if MODE not in CONFIGS:
  101. raise SystemExit("usage: python3 run_bucket_source_offline_5day_report.py [offline|offline_20260713_20260716|offline_20260612_1508|offline_20260617_all|offline_20260617_1512|offline_20260721_1578|offline_20260722_all|offline_20260717_20260721_exclude_qywx|offline_20260721_1578_exclude_qywx] [--rebuild-with-dau]")
  102. CONFIG = CONFIGS[MODE]
  103. SQL_FILE, OUTPUT_FILE = CONFIG["sql"], CONFIG["output"]
  104. DAU_SQL_FILE = BASE / "offline_dau_by_bucket_20260716_20260720_apptype4.sql"
  105. def build_daily_report(data, stat_date):
  106. daily = data[data["日期"] == stat_date].copy()
  107. app_types = daily["产品类型"].astype(str).unique()
  108. versions = daily["版本号"].astype(str).unique()
  109. if len(app_types) != 1 or len(versions) != 1:
  110. raise ValueError(f"日期{stat_date}的产品或版本不唯一:apptype={app_types}, versions={versions}")
  111. app_type, version = app_types[0], versions[0]
  112. daily["尾号"] = daily["尾号"].astype(str)
  113. daily["行类型"] = "尾号明细"
  114. experiment_buckets = CONFIG["experiment_buckets"]
  115. control_buckets = set("0123456789abcdef") - experiment_buckets
  116. experiment_label = CONFIG["experiment_label"]
  117. control_label = CONFIG["control_label"]
  118. tail_labels = {
  119. experiment_label: CONFIG["experiment_tail_label"],
  120. control_label: CONFIG["control_tail_label"],
  121. }
  122. daily["分组"] = daily["尾号"].map(lambda bucket: experiment_label if bucket in experiment_buckets else control_label)
  123. detail = daily[["日期", "产品类型", "版本号", "行类型", "分组", "尾号", *FACT_COLUMNS]]
  124. aggregate_rows, mean_rows = [], []
  125. for group, buckets in [(experiment_label, experiment_buckets), (control_label, control_buckets)]:
  126. selected = detail[detail["尾号"].isin(buckets)]
  127. common = {"日期": stat_date, "产品类型": app_type, "版本号": version, "分组": group}
  128. aggregate_rows.append({
  129. **common,
  130. "行类型": "分组聚合",
  131. "尾号": tail_labels[group],
  132. **selected[FACT_COLUMNS].sum().to_dict(),
  133. })
  134. mean_rows.append({
  135. **common,
  136. "行类型": "每桶均值",
  137. "尾号": f"{len(buckets)}桶均值",
  138. **selected[FACT_COLUMNS].mean().to_dict(),
  139. })
  140. result = pd.concat([detail, pd.DataFrame(aggregate_rows), pd.DataFrame(mean_rows)], ignore_index=True)
  141. result = add_relative_changes(
  142. add_rates(result),
  143. experiment_label=experiment_label,
  144. control_label=control_label,
  145. experiment_bucket_count=len(experiment_buckets),
  146. control_bucket_count=len(control_buckets),
  147. )
  148. return reorder_columns(format_fact_columns_as_integers(format_rate_columns(result)))
  149. def main():
  150. if "--rebuild-with-dau" in sys.argv:
  151. if MODE != "offline":
  152. raise ValueError("--rebuild-with-dau 仅适用于五天全版本离线报表")
  153. existing = pd.read_csv(OUTPUT_FILE, dtype={"日期": str, "尾号": str})
  154. detail = existing[existing["行类型"] == "尾号明细"][[
  155. "日期", "产品类型", "版本号", "尾号", *FACT_COLUMNS
  156. ]].copy()
  157. detail = detail.drop(columns=["DAU"])
  158. dau = ODPSClient().execute_sql(DAU_SQL_FILE.read_text(encoding="utf-8")).rename(
  159. columns={"stat_date": "日期", "bucket": "尾号", "dau": "DAU"}
  160. )
  161. dau["日期"] = dau["日期"].astype(str)
  162. dau["尾号"] = dau["尾号"].astype(str)
  163. facts = detail.merge(dau, on=["日期", "尾号"], how="left", validate="one_to_one")
  164. if len(facts) != 80 or facts["DAU"].isna().any():
  165. raise ValueError(f"DAU合并异常:rows={len(facts)}, missing_dau={facts['DAU'].isna().sum()}")
  166. reports = [build_daily_report(facts, stat_date) for stat_date in sorted(facts["日期"].unique(), reverse=True)]
  167. result = pd.concat(reports, ignore_index=True)
  168. result.to_csv(OUTPUT_FILE, index=False, encoding="utf-8-sig")
  169. print(f"[CSV REBUILT WITH OFFLINE DAU] {OUTPUT_FILE}", flush=True)
  170. print(f"[ROWS] {len(result)}", flush=True)
  171. print(result[result["行类型"] == "分组聚合"][["日期", "分组", "DAU", "DAU相对对照组变化率", "全部曝光PV/DAU", "全部分享PV/DAU", "全部回流UV/DAU", "全部ROV(回流UV/曝光PV)"]].to_string(index=False), flush=True)
  172. return
  173. facts = ODPSClient().execute_sql(SQL_FILE.read_text(encoding="utf-8")).rename(columns=column_mapping())
  174. facts["日期"] = facts["日期"].astype(str)
  175. reports = [build_daily_report(facts, stat_date) for stat_date in sorted(facts["日期"].unique(), reverse=True)]
  176. result = pd.concat(reports, ignore_index=True)
  177. result.to_csv(OUTPUT_FILE, index=False, encoding="utf-8-sig")
  178. print(f"[CSV] {OUTPUT_FILE}", flush=True)
  179. print(f"[ROWS] {len(result)}", flush=True)
  180. print(result[result["行类型"] == "分组聚合"][["日期", "分组", "DAU", "DAU相对对照组变化率", "全部曝光PV/DAU", "全部分享PV/DAU", "全部ROV(回流UV/曝光PV)"]].to_string(index=False), flush=True)
  181. if __name__ == "__main__":
  182. main()