| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162 |
- """
- 补充 demand_belong_category 与 multi_demand_pool_di 的匹配关系,
- 并回填词级 video_list(匹配池视频去重截断最多 10 个)。
- """
- from __future__ import annotations
- import json
- import logging
- from typing import Any
- from supply_infra.db.repositories.demand_belong_category_repo import (
- DemandBelongCategoryRepository,
- )
- from supply_infra.db.repositories.demand_belong_pool_rel_repo import (
- DemandBelongPoolRelRepository,
- )
- from supply_infra.db.repositories.multi_demand_pool_di_repo import (
- MultiDemandPoolDiRepository,
- )
- from supply_infra.db.session import get_session
- logger = logging.getLogger(__name__)
- _VIDEO_LIST_LIMIT = 10
- def _parse_video_ids(raw: str | None) -> list[str]:
- if not raw:
- return []
- try:
- parsed = json.loads(raw)
- except json.JSONDecodeError:
- return []
- if not isinstance(parsed, list):
- return []
- return [str(v).strip() for v in parsed if v is not None and str(v).strip()]
- def _merge_video_ids(video_lists: list[str | None], limit: int = _VIDEO_LIST_LIMIT) -> str | None:
- """按出现顺序去重,截断到 limit。"""
- seen: set[str] = set()
- ordered: list[str] = []
- for raw in video_lists:
- for vid in _parse_video_ids(raw):
- if vid in seen:
- continue
- seen.add(vid)
- ordered.append(vid)
- if len(ordered) >= limit:
- return json.dumps(ordered, ensure_ascii=False)
- if not ordered:
- return None
- return json.dumps(ordered, ensure_ascii=False)
- def _matches_explicit_token(name: str, demand_name: str) -> bool:
- """只接受上游空格分词的显式词,不把任意子串直接升级为正式关系。"""
- normalized_name = name.strip()
- normalized_demand = demand_name.strip()
- if not normalized_name or not normalized_demand:
- return False
- return normalized_name in {
- token.strip() for token in normalized_demand.split() if token.strip()
- }
- def sync_demand_belong_pool_rel(biz_dt: str) -> dict[str, Any]:
- """
- 增量建立词 ↔ 需求池匹配边,并更新词级 video_list。
- - 只读取 biz_dt 当日需求池,禁止跨日串数
- - name 必须是 demand_name 的显式空格分词,不接受任意子串
- - 当日关系先清理再精确重建,源内容修订后不会保留陈旧边
- - video_list:匹配行视频按顺序去重,最多 10 个
- """
- if len(str(biz_dt)) != 8 or not str(biz_dt).isdigit():
- raise ValueError(f"biz_dt 格式无效,应为 YYYYMMDD: {biz_dt!r}")
- logger.info("Starting demand_belong_pool_rel sync: biz_dt=%s", biz_dt)
- with get_session() as session:
- words = DemandBelongCategoryRepository(session).list_active_id_name()
- pool_rows = MultiDemandPoolDiRepository(session).list_id_name_video_lists(biz_dt)
- logger.info("Loaded words=%d pool_rows=%d", len(words), len(pool_rows))
- if not words or not pool_rows:
- result = {
- "words": len(words),
- "pool_rows": len(pool_rows),
- "biz_dt": biz_dt,
- "matched_edges": 0,
- "rejected_substring_candidates": 0,
- "replaced_edges": 0,
- "inserted": 0,
- "video_updated": 0,
- }
- logger.info("Nothing to sync: %s", result)
- return result
- candidate_pairs: list[tuple[int, int]] = []
- video_updates: dict[int, str | None] = {}
- matched_edges = 0
- rejected_substring_candidates = 0
- for belong_id, name in words:
- matched_pool_ids: list[int] = []
- matched_video_lists: list[str | None] = []
- for pool_id, demand_name, video_list in pool_rows:
- if _matches_explicit_token(name, demand_name):
- matched_pool_ids.append(pool_id)
- matched_video_lists.append(video_list)
- elif name and name in demand_name:
- rejected_substring_candidates += 1
- if not matched_pool_ids:
- video_updates[belong_id] = None
- continue
- matched_edges += len(matched_pool_ids)
- for pool_id in matched_pool_ids:
- candidate_pairs.append((belong_id, pool_id))
- video_updates[belong_id] = _merge_video_ids(matched_video_lists)
- with get_session() as session:
- rel_repo = DemandBelongPoolRelRepository(session)
- replaced_edges = rel_repo.delete_by_pool_ids(
- pool_id for pool_id, _, _ in pool_rows
- )
- insert_rows = [
- {
- "demand_belong_category_id": belong_id,
- "multi_demand_pool_di_id": pool_id,
- "biz_dt": biz_dt,
- "relation_type": "explicit_token",
- "relation_source": "demand_pool_name_tokens",
- "reason": "需求归属词完整命中上游 demand_name 的显式空格分词",
- "confidence": 1.0,
- "is_inferred": True,
- "status": "active",
- "valid_from_biz_dt": biz_dt,
- "valid_to_biz_dt": None,
- }
- for belong_id, pool_id in candidate_pairs
- ]
- inserted = rel_repo.bulk_insert_ignore(insert_rows)
- video_updated = DemandBelongCategoryRepository(session).update_video_lists(
- video_updates
- )
- result = {
- "biz_dt": biz_dt,
- "words": len(words),
- "pool_rows": len(pool_rows),
- "matched_edges": matched_edges,
- "candidate_pairs": len(candidate_pairs),
- "rejected_substring_candidates": rejected_substring_candidates,
- "replaced_edges": replaced_edges,
- "inserted": inserted,
- "video_updated": video_updated,
- }
- logger.info("demand_belong_pool_rel sync completed: %s", result)
- return result
|