| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091 |
- #!/usr/bin/env python3
- from pathlib import Path
- import sys
- import pandas as pd
- from odps_module import ODPSClient
- from run_bucket_source_pv_uv_return_realtime import (
- FACT_COLUMNS,
- add_rates,
- add_relative_changes,
- column_mapping,
- format_fact_columns_as_integers,
- format_rate_columns,
- group_name,
- reorder_columns,
- )
- BASE = Path(__file__).resolve().parent
- SQL_FILE = BASE / "bucket_source_pv_uv_return_offline_20260716_20260720_apptype4_all_versions.sql"
- DAU_SQL_FILE = BASE / "offline_dau_by_bucket_20260716_20260720_apptype4.sql"
- OUTPUT_FILE = BASE / "bucket_source_full_report_offline_20260716_20260720_apptype4_all_versions.csv"
- def build_daily_report(data, stat_date):
- daily = data[data["日期"] == stat_date].copy()
- daily["尾号"] = daily["尾号"].astype(str)
- daily["行类型"] = "尾号明细"
- daily["分组"] = daily["尾号"].map(group_name)
- detail = daily[["日期", "产品类型", "版本号", "行类型", "分组", "尾号", *FACT_COLUMNS]]
- aggregate_rows, mean_rows = [], []
- for group, buckets in [("实验组(0、1)", {"0", "1"}), ("对照组(其余14桶)", set("23456789abcdef"))]:
- selected = detail[detail["尾号"].isin(buckets)]
- common = {"日期": stat_date, "产品类型": "4", "版本号": "全部", "分组": group}
- aggregate_rows.append({
- **common,
- "行类型": "分组聚合",
- "尾号": "0、1" if len(buckets) == 2 else "其余14桶",
- **selected[FACT_COLUMNS].sum().to_dict(),
- })
- mean_rows.append({
- **common,
- "行类型": "每桶均值",
- "尾号": "2桶均值" if len(buckets) == 2 else "14桶均值",
- **selected[FACT_COLUMNS].mean().to_dict(),
- })
- result = pd.concat([detail, pd.DataFrame(aggregate_rows), pd.DataFrame(mean_rows)], ignore_index=True)
- result = add_relative_changes(add_rates(result))
- return reorder_columns(format_fact_columns_as_integers(format_rate_columns(result)))
- def main():
- if "--rebuild-with-dau" in sys.argv:
- existing = pd.read_csv(OUTPUT_FILE, dtype={"日期": str, "尾号": str})
- detail = existing[existing["行类型"] == "尾号明细"][[
- "日期", "产品类型", "版本号", "尾号", *FACT_COLUMNS
- ]].copy()
- detail = detail.drop(columns=["DAU"])
- dau = ODPSClient().execute_sql(DAU_SQL_FILE.read_text(encoding="utf-8")).rename(
- columns={"stat_date": "日期", "bucket": "尾号", "dau": "DAU"}
- )
- dau["日期"] = dau["日期"].astype(str)
- dau["尾号"] = dau["尾号"].astype(str)
- facts = detail.merge(dau, on=["日期", "尾号"], how="left", validate="one_to_one")
- if len(facts) != 80 or facts["DAU"].isna().any():
- raise ValueError(f"DAU合并异常:rows={len(facts)}, missing_dau={facts['DAU'].isna().sum()}")
- reports = [build_daily_report(facts, stat_date) for stat_date in sorted(facts["日期"].unique(), reverse=True)]
- result = pd.concat(reports, ignore_index=True)
- result.to_csv(OUTPUT_FILE, index=False, encoding="utf-8-sig")
- print(f"[CSV REBUILT WITH OFFLINE DAU] {OUTPUT_FILE}", flush=True)
- print(f"[ROWS] {len(result)}", flush=True)
- print(result[result["行类型"] == "分组聚合"][["日期", "分组", "DAU", "DAU相对对照组变化率", "全部曝光PV/DAU", "全部分享PV/DAU", "全部回流UV/DAU", "全部ROV(回流UV/曝光PV)"]].to_string(index=False), flush=True)
- return
- facts = ODPSClient().execute_sql(SQL_FILE.read_text(encoding="utf-8")).rename(columns=column_mapping())
- facts["日期"] = facts["日期"].astype(str)
- reports = [build_daily_report(facts, stat_date) for stat_date in sorted(facts["日期"].unique(), reverse=True)]
- result = pd.concat(reports, ignore_index=True)
- result.to_csv(OUTPUT_FILE, index=False, encoding="utf-8-sig")
- print(f"[CSV] {OUTPUT_FILE}", flush=True)
- print(f"[ROWS] {len(result)}", flush=True)
- print(result[result["行类型"] == "分组聚合"][["日期", "分组", "DAU", "DAU相对对照组变化率", "全部曝光PV/DAU", "全部分享PV/DAU", "全部ROV(回流UV/曝光PV)"]].to_string(index=False), flush=True)
- if __name__ == "__main__":
- main()
|