| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278 |
- #!/usr/bin/env python3
- """Execute V2 feature SQL modules and produce a Feishu-shaped user detail CSV."""
- from __future__ import annotations
- import argparse
- import json
- import sys
- from pathlib import Path
- import pandas as pd
- PROJECT_DIR = Path(__file__).resolve().parents[1]
- WORKSPACE_ROOT = Path(__file__).resolve().parents[2]
- if str(WORKSPACE_ROOT) not in sys.path:
- sys.path.insert(0, str(WORKSPACE_ROOT))
- if str(PROJECT_DIR) not in sys.path:
- sys.path.insert(0, str(PROJECT_DIR))
- from odps_module import ODPSClient # noqa: E402
- from src.feature_sql import ( # noqa: E402
- build_attributes_sql,
- build_background_return_sql,
- build_capture_sql,
- build_core_sql,
- build_input_users_sql,
- build_landing_return_sql,
- build_return_sql,
- )
- from src.user_input import load_users as load_user_input # noqa: E402
- BASE_DIR = PROJECT_DIR
- DEFAULT_INPUT = BASE_DIR / "samples" / "v2_220_20260810.csv"
- OUTPUT_COLUMNS = [
- "用户分类", "子分类", "统计周期", "用户id", "操作系统", "机型", "机型数",
- "地域", "地域数", "活跃天数", "点击卡片去重次数", "点击卡片不去重次数",
- "有效播放次数(进度30%或播放时长大于20秒)", "分享次数",
- "过去180天分享带回去重回流人数",
- "点击广告次数", "长按扫码次数", "广告播放中截图次数",
- "广告播放中切后台再返回次数", "广告落地页截图次数",
- "无扫码情况下落地页切换后台再返回", "来源个人分享(场景值1007)",
- "来源群分享(场景值1008)", "来源公众号文章(场景值1058)",
- "来源公众号即转(场景值1074)", "来源小程序投流(1067+1095)",
- "来源其他场景值", "来源上游用户shareid", "来源群id",
- ]
- RENAME = {
- "user_type": "用户分类",
- "sub_category": "子分类",
- "mid": "用户id",
- "operating_system_list": "操作系统",
- "device_model_list": "机型",
- "device_model_cnt": "机型数",
- "city_list": "地域",
- "city_cnt": "地域数",
- "active_day_cnt": "活跃天数",
- "card_click_object_uv": "点击卡片去重次数",
- "card_click_object_pv": "点击卡片不去重次数",
- "real_play_cnt": "有效播放次数(进度30%或播放时长大于20秒)",
- "video_share_cnt": "分享次数",
- "return_people_cnt": "过去180天分享带回去重回流人数",
- "own_ad_click_cnt": "点击广告次数",
- "long_press_scan_cnt": "长按扫码次数",
- "ad_play_capture_cnt": "广告播放中截图次数",
- "ad_play_background_return_cnt": "广告播放中切后台再返回次数",
- "ad_landing_capture_cnt": "广告落地页截图次数",
- "landing_no_scan_background_return_cnt": "无扫码情况下落地页切换后台再返回",
- "source_1007_cnt": "来源个人分享(场景值1007)",
- "source_1008_cnt": "来源群分享(场景值1008)",
- "source_1058_cnt": "来源公众号文章(场景值1058)",
- "source_1074_cnt": "来源公众号即转(场景值1074)",
- "source_1067_1095_cnt": "来源小程序投流(1067+1095)",
- "other_scene_list": "来源其他场景值",
- "source_shareid_list": "来源上游用户shareid",
- "source_opengid_list": "来源群id",
- }
- def load_users(path: Path, sample_per_group: int, lookback_days: int) -> pd.DataFrame:
- source = load_user_input(path, lookback_days=lookback_days)
- if sample_per_group > 0:
- source = (
- source.sort_values("mid", kind="stable")
- .groupby(["user_type", "sub_category"], as_index=False, group_keys=False)
- .head(sample_per_group)
- .reset_index(drop=True)
- )
- return source
- def latest_partition(client: ODPSClient, table_name: str) -> str:
- table = client.odps.get_table(table_name)
- values = [str(part.partition_spec["dt"]) for part in table.partitions]
- if not values:
- raise ValueError(f"{table_name}没有dt分区")
- return max(values)
- def build_queries(users: pd.DataFrame) -> dict[str, str]:
- records = users.to_dict("records")
- input_sql = build_input_users_sql(records)
- global_start = users["window_start_dt"].min()
- global_end = users["window_end_dt"].max()
- return {
- "core": build_core_sql(input_sql, global_start, global_end),
- "return": build_return_sql(
- input_sql,
- global_start,
- global_end,
- ),
- "attributes": build_attributes_sql(input_sql, global_start, global_end),
- "capture": build_capture_sql(input_sql, global_start, global_end),
- "background_return": build_background_return_sql(input_sql, global_start, global_end),
- "landing_return": build_landing_return_sql(input_sql, global_start, global_end),
- }
- def submit_all(client: ODPSClient, queries: dict[str, str]) -> tuple[dict[str, object], list[dict[str, str]]]:
- instances = {}
- metadata = []
- for name, sql in queries.items():
- instance = client.odps.run_sql(sql)
- instances[name] = instance
- item = {"name": name, "instance_id": instance.id, "logview": instance.get_logview_address()}
- metadata.append(item)
- print(f"[{name}] instance={item['instance_id']}")
- print(f"[{name}] logview={item['logview']}")
- return instances, metadata
- def read_all(instances: dict[str, object], raw_dir: Path) -> dict[str, pd.DataFrame]:
- results = {}
- for name, instance in instances.items():
- instance.wait_for_success()
- with instance.open_reader(tunnel=True) as reader:
- frame = reader.to_pandas()
- frame.columns = [str(column).lower() for column in frame.columns]
- frame.to_csv(raw_dir / f"{name}.csv", index=False)
- results[name] = frame
- print(f"[{name}] rows={len(frame)}")
- return results
- def merge_results(results: dict[str, pd.DataFrame], expected_users: int) -> pd.DataFrame:
- merged = results["core"]
- for name in ("return", "attributes", "capture", "background_return", "landing_return"):
- if results[name]["mid"].duplicated().any():
- raise AssertionError(f"{name}结果存在重复MID")
- if name == "return" and "return_people_cnt" in merged:
- merged = merged.drop(columns=["return_people_cnt"])
- merged = merged.merge(results[name], on="mid", how="left", validate="one_to_one")
- if len(merged) != expected_users or merged["mid"].nunique() != expected_users:
- raise AssertionError("合并后人数与输入不一致")
- merged["统计周期"] = (
- merged["window_start_dt"].astype(str)
- + "-"
- + merged["window_end_dt"].astype(str)
- )
- merged = merged.rename(columns=RENAME)
- for column in OUTPUT_COLUMNS:
- if column not in merged:
- merged[column] = 0
- text_columns = {
- "用户分类", "子分类", "统计周期", "用户id", "操作系统", "机型", "地域",
- "来源其他场景值", "来源上游用户shareid", "来源群id",
- }
- for column in OUTPUT_COLUMNS:
- if column in text_columns:
- merged[column] = merged[column].fillna("")
- else:
- merged[column] = pd.to_numeric(merged[column], errors="coerce").fillna(0).astype("int64")
- return merged[OUTPUT_COLUMNS].sort_values(["用户分类", "子分类", "用户id"], kind="stable")
- def main() -> None:
- parser = argparse.ArgumentParser(description=__doc__)
- parser.add_argument("--input", type=Path, default=DEFAULT_INPUT)
- parser.add_argument("--sample-per-group", type=int, default=0)
- parser.add_argument("--lookback-days", type=int, default=180)
- parser.add_argument("--sql-only", action="store_true")
- parser.add_argument("--run-name", default="full_220")
- parser.add_argument(
- "--modules",
- default="",
- help="Comma-separated modules to execute; missing module results are loaded from raw CSV files.",
- )
- parser.add_argument("--merge-only", action="store_true")
- parser.add_argument("--no-merge", action="store_true", help="Execute selected modules and save raw results only.")
- parser.add_argument(
- "--collect-existing",
- action="store_true",
- help="Collect results from the latest saved Instance ID for each selected module without resubmitting SQL.",
- )
- args = parser.parse_args()
- users = load_users(args.input, args.sample_per_group, args.lookback_days)
- output_dir = BASE_DIR / "output" / args.run_name / "behavior"
- sql_dir = BASE_DIR / "sql" / args.run_name / "behavior"
- raw_dir = output_dir / "raw"
- sql_dir.mkdir(parents=True, exist_ok=True)
- raw_dir.mkdir(parents=True, exist_ok=True)
- queries = build_queries(users)
- for name, sql in queries.items():
- (sql_dir / f"{name}.sql").write_text(sql + "\n", encoding="utf-8")
- users.to_csv(output_dir / "input_users.csv", index=False)
- if args.sql_only:
- print(f"rendered_sql={sql_dir}")
- return
- client = ODPSClient(project="loghubods")
- selected = [item.strip() for item in args.modules.split(",") if item.strip()]
- unknown = sorted(set(selected) - set(queries))
- if unknown:
- raise ValueError("未知查询模块: " + ", ".join(unknown))
- selected_queries = {name: queries[name] for name in selected} if selected else queries
- if args.merge_only:
- results = {}
- elif args.collect_existing:
- metadata_path = output_dir / "odps_runs.json"
- if not metadata_path.exists():
- raise ValueError(f"缺少ODPS执行记录: {metadata_path}")
- history = json.loads(metadata_path.read_text(encoding="utf-8"))
- latest_runs = {}
- for item in history:
- latest_runs[item["name"]] = item
- missing_runs = [name for name in selected_queries if name not in latest_runs]
- if missing_runs:
- raise ValueError("缺少模块Instance ID: " + ", ".join(missing_runs))
- instances = {
- name: client.odps.get_instance(latest_runs[name]["instance_id"])
- for name in selected_queries
- }
- results = read_all(instances, raw_dir)
- else:
- instances, metadata = submit_all(client, selected_queries)
- metadata_path = output_dir / "odps_runs.json"
- history = []
- if metadata_path.exists():
- history = json.loads(metadata_path.read_text(encoding="utf-8"))
- history.extend(metadata)
- metadata_path.write_text(
- json.dumps(history, ensure_ascii=False, indent=2) + "\n", encoding="utf-8"
- )
- results = read_all(instances, raw_dir)
- if args.no_merge:
- print(f"raw_results={raw_dir}")
- return
- for name in queries:
- if name in results:
- continue
- path = raw_dir / f"{name}.csv"
- if not path.exists():
- raise ValueError(f"缺少模块结果: {path}")
- frame = pd.read_csv(path)
- frame.columns = [str(column).lower() for column in frame.columns]
- results[name] = frame
- detail = merge_results(results, len(users))
- input_order = {mid: index for index, mid in enumerate(users["mid"])}
- detail["_input_order"] = detail["用户id"].map(input_order)
- detail = detail.sort_values("_input_order", kind="stable").drop(columns="_input_order")
- detail_path = output_dir / "user_detail.csv"
- detail.to_csv(detail_path, index=False)
- validation = detail.groupby(["用户分类", "子分类"], as_index=False).agg(用户数=("用户id", "nunique"))
- validation.to_csv(output_dir / "validation.csv", index=False)
- print(validation.to_string(index=False))
- print(f"detail={detail_path}")
- if __name__ == "__main__":
- main()
|