#!/usr/bin/env python # coding=utf-8 """单用户实时行为时间线:基于实时/5分钟日志输出 Excel。 用法:python3 user_timeline_realtime.py [日期yyyyMMdd] [apptype] [--output-dir 目录] 注意:simpleevent_log_flow 仅保留短实时窗口,因此结果代表当前可用实时数据, 不等同于离线表的整日完整回灌数据。 """ 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 = """ WITH v AS ( SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}' AND businesstype <> 'videoPreView' AND clienttimestamp IS NOT NULL AND clienttimestamp<>'' ), a AS ( SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}' AND clienttimestamp IS NOT NULL AND clienttimestamp<>'' ), p AS ( SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}' AND clienttimestamp IS NOT NULL AND clienttimestamp<>'' ), s AS ( SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}' AND (businesstype IS NULL OR businesstype <> 'openGIdError') AND clienttimestamp IS NOT NULL AND clienttimestamp<>'' ), u AS ( SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}' 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.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 main(): parser = argparse.ArgumentParser(description="查询单用户实时行为时间线") parser.add_argument("user_id", help="machinecode/mid") 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() query_date = datetime.strptime(args.date, "%Y%m%d") sql_args = { "mc": sql_text(args.user_id), "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.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(df.columns.get_loc("businesstype") + 1, "中文行为定义", df.apply(behavior_definition, axis=1)) output_dir = args.output_dir.expanduser() output_dir.mkdir(parents=True, exist_ok=True) out = output_dir / f"timeline_realtime_{safe_name(args.user_id[-12:])}_{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) for cell in writer.sheets["实时行为路径"]["A"][1:]: cell.number_format = "yyyy/mm/dd hh:mm:ss" print("注意:simpleevent_log_flow 仅保留当前短实时窗口。") counts = df["来源"].value_counts().to_dict() for source in ["video", "ad", "play", "simpleevent", "useractive"]: print(f"[ROWS] {source}={counts.get(source, 0)}", flush=True) print(f"[XLSX] 日期={args.date} 事件数={len(df)} -> {out.resolve()}", flush=True) if __name__ == "__main__": main()