#!/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()