user_timeline_realtime_batch.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104
  1. #!/usr/bin/env python3
  2. # coding=utf-8
  3. """多个实时用户的行为时间线,输出可筛选的单一 Excel。"""
  4. import argparse
  5. from datetime import datetime
  6. from pathlib import Path
  7. import pandas as pd
  8. from odps_module import ODPSClient
  9. from user_timeline import EXCLUDED_BUSINESSTYPES, add_device_info, behavior_definition, build_apptype_filter, safe_name, sql_text
  10. SQL_REALTIME_BATCH = """
  11. WITH v AS (
  12. SELECT CAST(clienttimestamp AS BIGINT) ts,apptype 产品apptype,mid 用户标识,'video' 来源,businesstype,CAST(NULL AS STRING) eventid,machineinfo_system 系统,pagesource,CAST(NULL AS STRING) endRoutePath,CAST(NULL AS STRING) objecttype,CAST(NULL AS STRING) topic,CAST(NULL AS STRING) shareid,CAST(NULL AS STRING) creativeCode,videoid 视频id,CAST(NULL AS STRING) rootPageSource,CAST(NULL AS STRING) objectId,CAST(NULL AS STRING) targetUid,CAST(NULL AS STRING) networkType,CAST(NULL AS STRING) operationOptions,CAST(NULL AS STRING) operationExtParams,CAST(NULL AS STRING) hotsencetype,CAST(NULL AS STRING) path,subsessionid,sessionid
  13. FROM loghubods.video_action_log_per5min WHERE dt>='{dt_start}' AND dt<='{dt_end}'{apptype_filter} AND mid IN ({mc_list}) AND businesstype<>'videoPreView' AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  14. ),a AS (
  15. SELECT CAST(clienttimestamp AS BIGINT) ts,apptype 产品apptype,machinecode 用户标识,'ad' 来源,businesstype,eventid,GET_JSON_OBJECT(machineinfo,'$.system') 系统,pagesource,CAST(NULL AS STRING) endRoutePath,CAST(NULL AS STRING) objecttype,CAST(NULL AS STRING) topic,CAST(NULL AS STRING) shareid,creativecode creativeCode,headvideoid 视频id,CAST(NULL AS STRING) rootPageSource,CAST(NULL AS STRING) objectId,CAST(NULL AS STRING) targetUid,CAST(NULL AS STRING) networkType,CAST(NULL AS STRING) operationOptions,CAST(NULL AS STRING) operationExtParams,hotsencetype,CAST(NULL AS STRING) path,subsessionid,sessionid
  16. FROM loghubods.ad_action_log_own_per5min WHERE dt>='{dt_start}' AND dt<='{dt_end}'{apptype_filter} AND machinecode IN ({mc_list}) AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  17. ),p AS (
  18. SELECT CAST(clienttimestamp AS BIGINT) ts,apptype 产品apptype,mid 用户标识,'play' 来源,businesstype,eventid,machineinfo_system 系统,pagesource,CAST(NULL AS STRING) endRoutePath,CAST(NULL AS STRING) objecttype,CAST(NULL AS STRING) topic,CAST(NULL AS STRING) shareid,CAST(NULL AS STRING) creativeCode,videoid 视频id,CAST(NULL AS STRING) rootPageSource,CAST(NULL AS STRING) objectId,CAST(NULL AS STRING) targetUid,CAST(NULL AS STRING) networkType,CAST(NULL AS STRING) operationOptions,CAST(NULL AS STRING) operationExtParams,CAST(NULL AS STRING) hotsencetype,CAST(NULL AS STRING) path,subsessionid,sessionid
  19. FROM loghubods.video_play_log_per5min WHERE dt>='{dt_start}' AND dt<='{dt_end}'{apptype_filter} AND mid IN ({mc_list}) AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  20. ),s AS (
  21. SELECT CAST(clienttimestamp AS BIGINT) ts,apptype 产品apptype,machinecode 用户标识,'simpleevent' 来源,businesstype,eventid,system 系统,pagesource,CAST(NULL AS STRING) endRoutePath,objecttype,CAST(NULL AS STRING) topic,CAST(NULL AS STRING) shareid,CAST(NULL AS STRING) creativeCode,videoid 视频id,CAST(NULL AS STRING) rootPageSource,CAST(NULL AS STRING) objectId,CAST(NULL AS STRING) targetUid,CAST(NULL AS STRING) networkType,CAST(NULL AS STRING) operationOptions,CAST(NULL AS STRING) operationExtParams,CAST(NULL AS STRING) hotsencetype,CAST(NULL AS STRING) path,subsessionid,sessionid
  22. FROM loghubods.simpleevent_log_flow WHERE year='{year}' AND month='{month}' AND day='{day_of_month}'{apptype_filter} AND machinecode IN ({mc_list}) AND (businesstype IS NULL OR businesstype<>'openGIdError') AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  23. ),u AS (
  24. SELECT CAST(clienttimestamp AS BIGINT) ts,apptype 产品apptype,machinecode 用户标识,'useractive' 来源,businesstype,eventid,system 系统,pagesource,CAST(NULL AS STRING) endRoutePath,CAST(NULL AS STRING) objecttype,CAST(NULL AS STRING) topic,CAST(NULL AS STRING) shareid,CAST(NULL AS STRING) creativeCode,CAST(NULL AS STRING) 视频id,CAST(NULL AS STRING) rootPageSource,CAST(NULL AS STRING) objectId,CAST(NULL AS STRING) targetUid,CAST(NULL AS STRING) networkType,CAST(NULL AS STRING) operationOptions,CAST(NULL AS STRING) operationExtParams,CAST(NULL AS STRING) hotsencetype,path,subsessionid,sessionid
  25. FROM loghubods.useractive_log_per5min WHERE dt>='{dt_start}' AND dt<='{dt_end}'{apptype_filter} AND machinecode IN ({mc_list}) AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  26. ),r AS (
  27. SELECT CAST(clienttimestamp AS BIGINT) ts,apptype 产品apptype,machinecode 用户标识,'share' 来源,type businesstype,eventid,CAST(NULL AS STRING) 系统,pagesource,CAST(NULL AS STRING) endRoutePath,CAST(NULL AS STRING) objecttype,topic,shareid,CAST(NULL AS STRING) creativeCode,CAST(NULL AS STRING) 视频id,rootpagesource rootPageSource,shareobjectid objectId,CAST(NULL AS STRING) targetUid,CAST(NULL AS STRING) networkType,CAST(NULL AS STRING) operationOptions,CAST(NULL AS STRING) operationExtParams,CAST(NULL AS STRING) hotsencetype,CAST(NULL AS STRING) path,subsessionid,sessionid
  28. FROM loghubods.user_share_log_per5min WHERE dt>='{dt_start}' AND dt<='{dt_end}'{apptype_filter} AND machinecode IN ({mc_list}) AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  29. ),o AS (
  30. SELECT CAST(clienttimestamp AS BIGINT) ts,apptype 产品apptype,machinecode 用户标识,'operation' 来源,CAST(NULL AS STRING) businesstype,eventid,machineinfo_system 系统,pagesource,CAST(NULL AS STRING) endRoutePath,CAST(NULL AS STRING) objecttype,CAST(NULL AS STRING) topic,CAST(NULL AS STRING) shareid,CAST(NULL AS STRING) creativeCode,CAST(NULL AS STRING) 视频id,rootpagesource rootPageSource,objectid objectId,targetuid targetUid,networktype networkType,options operationOptions,extparams operationExtParams,CAST(NULL AS STRING) hotsencetype,CAST(NULL AS STRING) path,subsessionid,sessionid
  31. FROM loghubods.operation_log_per5min WHERE dt>='{dt_start}' AND dt<='{dt_end}'{apptype_filter} AND machinecode IN ({mc_list}) AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  32. ),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 UNION ALL SELECT * FROM r UNION ALL SELECT * FROM o)
  33. SELECT t.ts,t.产品apptype,t.用户标识,t.来源,t.businesstype,t.eventid,t.系统,t.endRoutePath,t.objecttype,t.pagesource,t.topic,t.shareid,t.creativeCode,t.视频id,b.title 视频标题,t.rootPageSource,t.objectId,t.targetUid,t.networkType,t.operationOptions,t.operationExtParams,t.hotsencetype,t.path,t.subsessionid,t.sessionid
  34. FROM t LEFT JOIN videoods.dim_video b ON t.视频id=b.videoid ORDER BY t.ts
  35. """
  36. def load_user_strategies(path):
  37. users = pd.read_csv(path, dtype=str).fillna("")
  38. if "用户标识" not in users.columns:
  39. raise ValueError("用户清单必须包含“用户标识”列")
  40. strategy_column = "命中策略" if "命中策略" in users.columns else None
  41. strategies = {}
  42. for _, row in users.iterrows():
  43. user = row["用户标识"].strip()
  44. if user:
  45. strategy = row[strategy_column].strip() if strategy_column else ""
  46. strategies.setdefault(user, [])
  47. if strategy and strategy not in strategies[user]:
  48. strategies[user].append(strategy)
  49. if not strategies:
  50. raise ValueError("用户清单中没有有效的用户标识")
  51. return {user: ";".join(tags) for user, tags in strategies.items()}
  52. def sql_literals(values):
  53. return ",".join("'" + sql_text(value) + "'" for value in values)
  54. def main():
  55. parser = argparse.ArgumentParser(description="批量查询实时用户行为时间线")
  56. parser.add_argument("users_file", type=Path)
  57. parser.add_argument("date", nargs="?", default=datetime.now().strftime("%Y%m%d"), help="yyyyMMdd")
  58. parser.add_argument("apptype", nargs="?", default=None, help="可选;不传则查询全部产品")
  59. parser.add_argument("--output-dir", type=Path, default=Path("."))
  60. args = parser.parse_args()
  61. users = load_user_strategies(args.users_file)
  62. date = datetime.strptime(args.date, "%Y%m%d")
  63. sql_args = {"mc_list": sql_literals(users), "apptype_filter": build_apptype_filter(args.apptype), "dt_start": f"{args.date}000000", "dt_end": f"{args.date}235959", "year": date.strftime("%Y"), "month": date.strftime("%m"), "day_of_month": date.strftime("%d")}
  64. df = ODPSClient().execute_sql(SQL_REALTIME_BATCH.format(**sql_args))
  65. df = df[~df["businesstype"].isin(EXCLUDED_BUSINESSTYPES)].sort_values("ts", kind="stable").reset_index(drop=True)
  66. df = df.rename(columns={"creativecode": "creativeCode", "endroutepath": "endRoutePath", "rootpagesource": "rootPageSource", "objectid": "objectId", "targetuid": "targetUid", "networktype": "networkType", "operationoptions": "operationOptions", "operationextparams": "operationExtParams"})
  67. df.insert(0, "北京时间", pd.to_datetime(df["ts"], unit="ms", utc=True).dt.tz_convert("Asia/Shanghai").dt.tz_localize(None))
  68. df = add_device_info(df.drop(columns="ts"), group_column="用户标识")
  69. df.insert(2, "命中策略", df["用户标识"].map(users).fillna(""))
  70. df.insert(df.columns.get_loc("businesstype") + 1, "行为", df.apply(behavior_definition, axis=1))
  71. df = df.rename(columns={"businesstype": "事件类型", "eventid": "事件ID"})
  72. df.insert(0, "产品apptype", df.pop("产品apptype"))
  73. df.insert(1, "用户标识", df.pop("用户标识"))
  74. df.insert(2, "命中策略", df.pop("命中策略"))
  75. counts = df["用户标识"].value_counts()
  76. summary = pd.DataFrame({"用户标识": list(users), "命中策略": list(users.values()), "实时事件数": [counts.get(user, 0) for user in users]})
  77. args.output_dir.mkdir(parents=True, exist_ok=True)
  78. out = args.output_dir / f"timeline_realtime_{safe_name(args.users_file.stem)}_{args.date}.xlsx"
  79. with pd.ExcelWriter(out, engine="openpyxl", datetime_format="yyyy/mm/dd hh:mm:ss") as writer:
  80. df.to_excel(writer, sheet_name="实时行为路径", index=False)
  81. worksheet = writer.sheets["实时行为路径"]
  82. worksheet.freeze_panes = "A2"
  83. worksheet.auto_filter.ref = worksheet.dimensions
  84. for cell in worksheet["D"][1:]:
  85. cell.number_format = "yyyy/mm/dd hh:mm:ss"
  86. summary.to_excel(writer, sheet_name="用户摘要", index=False)
  87. writer.sheets["用户摘要"].freeze_panes = "A2"
  88. writer.sheets["用户摘要"].auto_filter.ref = writer.sheets["用户摘要"].dimensions
  89. print("注意:simpleevent_log_flow 仅保留当前短实时窗口。")
  90. print(f"[XLSX] 用户数={len(users)} 日期={args.date} 事件数={len(df)} -> {out.resolve()}", flush=True)
  91. if __name__ == "__main__":
  92. main()