user_timeline_realtime.py 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133
  1. #!/usr/bin/env python
  2. # coding=utf-8
  3. """单用户实时行为时间线:基于实时/5分钟日志输出 Excel。
  4. 用法:python3 user_timeline_realtime.py <machinecode/mid值> [日期yyyyMMdd] [apptype] [--output-dir 目录]
  5. 注意:simpleevent_log_flow 仅保留短实时窗口,因此结果代表当前可用实时数据,
  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 = """
  15. WITH v AS (
  16. SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}'
  24. AND businesstype <> 'videoPreView'
  25. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  26. ),
  27. a AS (
  28. SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}'
  34. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  35. ),
  36. p AS (
  37. SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}'
  45. AND clienttimestamp IS NOT NULL AND clienttimestamp<>''
  46. ),
  47. s AS (
  48. SELECT CAST(clienttimestamp AS BIGINT) AS ts, '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='{mc}'
  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, '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='{mc}'
  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.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 main():
  81. parser = argparse.ArgumentParser(description="查询单用户实时行为时间线")
  82. parser.add_argument("user_id", help="machinecode/mid")
  83. parser.add_argument("date", nargs="?", default=datetime.now().strftime("%Y%m%d"), help="yyyyMMdd")
  84. parser.add_argument("apptype", nargs="?", default="0")
  85. parser.add_argument("--output-dir", type=Path, default=Path("."))
  86. args = parser.parse_args()
  87. query_date = datetime.strptime(args.date, "%Y%m%d")
  88. sql_args = {
  89. "mc": sql_text(args.user_id),
  90. "apptype": sql_text(args.apptype),
  91. "dt_start": f"{args.date}000000",
  92. "dt_end": f"{args.date}235959",
  93. "year": query_date.strftime("%Y"),
  94. "month": query_date.strftime("%m"),
  95. "day_of_month": query_date.strftime("%d"),
  96. }
  97. df = ODPSClient().execute_sql(SQL_REALTIME.format(**sql_args))
  98. df = df[~df["businesstype"].isin(EXCLUDED_BUSINESSTYPES)]
  99. df = df.sort_values("ts", kind="stable").reset_index(drop=True)
  100. df = df.rename(columns={"pagestatus": "pageStatus", "newpage": "newPage"})
  101. df["pageStatus"] = df["pageStatus"].map(
  102. lambda value: f"{text(value)}_{PAGESTATUS_MEAN[text(value)]}" if text(value) in PAGESTATUS_MEAN else text(value)
  103. )
  104. df.insert(0, "北京时间", pd.to_datetime(df["ts"], unit="ms", utc=True).dt.tz_convert("Asia/Shanghai").dt.tz_localize(None))
  105. df = df.drop(columns="ts")
  106. df.insert(df.columns.get_loc("businesstype") + 1, "中文行为定义", df.apply(behavior_definition, axis=1))
  107. output_dir = args.output_dir.expanduser()
  108. output_dir.mkdir(parents=True, exist_ok=True)
  109. out = output_dir / f"timeline_realtime_{safe_name(args.user_id[-12:])}_{args.date}.xlsx"
  110. with pd.ExcelWriter(out, engine="openpyxl", datetime_format="yyyy/mm/dd hh:mm:ss") as writer:
  111. df.to_excel(writer, sheet_name="实时行为路径", index=False)
  112. for cell in writer.sheets["实时行为路径"]["A"][1:]:
  113. cell.number_format = "yyyy/mm/dd hh:mm:ss"
  114. print("注意:simpleevent_log_flow 仅保留当前短实时窗口。")
  115. counts = df["来源"].value_counts().to_dict()
  116. for source in ["video", "ad", "play", "simpleevent", "useractive"]:
  117. print(f"[ROWS] {source}={counts.get(source, 0)}", flush=True)
  118. print(f"[XLSX] 日期={args.date} 事件数={len(df)} -> {out.resolve()}", flush=True)
  119. if __name__ == "__main__":
  120. main()