data_source.py 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186
  1. """ODPS 数据读取与日聚合 SQL。"""
  2. from __future__ import annotations
  3. from datetime import datetime, timedelta
  4. from typing import Tuple
  5. from zoneinfo import ZoneInfo
  6. import pandas as pd
  7. from .metrics import GZH_CHANNEL, QIWEI_CHANNEL, SELF_CHANNEL
  8. from .odps_client import ODPSClient
  9. TABLE_NAME = "loghubods.opengid_base_data"
  10. SHANGHAI = ZoneInfo("Asia/Shanghai")
  11. def parse_yyyymmdd(value: str) -> datetime:
  12. try:
  13. return datetime.strptime(value, "%Y%m%d")
  14. except ValueError as exc:
  15. raise ValueError(f"日期必须是 YYYYMMDD,实际为: {value}") from exc
  16. def date_window(end_date: str) -> Tuple[str, str]:
  17. end = parse_yyyymmdd(end_date)
  18. return (end - timedelta(days=2)).strftime("%Y%m%d"), end.strftime("%Y%m%d")
  19. def resolve_end_date(client: ODPSClient, requested: str | None = None) -> str:
  20. """使用显式日期,或选择不晚于昨日的最新可用分区。"""
  21. if requested:
  22. parse_yyyymmdd(requested)
  23. return requested
  24. latest = (
  25. client.odps.get_table("opengid_base_data")
  26. .get_max_partition()
  27. .partition_spec["dt"]
  28. )
  29. parse_yyyymmdd(latest)
  30. yesterday = (datetime.now(SHANGHAI) - timedelta(days=1)).strftime("%Y%m%d")
  31. return min(latest, yesterday)
  32. def build_daily_sql(start_date: str, end_date: str) -> str:
  33. parse_yyyymmdd(start_date)
  34. parse_yyyymmdd(end_date)
  35. common_filter = f"""
  36. dt BETWEEN '{start_date}' AND '{end_date}'
  37. AND usersharedepth = '0'
  38. AND videoid IS NOT NULL
  39. AND NVL(hotsencetype, '') <> '1167'
  40. """
  41. optimize_goal = (
  42. "CASE WHEN 广告优化目标 IS NULL OR 广告优化目标='' "
  43. "OR 广告优化目标='null' THEN '' ELSE 广告优化目标 END"
  44. )
  45. return f"""
  46. SELECT
  47. dt,
  48. 'self' AS entity_type,
  49. channel,
  50. MAX(NVL(代理名称, '')) AS 代理名称,
  51. NVL(账号id, '') AS 账号id,
  52. MAX(NVL(账号名称, '')) AS 账号名称,
  53. NVL(广告id, '') AS 广告id,
  54. MAX(NVL(广告名称, '')) AS 广告名称,
  55. NVL(包名, '') AS 包名,
  56. {optimize_goal} AS 广告优化目标,
  57. NVL(创意id, '') AS 创意id,
  58. '' AS 合作方名,
  59. '' AS 公众号名,
  60. COUNT(DISTINCT mid) AS 首层UV,
  61. SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
  62. SUM(NVL(成本, 0)) AS 成本,
  63. SUM(NVL(效率收入, 0)) AS 效率收入,
  64. SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入
  65. FROM {TABLE_NAME}
  66. WHERE {common_filter}
  67. AND channel = '{SELF_CHANNEL}'
  68. AND 广告id IS NOT NULL
  69. AND 创意id IS NOT NULL
  70. GROUP BY
  71. dt, channel, 账号id, 广告id, 包名, {optimize_goal}, 创意id
  72. UNION ALL
  73. SELECT
  74. dt,
  75. 'gzh' AS entity_type,
  76. channel,
  77. '' AS 代理名称,
  78. '' AS 账号id,
  79. '' AS 账号名称,
  80. '' AS 广告id,
  81. '' AS 广告名称,
  82. '' AS 包名,
  83. '' AS 广告优化目标,
  84. '' AS 创意id,
  85. NVL(合作方名, '') AS 合作方名,
  86. NVL(公众号名, '') AS 公众号名,
  87. COUNT(DISTINCT mid) AS 首层UV,
  88. SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
  89. SUM(NVL(成本, 0)) AS 成本,
  90. SUM(NVL(效率收入, 0)) AS 效率收入,
  91. SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入
  92. FROM {TABLE_NAME}
  93. WHERE {common_filter}
  94. AND channel = '{GZH_CHANNEL}'
  95. AND 公众号名 IS NOT NULL
  96. GROUP BY dt, channel, 合作方名, 公众号名
  97. UNION ALL
  98. SELECT
  99. dt,
  100. 'qiwei' AS entity_type,
  101. channel,
  102. '' AS 代理名称,
  103. '' AS 账号id,
  104. '' AS 账号名称,
  105. '' AS 广告id,
  106. '' AS 广告名称,
  107. '' AS 包名,
  108. '' AS 广告优化目标,
  109. '' AS 创意id,
  110. NVL(合作方名, '') AS 合作方名,
  111. MAX(NVL(公众号名, '')) AS 公众号名,
  112. COUNT(DISTINCT mid) AS 首层UV,
  113. SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
  114. SUM(NVL(成本, 0)) AS 成本,
  115. SUM(NVL(效率收入, 0)) AS 效率收入,
  116. SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入
  117. FROM {TABLE_NAME}
  118. WHERE {common_filter}
  119. AND channel = '{QIWEI_CHANNEL}'
  120. AND 合作方名 IS NOT NULL
  121. GROUP BY dt, channel, 合作方名
  122. """.strip()
  123. def build_ad_age_sql(end_date: str, lookback_days: int = 30) -> str:
  124. end = parse_yyyymmdd(end_date)
  125. start_date = (end - timedelta(days=lookback_days - 1)).strftime("%Y%m%d")
  126. return f"""
  127. SELECT
  128. 广告id,
  129. MIN(dt) AS 首次出现日期
  130. FROM {TABLE_NAME}
  131. WHERE dt BETWEEN '{start_date}' AND '{end_date}'
  132. AND usersharedepth = '0'
  133. AND channel = '{SELF_CHANNEL}'
  134. AND videoid IS NOT NULL
  135. AND NVL(hotsencetype, '') <> '1167'
  136. AND 广告id IS NOT NULL
  137. GROUP BY 广告id
  138. """.strip()
  139. def fetch_daily_data(
  140. client: ODPSClient,
  141. start_date: str,
  142. end_date: str,
  143. ) -> pd.DataFrame:
  144. return client.execute_sql(build_daily_sql(start_date, end_date))
  145. def fetch_ad_age(
  146. client: ODPSClient,
  147. end_date: str,
  148. lookback_days: int = 30,
  149. ) -> pd.DataFrame:
  150. raw = client.execute_sql(build_ad_age_sql(end_date, lookback_days))
  151. if raw.empty:
  152. return pd.DataFrame(columns=["广告id", "广告age"])
  153. end = parse_yyyymmdd(end_date)
  154. result = raw.copy()
  155. result["广告id"] = result["广告id"].fillna("").astype(str)
  156. first_seen = pd.to_datetime(result["首次出现日期"].astype(str), format="%Y%m%d")
  157. result["广告age"] = (end - first_seen).dt.days + 1
  158. return result[["广告id", "广告age"]]