| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172 |
- #!/usr/bin/env python
- # coding=utf-8
- """多个实时风险用户的行为时间线,输出可按用户筛选的单一 Excel。
- 用法:python3 user_timeline_realtime_batch.py <用户清单.csv> [日期yyyyMMdd] [apptype] [--output-dir 目录]
- 清单必须包含“用户标识”列,可选“命中策略”列。相同用户会只查询一次,
- 并将多条策略标签合并到输出的“命中策略”列。
- """
- import argparse
- from datetime import datetime
- from pathlib import Path
- import pandas as pd
- from odps_module import ODPSClient
- from user_timeline import EXCLUDED_BUSINESSTYPES, PAGESTATUS_MEAN, behavior_definition, safe_name, sql_text, text
- SQL_REALTIME_BATCH = """
- WITH v AS (
- SELECT CAST(clienttimestamp AS BIGINT) AS ts, mid AS 用户标识, 'video' AS 来源,
- businesstype, pagesource, videoid AS 视频id,
- GET_JSON_OBJECT(extparams, '$.auto_enter') AS auto_enter,
- GET_JSON_OBJECT(extparams, '$.newPage') AS newPage,
- GET_JSON_OBJECT(extparams, '$.pageStatus') AS pageStatus,
- CAST(NULL AS STRING) AS hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
- FROM loghubods.video_action_log_per5min
- WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND mid IN ({mc_list})
- AND businesstype <> 'videoPreView'
- AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
- ),
- a AS (
- SELECT CAST(clienttimestamp AS BIGINT) AS ts, machinecode AS 用户标识, 'ad' AS 来源,
- businesstype, pagesource, headvideoid AS 视频id,
- CAST(NULL AS STRING) AS auto_enter, CAST(NULL AS STRING) AS newPage, CAST(NULL AS STRING) AS pageStatus,
- hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
- FROM loghubods.ad_action_log_own_per5min
- WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND machinecode IN ({mc_list})
- AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
- ),
- p AS (
- SELECT CAST(clienttimestamp AS BIGINT) AS ts, mid AS 用户标识, 'play' AS 来源,
- businesstype, pagesource, videoid AS 视频id,
- GET_JSON_OBJECT(extparams, '$.auto_enter') AS auto_enter,
- GET_JSON_OBJECT(extparams, '$.newPage') AS newPage,
- GET_JSON_OBJECT(extparams, '$.pageStatus') AS pageStatus,
- CAST(NULL AS STRING) AS hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
- FROM loghubods.video_play_log_per5min
- WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND mid IN ({mc_list})
- AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
- ),
- s AS (
- SELECT CAST(clienttimestamp AS BIGINT) AS ts, machinecode AS 用户标识, 'simpleevent' AS 来源,
- businesstype, pagesource, videoid AS 视频id,
- CAST(NULL AS STRING) AS auto_enter, CAST(NULL AS STRING) AS newPage, CAST(NULL AS STRING) AS pageStatus,
- CAST(NULL AS STRING) AS hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
- FROM loghubods.simpleevent_log_flow
- WHERE year='{year}' AND month='{month}' AND day='{day_of_month}'
- AND apptype='{apptype}' AND machinecode IN ({mc_list})
- AND (businesstype IS NULL OR businesstype <> 'openGIdError')
- AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
- ),
- u AS (
- SELECT CAST(clienttimestamp AS BIGINT) AS ts, machinecode AS 用户标识, 'useractive' AS 来源,
- businesstype, pagesource, CAST(NULL AS STRING) AS 视频id,
- CAST(NULL AS STRING) AS auto_enter, CAST(NULL AS STRING) AS newPage, CAST(NULL AS STRING) AS pageStatus,
- CAST(NULL AS STRING) AS hotsencetype, path, subsessionid, sessionid
- FROM loghubods.useractive_log_per5min
- WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND machinecode IN ({mc_list})
- AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
- ),
- t AS (
- SELECT * FROM v
- UNION ALL SELECT * FROM a
- UNION ALL SELECT * FROM p
- UNION ALL SELECT * FROM s
- UNION ALL SELECT * FROM u
- )
- SELECT t.ts, t.用户标识, t.来源, t.businesstype, t.pagesource, t.视频id, b.title AS 视频标题,
- t.auto_enter, t.newPage, t.pageStatus, t.hotsencetype, t.path, t.subsessionid, t.sessionid
- FROM t
- LEFT JOIN videoods.dim_video b ON t.视频id = b.videoid
- ORDER BY t.ts
- """
- def load_user_strategies(path):
- users = pd.read_csv(path, dtype=str).fillna("")
- if "用户标识" not in users.columns:
- raise ValueError("用户清单必须包含“用户标识”列")
- strategy_column = "命中策略" if "命中策略" in users.columns else None
- strategies = {}
- for _, row in users.iterrows():
- user = row["用户标识"].strip()
- if not user:
- continue
- strategies.setdefault(user, [])
- strategy = row[strategy_column].strip() if strategy_column else ""
- if strategy and strategy not in strategies[user]:
- strategies[user].append(strategy)
- if not strategies:
- raise ValueError("用户清单中没有有效的用户标识")
- return {user: ";".join(tags) for user, tags in strategies.items()}
- def sql_literals(values):
- return ",".join("'" + sql_text(value) + "'" for value in values)
- def main():
- parser = argparse.ArgumentParser(description="批量查询实时用户行为时间线")
- parser.add_argument("users_file", type=Path)
- parser.add_argument("date", nargs="?", default=datetime.now().strftime("%Y%m%d"), help="yyyyMMdd")
- parser.add_argument("apptype", nargs="?", default="0")
- parser.add_argument("--output-dir", type=Path, default=Path("."))
- args = parser.parse_args()
- user_strategies = load_user_strategies(args.users_file)
- query_date = datetime.strptime(args.date, "%Y%m%d")
- sql_args = {
- "mc_list": sql_literals(user_strategies),
- "apptype": sql_text(args.apptype),
- "dt_start": f"{args.date}000000",
- "dt_end": f"{args.date}235959",
- "year": query_date.strftime("%Y"),
- "month": query_date.strftime("%m"),
- "day_of_month": query_date.strftime("%d"),
- }
- df = ODPSClient().execute_sql(SQL_REALTIME_BATCH.format(**sql_args))
- df = df[~df["businesstype"].isin(EXCLUDED_BUSINESSTYPES)]
- df = df.sort_values("ts", kind="stable").reset_index(drop=True)
- df = df.rename(columns={"pagestatus": "pageStatus", "newpage": "newPage"})
- df["pageStatus"] = df["pageStatus"].map(
- lambda value: f"{text(value)}_{PAGESTATUS_MEAN[text(value)]}" if text(value) in PAGESTATUS_MEAN else text(value)
- )
- df.insert(0, "北京时间", pd.to_datetime(df["ts"], unit="ms", utc=True).dt.tz_convert("Asia/Shanghai").dt.tz_localize(None))
- df = df.drop(columns="ts")
- df.insert(2, "命中策略", df["用户标识"].map(user_strategies).fillna(""))
- df.insert(df.columns.get_loc("businesstype") + 1, "中文行为定义", df.apply(behavior_definition, axis=1))
- event_counts = df["用户标识"].value_counts()
- summary = pd.DataFrame(
- {
- "用户标识": list(user_strategies),
- "命中策略": list(user_strategies.values()),
- "实时事件数": [event_counts.get(user, 0) for user in user_strategies],
- }
- )
- output_dir = args.output_dir.expanduser()
- output_dir.mkdir(parents=True, exist_ok=True)
- cohort = safe_name(args.users_file.stem)
- out = output_dir / f"timeline_realtime_{cohort}_{args.date}.xlsx"
- with pd.ExcelWriter(out, engine="openpyxl", datetime_format="yyyy/mm/dd hh:mm:ss") as writer:
- df.to_excel(writer, sheet_name="实时行为路径", index=False)
- worksheet = writer.sheets["实时行为路径"]
- worksheet.freeze_panes = "A2"
- worksheet.auto_filter.ref = worksheet.dimensions
- for cell in worksheet["A"][1:]:
- cell.number_format = "yyyy/mm/dd hh:mm:ss"
- summary.to_excel(writer, sheet_name="用户摘要", index=False)
- summary_sheet = writer.sheets["用户摘要"]
- summary_sheet.freeze_panes = "A2"
- summary_sheet.auto_filter.ref = summary_sheet.dimensions
- print("注意:simpleevent_log_flow 仅保留当前短实时窗口。")
- print(f"[XLSX] 用户数={len(user_strategies)} 日期={args.date} 事件数={len(df)} -> {out.resolve()}", flush=True)
- if __name__ == "__main__":
- main()
|