material_recall.py 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694
  1. """创意素材召回接口适配层(vector 服务)。
  2. 业务模型(用户 2026-06-08 确认):
  3. 对一个承接视频(LandingVideo)执行 3 种策略召回素材,合并去重,按 score 取 top N。
  4. 接口:POST https://api-internal.piaoquantv.com/videoVector/recallTest/matchByText
  5. 鉴权:不需要(内部 API)
  6. 本质(用户 2026-06-08 确认):**多路召回 → 排序 → 去重**
  7. 实测确认素材库的 4 个有效 configCode(VIDEO_TITLE 和 ALL 对 MATERIAL 模态返回 0):
  8. - VIDEO_TOPIC (选题向量)
  9. - VIDEO_KEYPOINT (关键点向量)
  10. - VIDEO_INSPIRATION(灵感点向量)
  11. - VIDEO_PURPOSE (目的点向量)
  12. 策略:对一个承接视频,选合适的 queryText,**串行调 4 个 configCode**,合并按 materialId 去重(取 max score),按 score 降序返回 top N。
  13. queryText 选择优先级(选第一个非空且非占位符 "-"):
  14. 1. standard_element(标准化元素,实测效果最好)
  15. 2. demand_content_topic(选题,如果非 "-")
  16. 3. title(承接视频标题,兜底)
  17. """
  18. import logging
  19. import os
  20. from concurrent.futures import ThreadPoolExecutor, as_completed
  21. from dataclasses import dataclass, field
  22. from typing import Iterable, List, Optional
  23. import httpx
  24. from tools.video_recall import LandingVideo
  25. logger = logging.getLogger(__name__)
  26. VECTOR_BASE = os.getenv(
  27. "VECTOR_BASE_URL",
  28. "https://api-internal.piaoquantv.com/videoVector",
  29. )
  30. VECTOR_MATCH_BY_TEXT = f"{VECTOR_BASE}/recallTest/matchByText"
  31. VECTOR_BATCH_BY_TEXT = f"{VECTOR_BASE}/recallTest/batchByText"
  32. VECTOR_ALL_CONFIG_CODES = f"{VECTOR_BASE}/videoSearch/getAllConfigCodes"
  33. def fetch_all_config_codes() -> dict[str, str]:
  34. """动态获取 vector 服务支持的全部 configCode → 中文名。
  35. 返回 {configCode: 中文名}。
  36. 召回前先调一次,避免硬编码列表过时。
  37. """
  38. resp = httpx.get(VECTOR_ALL_CONFIG_CODES, timeout=15)
  39. resp.raise_for_status()
  40. data = resp.json()
  41. if data.get("code") not in (0, 200):
  42. raise RuntimeError(
  43. f"getAllConfigCodes 失败:code={data.get('code')} msg={data.get('msg')}"
  44. )
  45. return data.get("data") or {}
  46. # 字段维度策略(2026-06-08 用户最终确认)
  47. # - 维度 1 选题: queryText=demand_content_topic, configCode=VIDEO_TOPIC
  48. # - 维度 2 标准化元素: queryText=standard_element, configCode 按 point_type **唯一对应**:
  49. # - "灵感点" → INSPIRATION_SUBSTANCE
  50. # - "关键点" → KEYPOINT_SUBSTANCE
  51. # - "目的点" → PURPOSE_SUBSTANCE
  52. # - 其他/空 → 跳过(没办法确定走哪一路,no-guessing)
  53. # - 维度 3 标题: 不做(VIDEO_TITLE 对 MATERIAL 实测返回 0)
  54. POINT_TYPE_TO_SUBSTANCE = {
  55. "灵感点": "INSPIRATION_SUBSTANCE",
  56. "关键点": "KEYPOINT_SUBSTANCE",
  57. "目的点": "PURPOSE_SUBSTANCE",
  58. }
  59. POINT_TYPE_TO_VIDEO_CONFIG = {
  60. "灵感点": "VIDEO_INSPIRATION",
  61. "关键点": "VIDEO_KEYPOINT",
  62. "目的点": "VIDEO_PURPOSE",
  63. }
  64. SUBSTANCE_TO_VIDEO_CONFIG = {
  65. "INSPIRATION_SUBSTANCE": "VIDEO_INSPIRATION",
  66. "KEYPOINT_SUBSTANCE": "VIDEO_KEYPOINT",
  67. "PURPOSE_SUBSTANCE": "VIDEO_PURPOSE",
  68. }
  69. ODPS_FEATURE_TO_CONFIG = {
  70. ("解构选题", ""): "VIDEO_TOPIC",
  71. ("实质", "灵感点"): "INSPIRATION_SUBSTANCE",
  72. ("实质", "关键点"): "KEYPOINT_SUBSTANCE",
  73. ("实质", "目的点"): "PURPOSE_SUBSTANCE",
  74. }
  75. # fallback(动态接口失败时用)
  76. MATERIAL_EFFECTIVE_CONFIG_CODES = ["VIDEO_TOPIC"] + list(POINT_TYPE_TO_SUBSTANCE.values())
  77. # 占位符值视为无效(piaoquantv 数据里 demandContentTopic 大量为 "-")
  78. PLACEHOLDER_VALUES = {"", "-", None}
  79. # 单 configCode 召回 top N(合并前)
  80. DEFAULT_PER_CC_TOP_N = 100
  81. # 最终合并去重后返回 top N
  82. DEFAULT_FINAL_TOP_N = 20
  83. # Cover URL 黑名单 — 已知尺寸/比例不符合腾讯创意要求的来源
  84. # (2026-06-09 用户决策:踩坑驱动渐进收紧,踩到新模式再加)
  85. # - "auto_reply_cards_cover" :自动回复卡片业务的 cover,实测 reject code 1801159
  86. EXCLUDED_COVER_URL_PATTERNS = (
  87. "auto_reply_cards_cover",
  88. )
  89. @dataclass
  90. class RecallQuery:
  91. """One ODPS-derived material recall query."""
  92. query_text: str
  93. config_code: str
  94. element_dimension: str
  95. point_type: str
  96. standard_element: str
  97. contribution_score: float = 0.0
  98. dt: str = ""
  99. @property
  100. def strategy_name(self) -> str:
  101. """返回召回策略名称(解构选题 或 元素维度-点类型)。"""
  102. if self.element_dimension == "解构选题":
  103. return "解构选题"
  104. return f"{self.element_dimension}-{self.point_type}"
  105. def to_hit(self) -> dict:
  106. """把召回 query 转成命中记录 dict(写入素材的 recall_hit_queries)。"""
  107. return {
  108. "strategy": self.strategy_name,
  109. "element_dimension": self.element_dimension,
  110. "point_type": self.point_type,
  111. "standard_element": self.standard_element,
  112. "query_text": self.query_text,
  113. "config_code": self.config_code,
  114. "contribution_score": self.contribution_score,
  115. "dt": self.dt,
  116. }
  117. @dataclass
  118. class Material:
  119. """召回的创意素材。
  120. 2026-06-10 升级:从 batchByText 接口接收质量数据(cost/ctr/cvr/roi/impressions 等)。
  121. """
  122. material_id: str # 素材原始 ID
  123. score: float # 服务端综合分(wCtr=1 配置下 ≈ qualityScore,与 CTR 强相关)
  124. title: str = ""
  125. cover: str = ""
  126. video_url: str = ""
  127. # 投放质量数据(来自 materialDetail.quality)— 2026-06-10 升级新增
  128. cost: Optional[float] = None # 累计成本(元)
  129. ctr: Optional[float] = None # 点击率
  130. cvr: Optional[float] = None # 转化率
  131. roi: Optional[float] = None # ROI(收入/成本)
  132. impressions: Optional[int] = None # 累计曝光
  133. quality_score: Optional[float] = None # 服务端 qualityScore
  134. recall_strategy: str = ""
  135. recall_query_text: str = ""
  136. recall_config_code: str = ""
  137. recall_element_dimension: str = ""
  138. recall_point_type: str = ""
  139. recall_standard_element: str = ""
  140. recall_hit_queries: List[dict] = field(default_factory=list)
  141. # 原始 item dict(以备后用)
  142. raw: dict = field(default_factory=dict, repr=False)
  143. def _call_match_by_text(
  144. query_text: str,
  145. config_code: str,
  146. material_top_n: int = DEFAULT_PER_CC_TOP_N,
  147. source_labels: Optional[List[str]] = None,
  148. timeout: int = 30,
  149. ) -> List[dict]:
  150. """单次调 matchByText,只要 MATERIAL 模态,返回原始 items 列表。"""
  151. body = {
  152. "queryText": query_text,
  153. "configCode": config_code,
  154. "modalities": ["MATERIAL"],
  155. "videoTopN": 0,
  156. "articleTopN": 0,
  157. "materialTopN": material_top_n,
  158. "topN": material_top_n,
  159. "displayK": material_top_n,
  160. }
  161. if source_labels:
  162. body["sourceLabels"] = source_labels
  163. logger.info(
  164. "[material_recall] matchByText q=%r configCode=%s topN=%d",
  165. query_text[:40], config_code, material_top_n,
  166. )
  167. resp = httpx.post(
  168. VECTOR_MATCH_BY_TEXT,
  169. json=body,
  170. headers={"content-type": "application/json", "accept": "application/json"},
  171. timeout=timeout,
  172. )
  173. resp.raise_for_status()
  174. data = resp.json()
  175. code = data.get("code")
  176. # CommonResponse 可能用 0 或 200 表示成功;以 data 字段为准
  177. if code not in (0, 200, "0", "200"):
  178. raise RuntimeError(
  179. f"vector matchByText 失败:code={code} msg={data.get('msg') or data.get('message')}"
  180. )
  181. payload = data.get("data") or {}
  182. items = payload.get("items") or []
  183. # 只留 MATERIAL 模态(防御性 — 万一返回混入)
  184. return [it for it in items if it.get("modality") == "MATERIAL"]
  185. def _pick_query_text(landing: LandingVideo) -> Optional[str]:
  186. """按优先级选 queryText:standard_element > demand_content_topic > title。
  187. 占位符 "-" / 空串视为无效。"""
  188. for cand in (landing.standard_element, landing.demand_content_topic, landing.title):
  189. if cand and cand not in PLACEHOLDER_VALUES:
  190. return cand
  191. return None
  192. def _pick_query_text_for_batch(landing: LandingVideo) -> Optional[str]:
  193. """选 queryText:standard_element > demand_content_topic > title。(保留兼容)"""
  194. for cand in (landing.standard_element, landing.demand_content_topic, landing.title):
  195. if cand and cand not in PLACEHOLDER_VALUES:
  196. return cand
  197. return None
  198. def _pick_query_strategies_for_batch(landing: LandingVideo) -> List[tuple]:
  199. """返回 (queryText, configCode, strategy_name) 策略候选(2026-06-10 用户最终确认)。
  200. 只用 2 个维度:
  201. 维度 1 - 标准化元素:queryText=standard_element, configCode 按 point_type 唯一对应
  202. - "灵感点" → INSPIRATION_SUBSTANCE
  203. - "关键点" → KEYPOINT_SUBSTANCE
  204. - "目的点" → PURPOSE_SUBSTANCE
  205. 维度 2 - 选题:queryText=demand_content_topic, configCode=VIDEO_TOPIC
  206. prepare 阶段:维度 1 失败 → 试维度 2 → 仍失败 → 换 landing。
  207. """
  208. strategies = []
  209. seen = set()
  210. for feature in landing.raw.get("element_features") or []:
  211. if not isinstance(feature, dict):
  212. continue
  213. standard_element = (feature.get("standard_element") or "").strip()
  214. point_type = (feature.get("point_type") or "").strip()
  215. if standard_element in PLACEHOLDER_VALUES:
  216. continue
  217. cc = POINT_TYPE_TO_SUBSTANCE.get(point_type)
  218. if not cc:
  219. continue
  220. key = (standard_element, cc)
  221. if key in seen:
  222. continue
  223. seen.add(key)
  224. strategies.append((standard_element, cc, f"标准化元素-{point_type}"))
  225. # 维度 1: 标准化元素(优先)
  226. if landing.standard_element and landing.standard_element not in PLACEHOLDER_VALUES:
  227. cc = POINT_TYPE_TO_SUBSTANCE.get(landing.point_type)
  228. if cc and (landing.standard_element, cc) not in seen:
  229. strategies.append((landing.standard_element, cc, f"标准化元素-{landing.point_type}"))
  230. # 维度 2: 选题
  231. if landing.demand_content_topic and landing.demand_content_topic not in PLACEHOLDER_VALUES:
  232. strategies.append((landing.demand_content_topic, "VIDEO_TOPIC", "选题"))
  233. return strategies
  234. def _fallback_config_codes_for_strategy(config_code: str, landing: LandingVideo) -> List[str]:
  235. """batchByText 无结果时给 matchByText 的兼容 configCode。
  236. 2026-06-30 实测:batchByText 对 MATERIAL 返回 0 时,旧 matchByText 仍可在
  237. VIDEO_KEYPOINT/VIDEO_INSPIRATION 等历史向量字段命中素材。
  238. """
  239. out = [config_code]
  240. video_cc = SUBSTANCE_TO_VIDEO_CONFIG.get(config_code) or POINT_TYPE_TO_VIDEO_CONFIG.get(landing.point_type)
  241. if video_cc and video_cc not in out:
  242. out.append(video_cc)
  243. if config_code != "VIDEO_TOPIC" and "VIDEO_TOPIC" not in out:
  244. out.append("VIDEO_TOPIC")
  245. return out
  246. def _call_batch_by_text(
  247. query_text: str,
  248. config_codes: List[str],
  249. display_k: int,
  250. days: int,
  251. sim_threshold: float,
  252. alpha: float,
  253. w_ctr: float, w_cvr: float, w_roi: float,
  254. w_open_rate: float, w_fission_rate: float,
  255. deconstruct_boost: float,
  256. source_labels: List[str],
  257. modalities: List[str] = None,
  258. timeout: int = 30,
  259. ) -> List[dict]:
  260. """调用 batchByText 接口(2026-06-10 升级).
  261. 服务端单次 embedding + 多 configCode 并行 ANN + 跨模态过滤 + ranking 加权 + 去重。
  262. """
  263. body = {
  264. "queryText": query_text,
  265. "configCodes": config_codes,
  266. "displayK": display_k,
  267. "modalities": modalities or ["MATERIAL"],
  268. "sourceLabels": source_labels,
  269. "days": days,
  270. "ranking": {
  271. "simThreshold": sim_threshold,
  272. "alpha": alpha,
  273. "wCtr": w_ctr, "wCvr": w_cvr, "wRoi": w_roi,
  274. "wOpenRate": w_open_rate, "wFissionRate": w_fission_rate,
  275. "deconstructBoost": deconstruct_boost,
  276. },
  277. }
  278. logger.info(
  279. "[material_recall] batchByText q=%r configCodes=%d displayK=%d simT=%.2f wCtr=%.2f",
  280. query_text[:40], len(config_codes), display_k, sim_threshold, w_ctr,
  281. )
  282. resp = httpx.post(
  283. VECTOR_BATCH_BY_TEXT,
  284. json=body,
  285. headers={"content-type": "application/json", "accept": "application/json"},
  286. timeout=timeout,
  287. )
  288. resp.raise_for_status()
  289. data = resp.json()
  290. if data.get("code") not in (0, 200, "0", "200"):
  291. raise RuntimeError(
  292. f"batchByText 失败:code={data.get('code')} msg={data.get('msg') or data.get('message')}"
  293. )
  294. payload = data.get("data") or {}
  295. return payload.get("items") or []
  296. def _as_float(value, default: float = 0.0) -> float:
  297. """安全转 float,None 或非法值返回 default。"""
  298. try:
  299. if value is None:
  300. return default
  301. return float(value)
  302. except (TypeError, ValueError):
  303. return default
  304. def _as_int(value, default: int = 0) -> int:
  305. """安全转 int,None 或非法值返回 default。"""
  306. try:
  307. if value is None:
  308. return default
  309. return int(float(value))
  310. except (TypeError, ValueError):
  311. return default
  312. def _items_to_materials(items: List[dict], sim_threshold: float) -> tuple:
  313. """把召回 items 过滤 + 转 Material。
  314. 过滤条件:
  315. - modality=MATERIAL(防御性)
  316. - cover URL 不在黑名单
  317. - score >= sim_threshold
  318. CTR / impressions 只作为审批展示和兜底排序参考,不再作为硬筛。
  319. 返回 (materials, stats)。
  320. """
  321. out: List[Material] = []
  322. stats = {
  323. "blacklist": 0,
  324. "low_score": 0,
  325. "low_imp": 0,
  326. "low_ctr": 0,
  327. }
  328. for it in items:
  329. if it.get("modality") != "MATERIAL":
  330. continue
  331. mid = it.get("materialId") or (str(it["id"]) if it.get("id") is not None else None)
  332. if not mid:
  333. continue
  334. cover = it.get("cover") or ""
  335. if any(p in cover for p in EXCLUDED_COVER_URL_PATTERNS):
  336. stats["blacklist"] += 1
  337. continue
  338. score = _as_float(it.get("score"))
  339. if score < sim_threshold:
  340. stats["low_score"] += 1
  341. continue
  342. md = it.get("materialDetail") or {}
  343. q = md.get("quality") or {}
  344. out.append(Material(
  345. material_id=str(mid),
  346. score=score,
  347. title=it.get("title") or "",
  348. cover=cover,
  349. video_url=it.get("videoUrl") or "",
  350. cost=_as_float(q.get("cost")) if q.get("cost") is not None else None,
  351. ctr=_as_float(q.get("ctr")) if q.get("ctr") is not None else None,
  352. cvr=_as_float(q.get("cvr")) if q.get("cvr") is not None else None,
  353. roi=_as_float(q.get("roi")) if q.get("roi") is not None else None,
  354. impressions=_as_int(q.get("impressions")) if q.get("impressions") is not None else None,
  355. quality_score=_as_float(q.get("qualityScore")) if q.get("qualityScore") is not None else None,
  356. raw=it,
  357. ))
  358. return out, stats
  359. def _sort_materials_by_policy(materials: List[Material]) -> List[Material]:
  360. """生产排序策略:先相关性准入,再按历史消耗倒序。"""
  361. return sorted(
  362. materials,
  363. key=lambda m: (
  364. m.cost is not None,
  365. m.cost or 0,
  366. m.roi or 0,
  367. m.impressions or 0,
  368. m.ctr or 0,
  369. m.quality_score or 0,
  370. m.score or 0,
  371. ),
  372. reverse=True,
  373. )
  374. def _material_rank(material: Material) -> tuple:
  375. """生成素材排序键(与 _sort_materials_by_policy 同序:消耗优先)。"""
  376. return (
  377. material.cost is not None,
  378. material.cost or 0,
  379. material.roi or 0,
  380. material.impressions or 0,
  381. material.ctr or 0,
  382. material.quality_score or 0,
  383. material.score or 0,
  384. )
  385. def _feature_attr(feature, name: str, default=""):
  386. """兼容 dict 和对象两种 feature,统一取属性值。"""
  387. if isinstance(feature, dict):
  388. return feature.get(name, default)
  389. return getattr(feature, name, default)
  390. def _build_recall_queries_from_features(
  391. element_features: Iterable,
  392. query_limit: int,
  393. ) -> List[RecallQuery]:
  394. """Build recall queries from ODPS features.
  395. Supported dimensions:
  396. - 解构选题 -> VIDEO_TOPIC
  397. - 实质 + 灵感点/关键点/目的点 -> *_SUBSTANCE
  398. """
  399. queries: List[RecallQuery] = []
  400. seen = set()
  401. raw_features = list(element_features or [])
  402. sorted_features = sorted(
  403. raw_features,
  404. key=lambda f: (
  405. 0 if str(_feature_attr(f, "element_dimension") or "") == "解构选题" else 1,
  406. -float(_feature_attr(f, "contribution_score", 0) or 0),
  407. ),
  408. )
  409. for feature in sorted_features:
  410. element_dimension = str(_feature_attr(feature, "element_dimension") or "").strip()
  411. point_type = str(_feature_attr(feature, "point_type") or "").strip()
  412. standard_element = str(_feature_attr(feature, "standard_element") or "").strip()
  413. if standard_element in PLACEHOLDER_VALUES:
  414. continue
  415. config_code = ODPS_FEATURE_TO_CONFIG.get((element_dimension, point_type))
  416. if not config_code and element_dimension == "解构选题":
  417. config_code = ODPS_FEATURE_TO_CONFIG.get(("解构选题", ""))
  418. if not config_code:
  419. continue
  420. key = (standard_element, config_code)
  421. if key in seen:
  422. continue
  423. seen.add(key)
  424. queries.append(RecallQuery(
  425. query_text=standard_element,
  426. config_code=config_code,
  427. element_dimension=element_dimension,
  428. point_type=point_type,
  429. standard_element=standard_element,
  430. contribution_score=float(_feature_attr(feature, "contribution_score", 0) or 0),
  431. dt=str(_feature_attr(feature, "dt") or ""),
  432. ))
  433. if len(queries) >= query_limit:
  434. break
  435. return queries
  436. def _call_batch_for_recall_query(
  437. query: RecallQuery,
  438. *,
  439. display_k: int,
  440. days: int,
  441. sim_threshold: float,
  442. alpha: float,
  443. w_ctr: float,
  444. w_cvr: float,
  445. w_roi: float,
  446. w_open_rate: float,
  447. w_fission_rate: float,
  448. deconstruct_boost: float,
  449. source_labels: List[str],
  450. ) -> tuple[RecallQuery, List[dict]]:
  451. """对单个召回 query 调 batchByText,返回 (query, 原始 items)。供并行调用。"""
  452. items = _call_batch_by_text(
  453. query_text=query.query_text,
  454. config_codes=[query.config_code],
  455. display_k=display_k,
  456. days=days,
  457. sim_threshold=sim_threshold,
  458. alpha=alpha,
  459. w_ctr=w_ctr,
  460. w_cvr=w_cvr,
  461. w_roi=w_roi,
  462. w_open_rate=w_open_rate,
  463. w_fission_rate=w_fission_rate,
  464. deconstruct_boost=deconstruct_boost,
  465. source_labels=source_labels,
  466. modalities=["MATERIAL"],
  467. )
  468. return query, items
  469. def _merge_materials_by_policy(query_materials: List[tuple[RecallQuery, List[Material]]]) -> List[Material]:
  470. """多路召回结果按 material_id 去重合并(保留排序更优者),并记录命中 query,最后按策略排序。"""
  471. by_mid: dict[str, Material] = {}
  472. for query, materials in query_materials:
  473. hit = query.to_hit()
  474. for material in materials:
  475. material.recall_hit_queries = [hit]
  476. material.recall_strategy = query.strategy_name
  477. material.recall_query_text = query.query_text
  478. material.recall_config_code = query.config_code
  479. material.recall_element_dimension = query.element_dimension
  480. material.recall_point_type = query.point_type
  481. material.recall_standard_element = query.standard_element
  482. existing = by_mid.get(material.material_id)
  483. if existing is None:
  484. by_mid[material.material_id] = material
  485. continue
  486. merged_hits = existing.recall_hit_queries + [
  487. h for h in material.recall_hit_queries
  488. if h not in existing.recall_hit_queries
  489. ]
  490. if _material_rank(material) > _material_rank(existing):
  491. material.recall_hit_queries = merged_hits
  492. by_mid[material.material_id] = material
  493. else:
  494. existing.recall_hit_queries = merged_hits
  495. return _sort_materials_by_policy(list(by_mid.values()))
  496. def recall_materials_for_video(
  497. landing: LandingVideo,
  498. final_top_n: int = DEFAULT_FINAL_TOP_N,
  499. source_labels: Optional[List[str]] = None,
  500. element_features: Optional[Iterable] = None,
  501. ) -> List[Material]:
  502. """素材召回:用 ODPS 多维特征并行召回并合并排序。
  503. 流程:
  504. 1. 从 ODPS features 生成 query:解构选题 + 实质三点。
  505. 2. 多 query 并行调用 batchByText。
  506. 3. 汇总、material_id 去重、score>=阈值、按 cost 倒序。
  507. 当前硬筛只保留相似度阈值;曝光/CTR 进入审批表但不拦截。
  508. """
  509. from config import (
  510. RECALL_ALPHA, RECALL_DAYS, RECALL_DECONSTRUCT_BOOST,
  511. RECALL_DISPLAY_K, RECALL_PARALLEL_MAX_WORKERS, RECALL_QUERY_LIMIT_PER_VIDEO,
  512. RECALL_SIM_THRESHOLD, RECALL_SOURCE_LABELS,
  513. RECALL_W_CTR, RECALL_W_CVR, RECALL_W_FISSION_RATE,
  514. RECALL_W_OPEN_RATE, RECALL_W_ROI,
  515. )
  516. if element_features is None:
  517. element_features = landing.raw.get("element_features") or []
  518. queries = _build_recall_queries_from_features(
  519. element_features,
  520. query_limit=max(1, RECALL_QUERY_LIMIT_PER_VIDEO),
  521. )
  522. if not queries:
  523. logger.warning(
  524. "[material_recall] landing video_id=%d 无 ODPS 可用召回特征,返回空",
  525. landing.video_id,
  526. )
  527. return []
  528. logger.info(
  529. "[material_recall] landing video_id=%d 走 %d 个 ODPS 策略:%s",
  530. landing.video_id, len(queries),
  531. "; ".join(f"{q.strategy_name}:{q.query_text}->{q.config_code}" for q in queries),
  532. )
  533. query_materials: List[tuple[RecallQuery, List[Material]]] = []
  534. labels = source_labels or RECALL_SOURCE_LABELS
  535. max_workers = max(1, min(RECALL_PARALLEL_MAX_WORKERS, len(queries)))
  536. with ThreadPoolExecutor(max_workers=max_workers) as executor:
  537. futures = [
  538. executor.submit(
  539. _call_batch_for_recall_query,
  540. query,
  541. display_k=RECALL_DISPLAY_K,
  542. days=RECALL_DAYS,
  543. sim_threshold=RECALL_SIM_THRESHOLD,
  544. alpha=RECALL_ALPHA,
  545. w_ctr=RECALL_W_CTR,
  546. w_cvr=RECALL_W_CVR,
  547. w_roi=RECALL_W_ROI,
  548. w_open_rate=RECALL_W_OPEN_RATE,
  549. w_fission_rate=RECALL_W_FISSION_RATE,
  550. deconstruct_boost=RECALL_DECONSTRUCT_BOOST,
  551. source_labels=labels,
  552. )
  553. for query in queries
  554. ]
  555. for future in as_completed(futures):
  556. try:
  557. query, items = future.result()
  558. except Exception as e:
  559. logger.error("[material_recall] 并行 batchByText 失败:%s", e)
  560. continue
  561. mats, stats = _items_to_materials(items, RECALL_SIM_THRESHOLD)
  562. query_materials.append((query, mats))
  563. logger.info(
  564. "[material_recall] 策略=%s q=%r configCode=%s 返回 %d 条,⊘ 黑名单 %d,⊘ score<%.2f %d → 保留 %d",
  565. query.strategy_name, query.query_text[:30], query.config_code,
  566. len(items), stats["blacklist"], RECALL_SIM_THRESHOLD,
  567. stats["low_score"], len(mats),
  568. )
  569. merged = _merge_materials_by_policy(query_materials)
  570. logger.info(
  571. "[material_recall] landing video_id=%d 多维召回合并后保留 %d 条(cost desc)",
  572. landing.video_id, len(merged),
  573. )
  574. if merged:
  575. return merged[:final_top_n]
  576. for query in queries:
  577. # batchByText 当前可能对 MATERIAL 返回 0;降级到历史 matchByText 路径。
  578. for fallback_cc in _fallback_config_codes_for_strategy(query.config_code, landing):
  579. try:
  580. fallback_items = _call_match_by_text(
  581. query_text=query.query_text,
  582. config_code=fallback_cc,
  583. material_top_n=RECALL_DISPLAY_K,
  584. source_labels=labels,
  585. )
  586. except Exception as e:
  587. logger.error(
  588. "[material_recall] fallback matchByText %s 失败,试下一个:%s",
  589. fallback_cc, e,
  590. )
  591. continue
  592. fmats, fstats = _items_to_materials(fallback_items, RECALL_SIM_THRESHOLD)
  593. fquery = RecallQuery(
  594. query_text=query.query_text,
  595. config_code=fallback_cc,
  596. element_dimension=query.element_dimension,
  597. point_type=query.point_type,
  598. standard_element=query.standard_element,
  599. contribution_score=query.contribution_score,
  600. dt=query.dt,
  601. )
  602. fmats = _merge_materials_by_policy([(fquery, fmats)])
  603. logger.info(
  604. "[material_recall] fallback=%s 返回 %d 条,⊘ 黑名单 %d,⊘ score<%.2f %d → 保留 %d(cost desc)",
  605. fallback_cc, len(fallback_items), fstats["blacklist"],
  606. RECALL_SIM_THRESHOLD, fstats["low_score"],
  607. len(fmats),
  608. )
  609. if fmats:
  610. return fmats[:final_top_n]
  611. logger.info(
  612. "[material_recall] landing video_id=%d 所有 %d ODPS 策略全失败,返回空(上层换 landing)",
  613. landing.video_id, len(queries),
  614. )
  615. return []