odps_source.py 2.1 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374
  1. """ODPS source for the overall hourly CPM signal."""
  2. from __future__ import annotations
  3. import os
  4. from dataclasses import dataclass
  5. from datetime import date
  6. import pandas as pd
  7. from odps import ODPS
  8. @dataclass(frozen=True)
  9. class HourlyCpm:
  10. partition: str
  11. hour_of_day: int
  12. impressions: int
  13. cpm: float
  14. def build_odps_client() -> ODPS:
  15. missing = [
  16. key
  17. for key in ("ODPS_ACCESS_ID", "ODPS_ACCESS_SECRET")
  18. if not os.getenv(key)
  19. ]
  20. if missing:
  21. raise RuntimeError(f"Missing ODPS environment variables: {', '.join(missing)}")
  22. return ODPS(
  23. os.environ["ODPS_ACCESS_ID"],
  24. os.environ["ODPS_ACCESS_SECRET"],
  25. os.getenv("ODPS_PROJECT", "loghubods"),
  26. endpoint=os.getenv("ODPS_ENDPOINT", "http://service.odps.aliyun.com/api"),
  27. )
  28. def fetch_hourly_cpm(client: ODPS, data_date: str) -> pd.DataFrame:
  29. sql = f"""
  30. SELECT dt,
  31. CAST(SUBSTR(dt, 9, 2) AS BIGINT) AS hour_of_day,
  32. `曝光次数_总` AS impressions,
  33. `真实cpm_总` AS cpm
  34. FROM loghubods.advertiser_data_da_hour
  35. WHERE dt LIKE '{data_date}%'
  36. AND `名称` = 'SUM'
  37. AND advertisercode = 'SUM'
  38. AND company = 'SUM'
  39. AND `行业` = 'SUM'
  40. AND `客户` = 'SUM'
  41. AND `落地页类型` = 'SUM'
  42. ORDER BY dt
  43. """
  44. instance = client.execute_sql(sql, hints={"odps.sql.submit.mode": "script"})
  45. with instance.open_reader(tunnel=True) as reader:
  46. frame = reader.to_pandas()
  47. if frame.empty:
  48. return frame
  49. frame["hour_of_day"] = pd.to_numeric(frame["hour_of_day"]).astype("int64")
  50. frame["impressions"] = pd.to_numeric(frame["impressions"]).astype("int64")
  51. frame["cpm"] = pd.to_numeric(frame["cpm"]).astype(float)
  52. return frame
  53. def fetch_latest_cpm(client: ODPS, data_date: date) -> HourlyCpm | None:
  54. frame = fetch_hourly_cpm(client, data_date.strftime("%Y%m%d"))
  55. if frame.empty:
  56. return None
  57. row = frame.iloc[-1]
  58. return HourlyCpm(
  59. partition=str(row["dt"]),
  60. hour_of_day=int(row["hour_of_day"]),
  61. impressions=int(row["impressions"]),
  62. cpm=float(row["cpm"]),
  63. )