video_recall.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401
  1. """承接视频获取接口适配层(piaoquantv API)。
  2. 业务模型(用户 2026-06-08 确认):
  3. - "承接视频" = 用户点开创意/进入小程序后看到的视频内容
  4. - 接口已按 ROV 降序 + 访问量过滤,我们拿前 N 即可,**无需自己排序**
  5. - 这一步**不是**创意素材召回 — 创意素材在后续 material_recall 步用本视频字段
  6. (title / standardElement / categoryName / demandContentTitle 等)做 query 召回
  7. 接口:POST https://tp-open.piaoquantv.com/contentPlatform/plan/videoContentList
  8. 认证:token header(沿用用户给的值,可通过 env PIAOQUANTV_TOKEN 覆盖)
  9. 字段参考(2026-06-08 反推自 sample 调用):
  10. 返回 data.objs[*] 39 个字段,关键:
  11. videoId : 业务侧 8 位 video ID(注:非腾讯素材库 video_id,11 位)
  12. title / cover / video : 标题 + 封面 URL + 视频 URL
  13. score / rov / sim / visitUv : 推荐分(已排序)
  14. standardElement / categoryName / demandContentTitle / dimension : 语义维度
  15. """
  16. import logging
  17. import os
  18. import json
  19. from dataclasses import dataclass, field
  20. from pathlib import Path
  21. from typing import List, Optional
  22. import httpx
  23. try:
  24. from dotenv import load_dotenv
  25. load_dotenv(Path(__file__).resolve().parents[1] / ".env")
  26. except Exception:
  27. pass
  28. logger = logging.getLogger(__name__)
  29. PIAOQUANTV_VIDEO_API = "https://tp-open.piaoquantv.com/contentPlatform/plan/videoContentList"
  30. PIAOQUANTV_TOKEN = os.getenv(
  31. "PIAOQUANTV_TOKEN",
  32. "f6e07dba7fe3476cb31fd3733d607c5c", # 用户 2026-06-08 提供
  33. )
  34. PIAOQUANTV_VIDEO_SOURCE = os.getenv("PIAOQUANTV_VIDEO_SOURCE", "")
  35. _PIAOQUANTV_VIDEO_TYPE_RAW = os.getenv("PIAOQUANTV_VIDEO_TYPE", "").strip()
  36. PIAOQUANTV_VIDEO_TYPE = int(_PIAOQUANTV_VIDEO_TYPE_RAW) if _PIAOQUANTV_VIDEO_TYPE_RAW else None
  37. PIAOQUANTV_HOT_FALLBACK_ENABLED = os.getenv(
  38. "PIAOQUANTV_HOT_FALLBACK_ENABLED", "true"
  39. ).strip().lower() not in {"0", "false", "no", "off"}
  40. PIAOQUANTV_HOT_FALLBACK_SOURCE = os.getenv("PIAOQUANTV_HOT_FALLBACK_SOURCE", "hot")
  41. PIAOQUANTV_VIDEO_MAX_PAGES = int(os.getenv("PIAOQUANTV_VIDEO_MAX_PAGES", "3"))
  42. def _load_video_crowd_package_map() -> dict[str, str]:
  43. """从环境变量 VIDEO_RECALL_CROWD_PACKAGE_MAP 加载人群包映射,非法或未配置时用默认映射。"""
  44. raw = os.getenv("VIDEO_RECALL_CROWD_PACKAGE_MAP", "").strip()
  45. if raw:
  46. try:
  47. data = json.loads(raw)
  48. if isinstance(data, dict):
  49. return {
  50. str(k).strip(): str(v).strip()
  51. for k, v in data.items()
  52. if str(k).strip() and str(v).strip()
  53. }
  54. except json.JSONDecodeError:
  55. logger.warning("[video_recall] VIDEO_RECALL_CROWD_PACKAGE_MAP 不是合法 JSON,使用默认映射")
  56. return {
  57. "cell*year*商业": "wx*商业",
  58. "回流330以上人群": "R_330+",
  59. }
  60. VIDEO_RECALL_CROWD_PACKAGE_MAP = _load_video_crowd_package_map()
  61. def _load_video_source_map() -> dict[str, str]:
  62. raw = os.getenv("VIDEO_RECALL_SOURCE_MAP", "").strip()
  63. if raw:
  64. try:
  65. data = json.loads(raw)
  66. if isinstance(data, dict):
  67. return {
  68. str(k).strip(): str(v).strip()
  69. for k, v in data.items()
  70. if str(k).strip()
  71. }
  72. except json.JSONDecodeError:
  73. logger.warning("[video_recall] VIDEO_RECALL_SOURCE_MAP 不是合法 JSON,使用默认映射")
  74. return {}
  75. VIDEO_RECALL_SOURCE_MAP = _load_video_source_map()
  76. def map_crowd_package_for_video_recall(crowd_package: str) -> str:
  77. """只映射内容服务 videoContentList 的 crowdPackage,不影响腾讯投放人群包。"""
  78. return VIDEO_RECALL_CROWD_PACKAGE_MAP.get(crowd_package, crowd_package)
  79. def map_source_for_video_recall(
  80. crowd_package: str,
  81. video_crowd_package: str,
  82. source: str,
  83. ) -> str:
  84. """只映射内容服务 videoContentList 的 source,不影响 hot 兜底语义。"""
  85. if source == PIAOQUANTV_HOT_FALLBACK_SOURCE:
  86. return source
  87. if crowd_package in VIDEO_RECALL_SOURCE_MAP:
  88. return VIDEO_RECALL_SOURCE_MAP[crowd_package]
  89. if video_crowd_package in VIDEO_RECALL_SOURCE_MAP:
  90. return VIDEO_RECALL_SOURCE_MAP[video_crowd_package]
  91. return source
  92. def _normalize_category(raw) -> str:
  93. """把 category 字段规范化成字符串(list 用逗号拼接,None 返回空串)。"""
  94. if raw is None:
  95. return ""
  96. if isinstance(raw, list):
  97. return ",".join(str(v).strip() for v in raw if str(v).strip())
  98. return str(raw).strip()
  99. def _enrich_hot_video_elements(videos: List["LandingVideo"]) -> None:
  100. """Fill pointType / standardElement for hot videos from ODPS latest partition."""
  101. if not videos:
  102. return
  103. try:
  104. from tools.video_feature_query import fetch_video_element_features
  105. except Exception as e:
  106. logger.warning("[video_recall] hot 元素补全模块导入失败,跳过:%s", e)
  107. return
  108. features_by_vid = fetch_video_element_features(v.video_id for v in videos)
  109. enriched = 0
  110. for v in videos:
  111. features = features_by_vid.get(v.video_id) or []
  112. if not features:
  113. continue
  114. feature_payload = [
  115. {
  116. "element_dimension": f.element_dimension,
  117. "point_type": f.point_type,
  118. "standard_element": f.standard_element,
  119. "contribution_score": f.contribution_score,
  120. "dt": f.dt,
  121. }
  122. for f in features
  123. ]
  124. v.raw["element_features"] = feature_payload
  125. # Keep the top contribution element on first-class fields so existing
  126. # candidate checks continue to work.
  127. v.point_type = v.point_type or features[0].point_type
  128. v.standard_element = v.standard_element or features[0].standard_element
  129. enriched += 1
  130. logger.info(
  131. "[video_recall] hot 元素补全完成 videos=%d/%d element_rows=%d",
  132. enriched, len(videos), sum(len(v) for v in features_by_vid.values()),
  133. )
  134. @dataclass
  135. class LandingVideo:
  136. """承接视频(landing video)— 用户进小程序后看到的视频。
  137. 注意:`video_id` 是 piaoquantv 业务侧 ID(8 位),**不**是腾讯素材库 video_id(11 位)。
  138. 挂创意时腾讯 creative_components.video 需要的是腾讯 video_id,二者映射逻辑待补。
  139. """
  140. video_id: int
  141. title: str
  142. cover_url: str
  143. video_url: str
  144. # 推荐分数(接口已用 rov 降序)
  145. score: float
  146. rov: float
  147. sim: float
  148. visit_uv: int
  149. # 召回素材时用的语义维度(对应 material_recall 三策略)
  150. category: str # 内容品类,例:"早中晚好";用于内容侧过滤/多样性
  151. standard_element: str # 例:"煽动性"
  152. category_name: str # 例:"情感强度控制"
  153. demand_content_title: str # 关联原视频标题
  154. demand_content_topic: str # 关联原视频选题(对应 VIDEO_TOPIC 召回)
  155. demand_content_id: str # 关联原视频 ID
  156. demand_type: str # 需求类型
  157. point_type: str # 点类型("关键点"/"灵感点"/"目的点"),决定 standardElement 召回的 configCode
  158. dimension: str # 例:"传播的头部"
  159. experiment_id: str # 实验 ID(透传给 xcx/save,接口会嵌入 pageUrl)
  160. # 原始 raw(以备未来用其他字段)
  161. raw: dict = field(default_factory=dict, repr=False)
  162. def fetch_landing_videos(
  163. crowd_package: str,
  164. page_size: int = 10,
  165. page_num: int = 1,
  166. video_business_type: Optional[int] = PIAOQUANTV_VIDEO_TYPE,
  167. source: str = PIAOQUANTV_VIDEO_SOURCE,
  168. title_filter: str = "",
  169. sort: int = 0,
  170. timeout: int = 30,
  171. ) -> List[LandingVideo]:
  172. """从 piaoquantv 拉承接视频列表(按 ROV 降序,直接取 top N)。
  173. Args:
  174. crowd_package: 人群包名,例:"泛人群" / "回流330以上人群"
  175. page_size: 取前 K 条
  176. page_num: 分页(默认 1)
  177. video_business_type: 5=小程序投流(默认,SOP 不变)
  178. source: 投放人群需求类型(默认 "prior")
  179. title_filter: 标题过滤(空则不过滤)
  180. """
  181. body = {
  182. "title": title_filter,
  183. "sort": sort,
  184. "pageNum": page_num,
  185. "pageSize": page_size,
  186. "crowdPackage": crowd_package,
  187. }
  188. if video_business_type is not None:
  189. body["type"] = video_business_type
  190. if source is not None:
  191. body["source"] = source
  192. headers = {
  193. "content-type": "application/json",
  194. "token": PIAOQUANTV_TOKEN,
  195. "accept": "application/json",
  196. }
  197. logger.info(
  198. "[video_recall] fetch crowd=%r source=%r page=%d size=%d",
  199. crowd_package, source, page_num, page_size,
  200. )
  201. resp = httpx.post(PIAOQUANTV_VIDEO_API, json=body, headers=headers, timeout=timeout)
  202. resp.raise_for_status()
  203. data = resp.json()
  204. if data.get("code") != 0 or not data.get("success"):
  205. raise RuntimeError(
  206. f"piaoquantv 接口失败:code={data.get('code')} msg={data.get('msg')}"
  207. )
  208. payload = data.get("data") or {}
  209. objs = payload.get("objs") or []
  210. videos: List[LandingVideo] = []
  211. for obj in objs:
  212. if not isinstance(obj, dict) or "videoId" not in obj:
  213. logger.warning("[video_recall] 跳过缺 videoId 的项:%s", obj)
  214. continue
  215. videos.append(
  216. LandingVideo(
  217. video_id=int(obj["videoId"]),
  218. title=obj.get("title") or "",
  219. cover_url=obj.get("cover") or "",
  220. video_url=obj.get("video") or "",
  221. score=float(obj.get("score") or 0.0),
  222. rov=float(obj.get("rov") or 0.0),
  223. sim=float(obj.get("sim") or 0.0),
  224. visit_uv=int(obj.get("visitUv") or 0),
  225. category=_normalize_category(obj.get("category")),
  226. standard_element=obj.get("standardElement") or "",
  227. category_name=obj.get("categoryName") or "",
  228. demand_content_title=obj.get("demandContentTitle") or "",
  229. demand_content_topic=obj.get("demandContentTopic") or "",
  230. demand_content_id=obj.get("demandContentId") or "",
  231. demand_type=obj.get("demandType") or "",
  232. point_type=obj.get("pointType") or "",
  233. dimension=obj.get("dimension") or "",
  234. experiment_id=obj.get("experimentId") or "",
  235. raw=obj,
  236. )
  237. )
  238. logger.info(
  239. "[video_recall] 返回 %d 条承接视频(总池 %d / 当前页 %d)",
  240. len(videos), payload.get("totalSize") or 0, payload.get("currentPage") or 1,
  241. )
  242. if source == PIAOQUANTV_HOT_FALLBACK_SOURCE:
  243. _enrich_hot_video_elements(videos)
  244. return videos
  245. def _fetch_landing_video_pages(
  246. *,
  247. crowd_package: str,
  248. page_size: int,
  249. source: str,
  250. max_pages: int,
  251. ) -> List[LandingVideo]:
  252. max_pages = max(1, int(max_pages or 1))
  253. merged: List[LandingVideo] = []
  254. seen: set[int] = set()
  255. for page_num in range(1, max_pages + 1):
  256. page = fetch_landing_videos(
  257. crowd_package=crowd_package,
  258. page_size=page_size,
  259. page_num=page_num,
  260. source=source,
  261. )
  262. if not page:
  263. break
  264. for video in page:
  265. if video.video_id in seen:
  266. continue
  267. merged.append(video)
  268. seen.add(video.video_id)
  269. if len(page) < page_size:
  270. break
  271. logger.info(
  272. "[video_recall] 分页合并 crowd=%r source=%r pages<=%d merged=%d",
  273. crowd_package, source, max_pages, len(merged),
  274. )
  275. return merged
  276. def get_account_crowd_package(account_id: int) -> str:
  277. """从 account_whitelist 读账户级 crowd_package。
  278. crowd_package 决定:
  279. - videoContentList 拉哪批承接视频(本模块)
  280. - xcx/save 注册落地计划的 audiencePackage(landing_plan 模块)
  281. """
  282. from db.connection import get_connection
  283. conn = get_connection()
  284. try:
  285. with conn.cursor() as cur:
  286. cur.execute(
  287. "SELECT crowd_package FROM account_whitelist WHERE account_id=%s",
  288. (account_id,),
  289. )
  290. row = cur.fetchone()
  291. finally:
  292. conn.close()
  293. if not row or not row.get("crowd_package"):
  294. raise ValueError(
  295. f"account_id {account_id} 的 crowd_package 未在 account_whitelist 配置"
  296. )
  297. return row["crowd_package"]
  298. def fetch_landing_videos_for_account(
  299. account_id: int,
  300. page_size: int = 10,
  301. source: Optional[str] = None,
  302. enable_hot_fallback: bool = True,
  303. max_pages: int = PIAOQUANTV_VIDEO_MAX_PAGES,
  304. ) -> List[LandingVideo]:
  305. """根据账户的 crowd_package 字段拉视频。
  306. 主池不足时用同一个 crowdPackage + source=hot 补充热门视频。
  307. """
  308. crowd_package = get_account_crowd_package(account_id)
  309. video_crowd_package = map_crowd_package_for_video_recall(crowd_package)
  310. requested_source = PIAOQUANTV_VIDEO_SOURCE if source is None else source
  311. selected_source = map_source_for_video_recall(
  312. crowd_package,
  313. video_crowd_package,
  314. requested_source,
  315. )
  316. primary = _fetch_landing_video_pages(
  317. crowd_package=video_crowd_package,
  318. page_size=page_size,
  319. source=selected_source,
  320. max_pages=max_pages,
  321. )
  322. if (
  323. source is not None
  324. or not enable_hot_fallback
  325. or not PIAOQUANTV_HOT_FALLBACK_ENABLED
  326. or len(primary) >= page_size
  327. or selected_source == PIAOQUANTV_HOT_FALLBACK_SOURCE
  328. ):
  329. return primary
  330. missing = page_size - len(primary)
  331. logger.info(
  332. "[video_recall] primary 不足,用 hot 兜底: account=%d crowd=%r video_crowd=%r primary=%d need=%d",
  333. account_id, crowd_package, video_crowd_package, len(primary), missing,
  334. )
  335. fallback = _fetch_landing_video_pages(
  336. crowd_package=video_crowd_package,
  337. page_size=missing,
  338. source=PIAOQUANTV_HOT_FALLBACK_SOURCE,
  339. max_pages=max_pages,
  340. )
  341. seen = {v.video_id for v in primary}
  342. merged = list(primary)
  343. for v in fallback:
  344. if v.video_id in seen:
  345. continue
  346. merged.append(v)
  347. seen.add(v.video_id)
  348. logger.info(
  349. "[video_recall] primary + hot 合并后 %d 条(account=%d crowd=%r video_crowd=%r)",
  350. len(merged), account_id, crowd_package, video_crowd_package,
  351. )
  352. return merged