"""承接视频获取接口适配层(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