data_source.py 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179
  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. return f"""
  42. SELECT
  43. dt,
  44. 'self' AS entity_type,
  45. channel,
  46. MAX(NVL(代理名称, '')) AS 代理名称,
  47. NVL(账号id, '') AS 账号id,
  48. MAX(NVL(账号名称, '')) AS 账号名称,
  49. NVL(广告id, '') AS 广告id,
  50. MAX(NVL(广告名称, '')) AS 广告名称,
  51. NVL(包名, '') AS 包名,
  52. NVL(创意id, '') AS 创意id,
  53. '' AS 合作方名,
  54. '' AS 公众号名,
  55. COUNT(DISTINCT mid) AS 首层UV,
  56. SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
  57. SUM(NVL(成本, 0)) AS 成本,
  58. SUM(NVL(效率收入, 0)) AS 效率收入,
  59. SUM(NVL(多层t0裂变收入, 0)) AS 多层裂变收入
  60. FROM {TABLE_NAME}
  61. WHERE {common_filter}
  62. AND channel = '{SELF_CHANNEL}'
  63. AND 广告id IS NOT NULL
  64. AND 创意id IS NOT NULL
  65. GROUP BY
  66. dt, channel, 账号id, 广告id, 包名, 创意id
  67. UNION ALL
  68. SELECT
  69. dt,
  70. 'gzh' AS entity_type,
  71. channel,
  72. '' AS 代理名称,
  73. '' AS 账号id,
  74. '' AS 账号名称,
  75. '' AS 广告id,
  76. '' AS 广告名称,
  77. '' AS 包名,
  78. '' AS 创意id,
  79. NVL(合作方名, '') AS 合作方名,
  80. NVL(公众号名, '') AS 公众号名,
  81. COUNT(DISTINCT mid) AS 首层UV,
  82. SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
  83. SUM(NVL(成本, 0)) AS 成本,
  84. SUM(NVL(效率收入, 0)) AS 效率收入,
  85. SUM(NVL(多层t0裂变收入, 0)) AS 多层裂变收入
  86. FROM {TABLE_NAME}
  87. WHERE {common_filter}
  88. AND channel = '{GZH_CHANNEL}'
  89. AND 公众号名 IS NOT NULL
  90. GROUP BY dt, channel, 合作方名, 公众号名
  91. UNION ALL
  92. SELECT
  93. dt,
  94. 'qiwei' AS entity_type,
  95. channel,
  96. '' AS 代理名称,
  97. '' AS 账号id,
  98. '' AS 账号名称,
  99. '' AS 广告id,
  100. '' AS 广告名称,
  101. '' AS 包名,
  102. '' AS 创意id,
  103. NVL(合作方名, '') AS 合作方名,
  104. MAX(NVL(公众号名, '')) AS 公众号名,
  105. COUNT(DISTINCT mid) AS 首层UV,
  106. SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
  107. SUM(NVL(成本, 0)) AS 成本,
  108. SUM(NVL(效率收入, 0)) AS 效率收入,
  109. SUM(NVL(多层t0裂变收入, 0)) AS 多层裂变收入
  110. FROM {TABLE_NAME}
  111. WHERE {common_filter}
  112. AND channel = '{QIWEI_CHANNEL}'
  113. AND 合作方名 IS NOT NULL
  114. GROUP BY dt, channel, 合作方名
  115. """.strip()
  116. def build_ad_age_sql(end_date: str, lookback_days: int = 30) -> str:
  117. end = parse_yyyymmdd(end_date)
  118. start_date = (end - timedelta(days=lookback_days - 1)).strftime("%Y%m%d")
  119. return f"""
  120. SELECT
  121. 广告id,
  122. MIN(dt) AS 首次出现日期
  123. FROM {TABLE_NAME}
  124. WHERE dt BETWEEN '{start_date}' AND '{end_date}'
  125. AND usersharedepth = '0'
  126. AND channel = '{SELF_CHANNEL}'
  127. AND videoid IS NOT NULL
  128. AND NVL(hotsencetype, '') <> '1167'
  129. AND 广告id IS NOT NULL
  130. GROUP BY 广告id
  131. """.strip()
  132. def fetch_daily_data(
  133. client: ODPSClient,
  134. start_date: str,
  135. end_date: str,
  136. ) -> pd.DataFrame:
  137. return client.execute_sql(build_daily_sql(start_date, end_date))
  138. def fetch_ad_age(
  139. client: ODPSClient,
  140. end_date: str,
  141. lookback_days: int = 30,
  142. ) -> pd.DataFrame:
  143. raw = client.execute_sql(build_ad_age_sql(end_date, lookback_days))
  144. if raw.empty:
  145. return pd.DataFrame(columns=["广告id", "广告age"])
  146. end = parse_yyyymmdd(end_date)
  147. result = raw.copy()
  148. result["广告id"] = result["广告id"].fillna("").astype(str)
  149. first_seen = pd.to_datetime(result["首次出现日期"].astype(str), format="%Y%m%d")
  150. result["广告age"] = (end - first_seen).dt.days + 1
  151. return result[["广告id", "广告age"]]