| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246 |
- """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, SELF_CHANNEL
- from .odps_client import ODPSClient
- TABLE_NAME = "loghubods.opengid_base_data"
- SHANGHAI = ZoneInfo("Asia/Shanghai")
- class SourceDataNotReadyError(RuntimeError):
- """The exact ROI source partition required for this run is not ready."""
- 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 build_source_readiness_sql(end_date: str) -> str:
- parse_yyyymmdd(end_date)
- return f"""
- SELECT
- COUNT(1) AS row_count,
- SUM(CASE WHEN channel = '{SELF_CHANNEL}' THEN 1 ELSE 0 END) AS self_rows,
- SUM(CASE WHEN channel = '{GZH_CHANNEL}' THEN 1 ELSE 0 END) AS gzh_rows
- FROM {TABLE_NAME}
- WHERE dt = '{end_date}'
- AND usersharedepth <= 1
- AND videoid IS NOT NULL
- AND NVL(hotsencetype, '') <> '1167'
- """.strip()
- def validate_source_ready(client: ODPSClient, end_date: str) -> dict[str, int]:
- """Require an exact, non-empty T-1 partition for every reported channel."""
- max_partition = client.odps.get_table("opengid_base_data").get_max_partition()
- if max_partition is None:
- raise SourceDataNotReadyError(
- f"ROI来源表没有可用分区: table={TABLE_NAME}, required_dt={end_date}"
- )
- latest = max_partition.partition_spec["dt"]
- parse_yyyymmdd(latest)
- if latest < end_date:
- raise SourceDataNotReadyError(
- f"ROI来源表未同步目标分区: table={TABLE_NAME}, "
- f"required_dt={end_date}, latest_dt={latest}"
- )
- frame = client.execute_sql(build_source_readiness_sql(end_date))
- if frame.empty:
- raise SourceDataNotReadyError(
- f"ROI来源表目标分区为空: table={TABLE_NAME}, dt={end_date}"
- )
- row = frame.iloc[0]
- counts = {}
- for name in ("row_count", "self_rows", "gzh_rows"):
- value = row.get(name)
- counts[name] = 0 if pd.isna(value) else int(value)
- missing = [name for name, value in counts.items() if value <= 0]
- if missing:
- raise SourceDataNotReadyError(
- f"ROI来源表目标分区业务数据未就绪: table={TABLE_NAME}, "
- f"dt={end_date}, missing={','.join(missing)}"
- )
- return counts
- def resolve_end_date(
- client: ODPSClient,
- requested: str | None = None,
- *,
- now: datetime | None = None,
- ) -> str:
- """Resolve exactly the requested date or T-1, then require source readiness."""
- if requested:
- parse_yyyymmdd(requested)
- target = requested
- else:
- current = now or datetime.now(SHANGHAI)
- if current.tzinfo is None:
- current = current.replace(tzinfo=SHANGHAI)
- target = (current.astimezone(SHANGHAI) - timedelta(days=1)).strftime(
- "%Y%m%d"
- )
- validate_source_ready(client, target)
- return target
- 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 <= 1
- 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_fission_uv_root, 0)) AS T0裂变数,
- SUM(NVL(成本, 0)) AS 成本,
- SUM(NVL(效率收入, 0)) AS 效率收入,
- SUM(NVL(裂变效率收入, 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,
- 'self_ad' 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 广告优化目标,
- '' AS 创意id,
- '' AS 合作方名,
- '' AS 公众号名,
- COUNT(DISTINCT mid) AS 首层UV,
- SUM(NVL(t0_fission_uv_root, 0)) AS T0裂变数,
- SUM(NVL(成本, 0)) AS 成本,
- SUM(NVL(效率收入, 0)) AS 效率收入,
- SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入
- FROM {TABLE_NAME}
- WHERE {common_filter}
- AND channel = '{SELF_CHANNEL}'
- AND 广告id IS NOT NULL
- GROUP BY
- dt, channel, 账号id, 广告id, 包名, {optimize_goal}
- 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_fission_uv_root, 0)) AS T0裂变数,
- SUM(NVL(成本, 0)) AS 成本,
- SUM(NVL(效率收入, 0)) AS 效率收入,
- SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入
- FROM {TABLE_NAME}
- WHERE {common_filter}
- AND channel = '{GZH_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"]]
|