run_bucket_source_offline_5day_report.py 4.2 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091
  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. group_name,
  14. reorder_columns,
  15. )
  16. BASE = Path(__file__).resolve().parent
  17. SQL_FILE = BASE / "bucket_source_pv_uv_return_offline_20260716_20260720_apptype4_all_versions.sql"
  18. DAU_SQL_FILE = BASE / "offline_dau_by_bucket_20260716_20260720_apptype4.sql"
  19. OUTPUT_FILE = BASE / "bucket_source_full_report_offline_20260716_20260720_apptype4_all_versions.csv"
  20. def build_daily_report(data, stat_date):
  21. daily = data[data["日期"] == stat_date].copy()
  22. daily["尾号"] = daily["尾号"].astype(str)
  23. daily["行类型"] = "尾号明细"
  24. daily["分组"] = daily["尾号"].map(group_name)
  25. detail = daily[["日期", "产品类型", "版本号", "行类型", "分组", "尾号", *FACT_COLUMNS]]
  26. aggregate_rows, mean_rows = [], []
  27. for group, buckets in [("实验组(0、1)", {"0", "1"}), ("对照组(其余14桶)", set("23456789abcdef"))]:
  28. selected = detail[detail["尾号"].isin(buckets)]
  29. common = {"日期": stat_date, "产品类型": "4", "版本号": "全部", "分组": group}
  30. aggregate_rows.append({
  31. **common,
  32. "行类型": "分组聚合",
  33. "尾号": "0、1" if len(buckets) == 2 else "其余14桶",
  34. **selected[FACT_COLUMNS].sum().to_dict(),
  35. })
  36. mean_rows.append({
  37. **common,
  38. "行类型": "每桶均值",
  39. "尾号": "2桶均值" if len(buckets) == 2 else "14桶均值",
  40. **selected[FACT_COLUMNS].mean().to_dict(),
  41. })
  42. result = pd.concat([detail, pd.DataFrame(aggregate_rows), pd.DataFrame(mean_rows)], ignore_index=True)
  43. result = add_relative_changes(add_rates(result))
  44. return reorder_columns(format_fact_columns_as_integers(format_rate_columns(result)))
  45. def main():
  46. if "--rebuild-with-dau" in sys.argv:
  47. existing = pd.read_csv(OUTPUT_FILE, dtype={"日期": str, "尾号": str})
  48. detail = existing[existing["行类型"] == "尾号明细"][[
  49. "日期", "产品类型", "版本号", "尾号", *FACT_COLUMNS
  50. ]].copy()
  51. detail = detail.drop(columns=["DAU"])
  52. dau = ODPSClient().execute_sql(DAU_SQL_FILE.read_text(encoding="utf-8")).rename(
  53. columns={"stat_date": "日期", "bucket": "尾号", "dau": "DAU"}
  54. )
  55. dau["日期"] = dau["日期"].astype(str)
  56. dau["尾号"] = dau["尾号"].astype(str)
  57. facts = detail.merge(dau, on=["日期", "尾号"], how="left", validate="one_to_one")
  58. if len(facts) != 80 or facts["DAU"].isna().any():
  59. raise ValueError(f"DAU合并异常:rows={len(facts)}, missing_dau={facts['DAU'].isna().sum()}")
  60. reports = [build_daily_report(facts, stat_date) for stat_date in sorted(facts["日期"].unique(), reverse=True)]
  61. result = pd.concat(reports, ignore_index=True)
  62. result.to_csv(OUTPUT_FILE, index=False, encoding="utf-8-sig")
  63. print(f"[CSV REBUILT WITH OFFLINE DAU] {OUTPUT_FILE}", flush=True)
  64. print(f"[ROWS] {len(result)}", flush=True)
  65. print(result[result["行类型"] == "分组聚合"][["日期", "分组", "DAU", "DAU相对对照组变化率", "全部曝光PV/DAU", "全部分享PV/DAU", "全部回流UV/DAU", "全部ROV(回流UV/曝光PV)"]].to_string(index=False), flush=True)
  66. return
  67. facts = ODPSClient().execute_sql(SQL_FILE.read_text(encoding="utf-8")).rename(columns=column_mapping())
  68. facts["日期"] = facts["日期"].astype(str)
  69. reports = [build_daily_report(facts, stat_date) for stat_date in sorted(facts["日期"].unique(), reverse=True)]
  70. result = pd.concat(reports, ignore_index=True)
  71. result.to_csv(OUTPUT_FILE, index=False, encoding="utf-8-sig")
  72. print(f"[CSV] {OUTPUT_FILE}", flush=True)
  73. print(f"[ROWS] {len(result)}", flush=True)
  74. print(result[result["行类型"] == "分组聚合"][["日期", "分组", "DAU", "DAU相对对照组变化率", "全部曝光PV/DAU", "全部分享PV/DAU", "全部ROV(回流UV/曝光PV)"]].to_string(index=False), flush=True)
  75. if __name__ == "__main__":
  76. main()