"""ODPS 数据读取与日聚合 SQL。""" from __future__ import annotations from datetime import datetime, timedelta from typing import Tuple from zoneinfo import ZoneInfo import pandas as pd from .metrics import GZH_CHANNEL, QIWEI_CHANNEL, SELF_CHANNEL from .odps_client import ODPSClient TABLE_NAME = "loghubods.opengid_base_data" SHANGHAI = ZoneInfo("Asia/Shanghai") def parse_yyyymmdd(value: str) -> datetime: try: return datetime.strptime(value, "%Y%m%d") except ValueError as exc: raise ValueError(f"日期必须是 YYYYMMDD,实际为: {value}") from exc def date_window(end_date: str) -> Tuple[str, str]: end = parse_yyyymmdd(end_date) return (end - timedelta(days=2)).strftime("%Y%m%d"), end.strftime("%Y%m%d") def resolve_end_date(client: ODPSClient, requested: str | None = None) -> str: """使用显式日期,或选择不晚于昨日的最新可用分区。""" if requested: parse_yyyymmdd(requested) return requested latest = ( client.odps.get_table("opengid_base_data") .get_max_partition() .partition_spec["dt"] ) parse_yyyymmdd(latest) yesterday = (datetime.now(SHANGHAI) - timedelta(days=1)).strftime("%Y%m%d") return min(latest, yesterday) def build_daily_sql(start_date: str, end_date: str) -> str: parse_yyyymmdd(start_date) parse_yyyymmdd(end_date) common_filter = f""" dt BETWEEN '{start_date}' AND '{end_date}' AND usersharedepth = '0' AND videoid IS NOT NULL AND NVL(hotsencetype, '') <> '1167' """ optimize_goal = ( "CASE WHEN 广告优化目标 IS NULL OR 广告优化目标='' " "OR 广告优化目标='null' THEN '' ELSE 广告优化目标 END" ) return f""" SELECT dt, 'self' AS entity_type, channel, MAX(NVL(代理名称, '')) AS 代理名称, NVL(账号id, '') AS 账号id, MAX(NVL(账号名称, '')) AS 账号名称, NVL(广告id, '') AS 广告id, MAX(NVL(广告名称, '')) AS 广告名称, NVL(包名, '') AS 包名, {optimize_goal} AS 广告优化目标, NVL(创意id, '') AS 创意id, '' AS 合作方名, '' AS 公众号名, COUNT(DISTINCT mid) AS 首层UV, SUM(NVL(t0裂变人数, 0)) AS T0裂变数, SUM(NVL(成本, 0)) AS 成本, SUM(NVL(效率收入, 0)) AS 效率收入, SUM(NVL(多层t0裂变收入, 0)) AS 多层裂变收入 FROM {TABLE_NAME} WHERE {common_filter} AND channel = '{SELF_CHANNEL}' AND 广告id IS NOT NULL AND 创意id IS NOT NULL GROUP BY dt, channel, 账号id, 广告id, 包名, {optimize_goal}, 创意id UNION ALL SELECT dt, 'gzh' AS entity_type, channel, '' AS 代理名称, '' AS 账号id, '' AS 账号名称, '' AS 广告id, '' AS 广告名称, '' AS 包名, '' AS 广告优化目标, '' AS 创意id, NVL(合作方名, '') AS 合作方名, NVL(公众号名, '') AS 公众号名, COUNT(DISTINCT mid) AS 首层UV, SUM(NVL(t0裂变人数, 0)) AS T0裂变数, SUM(NVL(成本, 0)) AS 成本, SUM(NVL(效率收入, 0)) AS 效率收入, SUM(NVL(多层t0裂变收入, 0)) AS 多层裂变收入 FROM {TABLE_NAME} WHERE {common_filter} AND channel = '{GZH_CHANNEL}' AND 公众号名 IS NOT NULL GROUP BY dt, channel, 合作方名, 公众号名 UNION ALL SELECT dt, 'qiwei' AS entity_type, channel, '' AS 代理名称, '' AS 账号id, '' AS 账号名称, '' AS 广告id, '' AS 广告名称, '' AS 包名, '' AS 广告优化目标, '' AS 创意id, NVL(合作方名, '') AS 合作方名, MAX(NVL(公众号名, '')) AS 公众号名, COUNT(DISTINCT mid) AS 首层UV, SUM(NVL(t0裂变人数, 0)) AS T0裂变数, SUM(NVL(成本, 0)) AS 成本, SUM(NVL(效率收入, 0)) AS 效率收入, SUM(NVL(多层t0裂变收入, 0)) AS 多层裂变收入 FROM {TABLE_NAME} WHERE {common_filter} AND channel = '{QIWEI_CHANNEL}' AND 合作方名 IS NOT NULL GROUP BY dt, channel, 合作方名 """.strip() def build_ad_age_sql(end_date: str, lookback_days: int = 30) -> str: end = parse_yyyymmdd(end_date) start_date = (end - timedelta(days=lookback_days - 1)).strftime("%Y%m%d") return f""" SELECT 广告id, MIN(dt) AS 首次出现日期 FROM {TABLE_NAME} WHERE dt BETWEEN '{start_date}' AND '{end_date}' AND usersharedepth = '0' AND channel = '{SELF_CHANNEL}' AND videoid IS NOT NULL AND NVL(hotsencetype, '') <> '1167' AND 广告id IS NOT NULL GROUP BY 广告id """.strip() def fetch_daily_data( client: ODPSClient, start_date: str, end_date: str, ) -> pd.DataFrame: return client.execute_sql(build_daily_sql(start_date, end_date)) def fetch_ad_age( client: ODPSClient, end_date: str, lookback_days: int = 30, ) -> pd.DataFrame: raw = client.execute_sql(build_ad_age_sql(end_date, lookback_days)) if raw.empty: return pd.DataFrame(columns=["广告id", "广告age"]) end = parse_yyyymmdd(end_date) result = raw.copy() result["广告id"] = result["广告id"].fillna("").astype(str) first_seen = pd.to_datetime(result["首次出现日期"].astype(str), format="%Y%m%d") result["广告age"] = (end - first_seen).dt.days + 1 return result[["广告id", "广告age"]]