| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401 |
- """承接视频获取接口适配层(piaoquantv API)。
- 业务模型(用户 2026-06-08 确认):
- - "承接视频" = 用户点开创意/进入小程序后看到的视频内容
- - 接口已按 ROV 降序 + 访问量过滤,我们拿前 N 即可,**无需自己排序**
- - 这一步**不是**创意素材召回 — 创意素材在后续 material_recall 步用本视频字段
- (title / standardElement / categoryName / demandContentTitle 等)做 query 召回
- 接口:POST https://tp-open.piaoquantv.com/contentPlatform/plan/videoContentList
- 认证:token header(沿用用户给的值,可通过 env PIAOQUANTV_TOKEN 覆盖)
- 字段参考(2026-06-08 反推自 sample 调用):
- 返回 data.objs[*] 39 个字段,关键:
- videoId : 业务侧 8 位 video ID(注:非腾讯素材库 video_id,11 位)
- title / cover / video : 标题 + 封面 URL + 视频 URL
- score / rov / sim / visitUv : 推荐分(已排序)
- standardElement / categoryName / demandContentTitle / dimension : 语义维度
- """
- import logging
- import os
- import json
- from dataclasses import dataclass, field
- from pathlib import Path
- from typing import List, Optional
- import httpx
- try:
- from dotenv import load_dotenv
- load_dotenv(Path(__file__).resolve().parents[1] / ".env")
- except Exception:
- pass
- logger = logging.getLogger(__name__)
- PIAOQUANTV_VIDEO_API = "https://tp-open.piaoquantv.com/contentPlatform/plan/videoContentList"
- PIAOQUANTV_TOKEN = os.getenv(
- "PIAOQUANTV_TOKEN",
- "f6e07dba7fe3476cb31fd3733d607c5c", # 用户 2026-06-08 提供
- )
- PIAOQUANTV_VIDEO_SOURCE = os.getenv("PIAOQUANTV_VIDEO_SOURCE", "")
- _PIAOQUANTV_VIDEO_TYPE_RAW = os.getenv("PIAOQUANTV_VIDEO_TYPE", "").strip()
- PIAOQUANTV_VIDEO_TYPE = int(_PIAOQUANTV_VIDEO_TYPE_RAW) if _PIAOQUANTV_VIDEO_TYPE_RAW else None
- PIAOQUANTV_HOT_FALLBACK_ENABLED = os.getenv(
- "PIAOQUANTV_HOT_FALLBACK_ENABLED", "true"
- ).strip().lower() not in {"0", "false", "no", "off"}
- PIAOQUANTV_HOT_FALLBACK_SOURCE = os.getenv("PIAOQUANTV_HOT_FALLBACK_SOURCE", "hot")
- PIAOQUANTV_VIDEO_MAX_PAGES = int(os.getenv("PIAOQUANTV_VIDEO_MAX_PAGES", "3"))
- def _load_video_crowd_package_map() -> dict[str, str]:
- """从环境变量 VIDEO_RECALL_CROWD_PACKAGE_MAP 加载人群包映射,非法或未配置时用默认映射。"""
- raw = os.getenv("VIDEO_RECALL_CROWD_PACKAGE_MAP", "").strip()
- if raw:
- try:
- data = json.loads(raw)
- if isinstance(data, dict):
- return {
- str(k).strip(): str(v).strip()
- for k, v in data.items()
- if str(k).strip() and str(v).strip()
- }
- except json.JSONDecodeError:
- logger.warning("[video_recall] VIDEO_RECALL_CROWD_PACKAGE_MAP 不是合法 JSON,使用默认映射")
- return {
- "cell*year*商业": "wx*商业",
- "回流330以上人群": "R_330+",
- }
- VIDEO_RECALL_CROWD_PACKAGE_MAP = _load_video_crowd_package_map()
- def _load_video_source_map() -> dict[str, str]:
- raw = os.getenv("VIDEO_RECALL_SOURCE_MAP", "").strip()
- if raw:
- try:
- data = json.loads(raw)
- if isinstance(data, dict):
- return {
- str(k).strip(): str(v).strip()
- for k, v in data.items()
- if str(k).strip()
- }
- except json.JSONDecodeError:
- logger.warning("[video_recall] VIDEO_RECALL_SOURCE_MAP 不是合法 JSON,使用默认映射")
- return {}
- VIDEO_RECALL_SOURCE_MAP = _load_video_source_map()
- def map_crowd_package_for_video_recall(crowd_package: str) -> str:
- """只映射内容服务 videoContentList 的 crowdPackage,不影响腾讯投放人群包。"""
- return VIDEO_RECALL_CROWD_PACKAGE_MAP.get(crowd_package, crowd_package)
- def map_source_for_video_recall(
- crowd_package: str,
- video_crowd_package: str,
- source: str,
- ) -> str:
- """只映射内容服务 videoContentList 的 source,不影响 hot 兜底语义。"""
- if source == PIAOQUANTV_HOT_FALLBACK_SOURCE:
- return source
- if crowd_package in VIDEO_RECALL_SOURCE_MAP:
- return VIDEO_RECALL_SOURCE_MAP[crowd_package]
- if video_crowd_package in VIDEO_RECALL_SOURCE_MAP:
- return VIDEO_RECALL_SOURCE_MAP[video_crowd_package]
- return source
- def _normalize_category(raw) -> str:
- """把 category 字段规范化成字符串(list 用逗号拼接,None 返回空串)。"""
- if raw is None:
- return ""
- if isinstance(raw, list):
- return ",".join(str(v).strip() for v in raw if str(v).strip())
- return str(raw).strip()
- def _enrich_hot_video_elements(videos: List["LandingVideo"]) -> None:
- """Fill pointType / standardElement for hot videos from ODPS latest partition."""
- if not videos:
- return
- try:
- from tools.video_feature_query import fetch_video_element_features
- except Exception as e:
- logger.warning("[video_recall] hot 元素补全模块导入失败,跳过:%s", e)
- return
- features_by_vid = fetch_video_element_features(v.video_id for v in videos)
- enriched = 0
- for v in videos:
- features = features_by_vid.get(v.video_id) or []
- if not features:
- continue
- feature_payload = [
- {
- "element_dimension": f.element_dimension,
- "point_type": f.point_type,
- "standard_element": f.standard_element,
- "contribution_score": f.contribution_score,
- "dt": f.dt,
- }
- for f in features
- ]
- v.raw["element_features"] = feature_payload
- # Keep the top contribution element on first-class fields so existing
- # candidate checks continue to work.
- v.point_type = v.point_type or features[0].point_type
- v.standard_element = v.standard_element or features[0].standard_element
- enriched += 1
- logger.info(
- "[video_recall] hot 元素补全完成 videos=%d/%d element_rows=%d",
- enriched, len(videos), sum(len(v) for v in features_by_vid.values()),
- )
- @dataclass
- class LandingVideo:
- """承接视频(landing video)— 用户进小程序后看到的视频。
- 注意:`video_id` 是 piaoquantv 业务侧 ID(8 位),**不**是腾讯素材库 video_id(11 位)。
- 挂创意时腾讯 creative_components.video 需要的是腾讯 video_id,二者映射逻辑待补。
- """
- video_id: int
- title: str
- cover_url: str
- video_url: str
- # 推荐分数(接口已用 rov 降序)
- score: float
- rov: float
- sim: float
- visit_uv: int
- # 召回素材时用的语义维度(对应 material_recall 三策略)
- category: str # 内容品类,例:"早中晚好";用于内容侧过滤/多样性
- standard_element: str # 例:"煽动性"
- category_name: str # 例:"情感强度控制"
- demand_content_title: str # 关联原视频标题
- demand_content_topic: str # 关联原视频选题(对应 VIDEO_TOPIC 召回)
- demand_content_id: str # 关联原视频 ID
- demand_type: str # 需求类型
- point_type: str # 点类型("关键点"/"灵感点"/"目的点"),决定 standardElement 召回的 configCode
- dimension: str # 例:"传播的头部"
- experiment_id: str # 实验 ID(透传给 xcx/save,接口会嵌入 pageUrl)
- # 原始 raw(以备未来用其他字段)
- raw: dict = field(default_factory=dict, repr=False)
- def fetch_landing_videos(
- crowd_package: str,
- page_size: int = 10,
- page_num: int = 1,
- video_business_type: Optional[int] = PIAOQUANTV_VIDEO_TYPE,
- source: str = PIAOQUANTV_VIDEO_SOURCE,
- title_filter: str = "",
- sort: int = 0,
- timeout: int = 30,
- ) -> List[LandingVideo]:
- """从 piaoquantv 拉承接视频列表(按 ROV 降序,直接取 top N)。
- Args:
- crowd_package: 人群包名,例:"泛人群" / "回流330以上人群"
- page_size: 取前 K 条
- page_num: 分页(默认 1)
- video_business_type: 5=小程序投流(默认,SOP 不变)
- source: 投放人群需求类型(默认 "prior")
- title_filter: 标题过滤(空则不过滤)
- """
- body = {
- "title": title_filter,
- "sort": sort,
- "pageNum": page_num,
- "pageSize": page_size,
- "crowdPackage": crowd_package,
- }
- if video_business_type is not None:
- body["type"] = video_business_type
- if source is not None:
- body["source"] = source
- headers = {
- "content-type": "application/json",
- "token": PIAOQUANTV_TOKEN,
- "accept": "application/json",
- }
- logger.info(
- "[video_recall] fetch crowd=%r source=%r page=%d size=%d",
- crowd_package, source, page_num, page_size,
- )
- resp = httpx.post(PIAOQUANTV_VIDEO_API, json=body, headers=headers, timeout=timeout)
- resp.raise_for_status()
- data = resp.json()
- if data.get("code") != 0 or not data.get("success"):
- raise RuntimeError(
- f"piaoquantv 接口失败:code={data.get('code')} msg={data.get('msg')}"
- )
- payload = data.get("data") or {}
- objs = payload.get("objs") or []
- videos: List[LandingVideo] = []
- for obj in objs:
- if not isinstance(obj, dict) or "videoId" not in obj:
- logger.warning("[video_recall] 跳过缺 videoId 的项:%s", obj)
- continue
- videos.append(
- LandingVideo(
- video_id=int(obj["videoId"]),
- title=obj.get("title") or "",
- cover_url=obj.get("cover") or "",
- video_url=obj.get("video") or "",
- score=float(obj.get("score") or 0.0),
- rov=float(obj.get("rov") or 0.0),
- sim=float(obj.get("sim") or 0.0),
- visit_uv=int(obj.get("visitUv") or 0),
- category=_normalize_category(obj.get("category")),
- standard_element=obj.get("standardElement") or "",
- category_name=obj.get("categoryName") or "",
- demand_content_title=obj.get("demandContentTitle") or "",
- demand_content_topic=obj.get("demandContentTopic") or "",
- demand_content_id=obj.get("demandContentId") or "",
- demand_type=obj.get("demandType") or "",
- point_type=obj.get("pointType") or "",
- dimension=obj.get("dimension") or "",
- experiment_id=obj.get("experimentId") or "",
- raw=obj,
- )
- )
- logger.info(
- "[video_recall] 返回 %d 条承接视频(总池 %d / 当前页 %d)",
- len(videos), payload.get("totalSize") or 0, payload.get("currentPage") or 1,
- )
- if source == PIAOQUANTV_HOT_FALLBACK_SOURCE:
- _enrich_hot_video_elements(videos)
- return videos
- def _fetch_landing_video_pages(
- *,
- crowd_package: str,
- page_size: int,
- source: str,
- max_pages: int,
- ) -> List[LandingVideo]:
- max_pages = max(1, int(max_pages or 1))
- merged: List[LandingVideo] = []
- seen: set[int] = set()
- for page_num in range(1, max_pages + 1):
- page = fetch_landing_videos(
- crowd_package=crowd_package,
- page_size=page_size,
- page_num=page_num,
- source=source,
- )
- if not page:
- break
- for video in page:
- if video.video_id in seen:
- continue
- merged.append(video)
- seen.add(video.video_id)
- if len(page) < page_size:
- break
- logger.info(
- "[video_recall] 分页合并 crowd=%r source=%r pages<=%d merged=%d",
- crowd_package, source, max_pages, len(merged),
- )
- return merged
- def get_account_crowd_package(account_id: int) -> str:
- """从 account_whitelist 读账户级 crowd_package。
- crowd_package 决定:
- - videoContentList 拉哪批承接视频(本模块)
- - xcx/save 注册落地计划的 audiencePackage(landing_plan 模块)
- """
- from db.connection import get_connection
- conn = get_connection()
- try:
- with conn.cursor() as cur:
- cur.execute(
- "SELECT crowd_package FROM account_whitelist WHERE account_id=%s",
- (account_id,),
- )
- row = cur.fetchone()
- finally:
- conn.close()
- if not row or not row.get("crowd_package"):
- raise ValueError(
- f"account_id {account_id} 的 crowd_package 未在 account_whitelist 配置"
- )
- return row["crowd_package"]
- def fetch_landing_videos_for_account(
- account_id: int,
- page_size: int = 10,
- source: Optional[str] = None,
- enable_hot_fallback: bool = True,
- max_pages: int = PIAOQUANTV_VIDEO_MAX_PAGES,
- ) -> List[LandingVideo]:
- """根据账户的 crowd_package 字段拉视频。
- 主池不足时用同一个 crowdPackage + source=hot 补充热门视频。
- """
- crowd_package = get_account_crowd_package(account_id)
- video_crowd_package = map_crowd_package_for_video_recall(crowd_package)
- requested_source = PIAOQUANTV_VIDEO_SOURCE if source is None else source
- selected_source = map_source_for_video_recall(
- crowd_package,
- video_crowd_package,
- requested_source,
- )
- primary = _fetch_landing_video_pages(
- crowd_package=video_crowd_package,
- page_size=page_size,
- source=selected_source,
- max_pages=max_pages,
- )
- if (
- source is not None
- or not enable_hot_fallback
- or not PIAOQUANTV_HOT_FALLBACK_ENABLED
- or len(primary) >= page_size
- or selected_source == PIAOQUANTV_HOT_FALLBACK_SOURCE
- ):
- return primary
- missing = page_size - len(primary)
- logger.info(
- "[video_recall] primary 不足,用 hot 兜底: account=%d crowd=%r video_crowd=%r primary=%d need=%d",
- account_id, crowd_package, video_crowd_package, len(primary), missing,
- )
- fallback = _fetch_landing_video_pages(
- crowd_package=video_crowd_package,
- page_size=missing,
- source=PIAOQUANTV_HOT_FALLBACK_SOURCE,
- max_pages=max_pages,
- )
- seen = {v.video_id for v in primary}
- merged = list(primary)
- for v in fallback:
- if v.video_id in seen:
- continue
- merged.append(v)
- seen.add(v.video_id)
- logger.info(
- "[video_recall] primary + hot 合并后 %d 条(account=%d crowd=%r video_crowd=%r)",
- len(merged), account_id, crowd_package, video_crowd_package,
- )
- return merged
|