user_timeline_realtime_batch.py 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172
  1. #!/usr/bin/env python
  2. # coding=utf-8
  3. """多个实时风险用户的行为时间线,输出可按用户筛选的单一 Excel。
  4. 用法:python3 user_timeline_realtime_batch.py <用户清单.csv> [日期yyyyMMdd] [apptype] [--output-dir 目录]
  5. 清单必须包含“用户标识”列,可选“命中策略”列。相同用户会只查询一次,
  6. 并将多条策略标签合并到输出的“命中策略”列。
  7. """
  8. import argparse
  9. from datetime import datetime
  10. from pathlib import Path
  11. import pandas as pd
  12. from odps_module import ODPSClient
  13. from user_timeline import EXCLUDED_BUSINESSTYPES, PAGESTATUS_MEAN, behavior_definition, safe_name, sql_text, text
  14. SQL_REALTIME_BATCH = """
  15. WITH v AS (
  16. SELECT CAST(clienttimestamp AS BIGINT) AS ts, mid AS 用户标识, 'video' AS 来源,
  17. businesstype, pagesource, videoid AS 视频id,
  18. GET_JSON_OBJECT(extparams, '$.auto_enter') AS auto_enter,
  19. GET_JSON_OBJECT(extparams, '$.newPage') AS newPage,
  20. GET_JSON_OBJECT(extparams, '$.pageStatus') AS pageStatus,
  21. CAST(NULL AS STRING) AS hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
  22. FROM loghubods.video_action_log_per5min
  23. WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND mid IN ({mc_list})
  24. AND businesstype <> 'videoPreView'
  25. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  26. ),
  27. a AS (
  28. SELECT CAST(clienttimestamp AS BIGINT) AS ts, machinecode AS 用户标识, 'ad' AS 来源,
  29. businesstype, pagesource, headvideoid AS 视频id,
  30. CAST(NULL AS STRING) AS auto_enter, CAST(NULL AS STRING) AS newPage, CAST(NULL AS STRING) AS pageStatus,
  31. hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
  32. FROM loghubods.ad_action_log_own_per5min
  33. WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND machinecode IN ({mc_list})
  34. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  35. ),
  36. p AS (
  37. SELECT CAST(clienttimestamp AS BIGINT) AS ts, mid AS 用户标识, 'play' AS 来源,
  38. businesstype, pagesource, videoid AS 视频id,
  39. GET_JSON_OBJECT(extparams, '$.auto_enter') AS auto_enter,
  40. GET_JSON_OBJECT(extparams, '$.newPage') AS newPage,
  41. GET_JSON_OBJECT(extparams, '$.pageStatus') AS pageStatus,
  42. CAST(NULL AS STRING) AS hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
  43. FROM loghubods.video_play_log_per5min
  44. WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND mid IN ({mc_list})
  45. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  46. ),
  47. s AS (
  48. SELECT CAST(clienttimestamp AS BIGINT) AS ts, machinecode AS 用户标识, 'simpleevent' AS 来源,
  49. businesstype, pagesource, videoid AS 视频id,
  50. CAST(NULL AS STRING) AS auto_enter, CAST(NULL AS STRING) AS newPage, CAST(NULL AS STRING) AS pageStatus,
  51. CAST(NULL AS STRING) AS hotsencetype, CAST(NULL AS STRING) AS path, subsessionid, sessionid
  52. FROM loghubods.simpleevent_log_flow
  53. WHERE year='{year}' AND month='{month}' AND day='{day_of_month}'
  54. AND apptype='{apptype}' AND machinecode IN ({mc_list})
  55. AND (businesstype IS NULL OR businesstype <> 'openGIdError')
  56. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  57. ),
  58. u AS (
  59. SELECT CAST(clienttimestamp AS BIGINT) AS ts, machinecode AS 用户标识, 'useractive' AS 来源,
  60. businesstype, pagesource, CAST(NULL AS STRING) AS 视频id,
  61. CAST(NULL AS STRING) AS auto_enter, CAST(NULL AS STRING) AS newPage, CAST(NULL AS STRING) AS pageStatus,
  62. CAST(NULL AS STRING) AS hotsencetype, path, subsessionid, sessionid
  63. FROM loghubods.useractive_log_per5min
  64. WHERE dt >= '{dt_start}' AND dt <= '{dt_end}' AND apptype='{apptype}' AND machinecode IN ({mc_list})
  65. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  66. ),
  67. t AS (
  68. SELECT * FROM v
  69. UNION ALL SELECT * FROM a
  70. UNION ALL SELECT * FROM p
  71. UNION ALL SELECT * FROM s
  72. UNION ALL SELECT * FROM u
  73. )
  74. SELECT t.ts, t.用户标识, t.来源, t.businesstype, t.pagesource, t.视频id, b.title AS 视频标题,
  75. t.auto_enter, t.newPage, t.pageStatus, t.hotsencetype, t.path, t.subsessionid, t.sessionid
  76. FROM t
  77. LEFT JOIN videoods.dim_video b ON t.视频id = b.videoid
  78. ORDER BY t.ts
  79. """
  80. def load_user_strategies(path):
  81. users = pd.read_csv(path, dtype=str).fillna("")
  82. if "用户标识" not in users.columns:
  83. raise ValueError("用户清单必须包含“用户标识”列")
  84. strategy_column = "命中策略" if "命中策略" in users.columns else None
  85. strategies = {}
  86. for _, row in users.iterrows():
  87. user = row["用户标识"].strip()
  88. if not user:
  89. continue
  90. strategies.setdefault(user, [])
  91. strategy = row[strategy_column].strip() if strategy_column else ""
  92. if strategy and strategy not in strategies[user]:
  93. strategies[user].append(strategy)
  94. if not strategies:
  95. raise ValueError("用户清单中没有有效的用户标识")
  96. return {user: ";".join(tags) for user, tags in strategies.items()}
  97. def sql_literals(values):
  98. return ",".join("'" + sql_text(value) + "'" for value in values)
  99. def main():
  100. parser = argparse.ArgumentParser(description="批量查询实时用户行为时间线")
  101. parser.add_argument("users_file", type=Path)
  102. parser.add_argument("date", nargs="?", default=datetime.now().strftime("%Y%m%d"), help="yyyyMMdd")
  103. parser.add_argument("apptype", nargs="?", default="0")
  104. parser.add_argument("--output-dir", type=Path, default=Path("."))
  105. args = parser.parse_args()
  106. user_strategies = load_user_strategies(args.users_file)
  107. query_date = datetime.strptime(args.date, "%Y%m%d")
  108. sql_args = {
  109. "mc_list": sql_literals(user_strategies),
  110. "apptype": sql_text(args.apptype),
  111. "dt_start": f"{args.date}000000",
  112. "dt_end": f"{args.date}235959",
  113. "year": query_date.strftime("%Y"),
  114. "month": query_date.strftime("%m"),
  115. "day_of_month": query_date.strftime("%d"),
  116. }
  117. df = ODPSClient().execute_sql(SQL_REALTIME_BATCH.format(**sql_args))
  118. df = df[~df["businesstype"].isin(EXCLUDED_BUSINESSTYPES)]
  119. df = df.sort_values("ts", kind="stable").reset_index(drop=True)
  120. df = df.rename(columns={"pagestatus": "pageStatus", "newpage": "newPage"})
  121. df["pageStatus"] = df["pageStatus"].map(
  122. lambda value: f"{text(value)}_{PAGESTATUS_MEAN[text(value)]}" if text(value) in PAGESTATUS_MEAN else text(value)
  123. )
  124. df.insert(0, "北京时间", pd.to_datetime(df["ts"], unit="ms", utc=True).dt.tz_convert("Asia/Shanghai").dt.tz_localize(None))
  125. df = df.drop(columns="ts")
  126. df.insert(2, "命中策略", df["用户标识"].map(user_strategies).fillna(""))
  127. df.insert(df.columns.get_loc("businesstype") + 1, "中文行为定义", df.apply(behavior_definition, axis=1))
  128. event_counts = df["用户标识"].value_counts()
  129. summary = pd.DataFrame(
  130. {
  131. "用户标识": list(user_strategies),
  132. "命中策略": list(user_strategies.values()),
  133. "实时事件数": [event_counts.get(user, 0) for user in user_strategies],
  134. }
  135. )
  136. output_dir = args.output_dir.expanduser()
  137. output_dir.mkdir(parents=True, exist_ok=True)
  138. cohort = safe_name(args.users_file.stem)
  139. out = output_dir / f"timeline_realtime_{cohort}_{args.date}.xlsx"
  140. with pd.ExcelWriter(out, engine="openpyxl", datetime_format="yyyy/mm/dd hh:mm:ss") as writer:
  141. df.to_excel(writer, sheet_name="实时行为路径", index=False)
  142. worksheet = writer.sheets["实时行为路径"]
  143. worksheet.freeze_panes = "A2"
  144. worksheet.auto_filter.ref = worksheet.dimensions
  145. for cell in worksheet["A"][1:]:
  146. cell.number_format = "yyyy/mm/dd hh:mm:ss"
  147. summary.to_excel(writer, sheet_name="用户摘要", index=False)
  148. summary_sheet = writer.sheets["用户摘要"]
  149. summary_sheet.freeze_panes = "A2"
  150. summary_sheet.auto_filter.ref = summary_sheet.dimensions
  151. print("注意:simpleevent_log_flow 仅保留当前短实时窗口。")
  152. print(f"[XLSX] 用户数={len(user_strategies)} 日期={args.date} 事件数={len(df)} -> {out.resolve()}", flush=True)
  153. if __name__ == "__main__":
  154. main()