فهرست منبع

增加全流程定时任务测试

xueyiming 1 هفته پیش
والد
کامیت
5ea5e65be3

+ 1 - 1
.env.example

@@ -35,7 +35,7 @@ ODPS_ACCESS_KEY=
 ODPS_PROJECT=
 ODPS_ENDPOINT=https://service.cn.maxcompute.aliyun.com/api
 
-# Scheduler(供给流水线每天 14:00 Asia/Shanghai 触发)
+# Scheduler(供给流水线每天 14:30 Asia/Shanghai 触发)
 SCHEDULER_ENABLED=true
 SCHEDULER_TIMEZONE=Asia/Shanghai
 # Aliyun OSS (agent 运行日志可视化上传;qwen 视频截断后片段也上传到此 bucket)

+ 14 - 0
agents/find_agent/demand_run.py

@@ -131,17 +131,29 @@ def _serialize_video(video: FindDemandVideo) -> dict[str, Any]:
     }
 
 
+def _grade_priority_key(row: Any) -> tuple[float, int, str, int]:
+    """S 优先于 A;同等级按 score 降序。"""
+    score = float(row.score) if row.score is not None else -1.0
+    grade_rank = 0 if str(row.grade) == "S" else 1
+    return (-score, grade_rank, str(row.demand_name), int(row.id))
+
+
 def load_find_demand_contexts(
     session,
     biz_dt: str,
     *,
     grades: Iterable[str] = ("S", "A"),
+    top_limit: int | None = None,
 ) -> list[FindDemandContext]:
     """加载指定业务日 S/A 需求,每个需求词组装为一条完整上下文。"""
     grades_rows = DemandGradeRepository(session).list_by_biz_dt_and_grades(biz_dt, grades)
     if not grades_rows:
         return []
 
+    grades_rows = sorted(grades_rows, key=_grade_priority_key)
+    if top_limit is not None and top_limit > 0:
+        grades_rows = grades_rows[: int(top_limit)]
+
     expansion_repo = DemandVideoExpansionRepository(session)
     detail_repo = MultiDemandVideoDetailRepository(session)
 
@@ -229,6 +241,7 @@ def list_find_demand_contexts(
     biz_dt: str | None = None,
     *,
     grades: Iterable[str] = ("S", "A"),
+    top_limit: int | None = None,
 ) -> tuple[str, list[FindDemandContext]]:
     """返回解析后的业务日与全部待执行上下文。"""
     resolved_biz_dt = _resolve_biz_dt(biz_dt)
@@ -237,6 +250,7 @@ def list_find_demand_contexts(
             session,
             resolved_biz_dt,
             grades=grades,
+            top_limit=top_limit,
         )
     return resolved_biz_dt, contexts
 

+ 5 - 1
jobs/discover_videos_from_demands.py

@@ -7,7 +7,7 @@
     python jobs/discover_videos_from_demands.py 20260721 3        # 指定业务日 + 3 并发
     python jobs/discover_videos_from_demands.py 20260721 1 --limit 1   # 只跑 1 条
     python jobs/discover_videos_from_demands.py 20260721 1 --offset 1  # 跳过第 0 条
-    python jobs/discover_videos_from_demands.py 20260721 1 --force     # 忽略已执行记录重跑
+    python jobs/discover_videos_from_demands.py 20260721 2 --top-limit 200   # top200 + 2 并发
 """
 from __future__ import annotations
 
@@ -39,6 +39,7 @@ def main(
     *,
     offset: int = 0,
     limit: int | None = None,
+    top_limit: int | None = None,
     skip_finished: bool = True,
     force: bool = False,
 ) -> dict:
@@ -48,6 +49,7 @@ def main(
         workers=workers,
         offset=offset,
         limit=limit,
+        top_limit=top_limit,
         skip_finished=skip_finished,
         force=force,
     )
@@ -64,6 +66,7 @@ if __name__ == "__main__":
 
     offset_arg = _read_flag_value(args, "--offset") or 0
     limit_arg = _read_flag_value(args, "--limit")
+    top_limit_arg = _read_flag_value(args, "--top-limit")
     skip_finished = "--force" not in args
     force = "--force" in args
 
@@ -72,6 +75,7 @@ if __name__ == "__main__":
         workers_arg,
         offset=offset_arg,
         limit=limit_arg,
+        top_limit=top_limit_arg,
         skip_finished=skip_finished,
         force=force,
     )

+ 1 - 1
jobs/publish_videos_from_discovery.py

@@ -1,7 +1,7 @@
 #!/usr/bin/env python3
 """从 video_discovery_candidate 均匀分发视频到各 AIGC 发布计划。
 
-仅发布 decision_bucket 为 primary / backup 的候选;不按品类匹配。
+仅发布 decision_bucket 为 primary / backup 且 aweme_id 非空的候选;不按品类匹配。
 
 用法:
     python jobs/publish_videos_from_discovery.py

+ 20 - 10
supply_infra/aigc/client.py

@@ -9,6 +9,8 @@ from zoneinfo import ZoneInfo
 
 import requests
 
+from supply_infra.config import get_infra_settings
+
 logger = logging.getLogger(__name__)
 
 AIGC_BASE_URL = "https://aigc-api.aiddit.com"
@@ -17,6 +19,7 @@ GET_PRODUCE_PLAN_DETAIL_BY_ID = f"{AIGC_BASE_URL}/aigc/produce/plan/detail"
 PRODUCE_PLAN_SAVE = f"{AIGC_BASE_URL}/aigc/produce/plan/save"
 DEFAULT_TIMEOUT = 60.0
 SHANGHAI_TZ = ZoneInfo("Asia/Shanghai")
+MAX_VIDEOS_PER_CRAWLER_PLAN = 10
 _INPUT_SOURCE_CHECK_KEYS = ("inputSourceModal", "inputSourceChannel", "contentType")
 
 
@@ -24,15 +27,20 @@ class AigcClient:
     """AIGC 平台 HTTP 客户端(视频爬取计划 + 生成计划绑定)。"""
 
     def __init__(self, token: str | None = None, *, dry_run: bool | None = None) -> None:
-        self.token = (token or os.getenv("AIGC_API_TOKEN") or "").strip()
+        settings = get_infra_settings()
+        self.token = (
+            token
+            or os.getenv("AIGC_API_TOKEN")
+            or settings.aigc_api_token
+            or ""
+        ).strip()
         if dry_run is None:
-            dry_run = os.getenv("AIGC_DRY_RUN", "").strip().lower() in {
-                "1",
-                "true",
-                "yes",
-                "on",
-            }
-        self.dry_run = dry_run
+            env_dry_run = os.getenv("AIGC_DRY_RUN", "").strip().lower()
+            if env_dry_run:
+                dry_run = env_dry_run in {"1", "true", "yes", "on"}
+            else:
+                dry_run = settings.aigc_dry_run
+        self.dry_run = bool(dry_run)
 
     def create_video_crawler_plan(
         self,
@@ -42,8 +50,10 @@ class AigcClient:
     ) -> dict[str, Any]:
         if not aweme_ids:
             raise ValueError("aweme_ids 不能为空")
-        if len(aweme_ids) > 100:
-            raise ValueError(f"单次最多 100 个视频,当前 {len(aweme_ids)} 个")
+        if len(aweme_ids) > MAX_VIDEOS_PER_CRAWLER_PLAN:
+            raise ValueError(
+                f"单次最多 {MAX_VIDEOS_PER_CRAWLER_PLAN} 个视频,当前 {len(aweme_ids)} 个"
+            )
 
         dt = datetime.now(SHANGHAI_TZ).strftime("%Y%m%d%H%M%S")
         crawler_plan_name = plan_name or f"【SupplyAgent】抖音视频直接抓取-{dt}-抖音"

+ 3 - 1
supply_infra/db/repositories/video_discovery_repo.py

@@ -323,7 +323,7 @@ class VideoDiscoveryRepository(BaseRepository[VideoDiscoveryRun]):
         skip_published: bool = True,
         limit: int | None = None,
     ) -> list[VideoDiscoveryCandidate]:
-        """查询待提交 AIGC 的候选视频;仅 primary / backup 分池。"""
+        """查询待提交 AIGC 的候选视频;仅 primary / backup 分池,且 aweme_id 非空。"""
         stmt = (
             select(VideoDiscoveryCandidate)
             .join(
@@ -331,6 +331,8 @@ class VideoDiscoveryRepository(BaseRepository[VideoDiscoveryRun]):
                 VideoDiscoveryRun.run_id == VideoDiscoveryCandidate.run_id,
             )
             .where(VideoDiscoveryCandidate.decision_bucket.in_(_PUBLISHABLE_BUCKETS))
+            .where(VideoDiscoveryCandidate.aweme_id.is_not(None))
+            .where(func.trim(VideoDiscoveryCandidate.aweme_id) != "")
             .order_by(
                 VideoDiscoveryCandidate.value_score.desc(),
                 VideoDiscoveryCandidate.share_count.desc(),

+ 11 - 9
supply_infra/scheduler/app.py

@@ -9,7 +9,12 @@ from apscheduler.schedulers.background import BackgroundScheduler
 from apscheduler.triggers.cron import CronTrigger
 
 from supply_infra.config import get_infra_settings
-from supply_infra.scheduler.constants import SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME
+from supply_infra.scheduler.constants import (
+    PIPELINE_CRON_HOUR,
+    PIPELINE_CRON_MINUTE,
+    SUPPLY_PIPELINE_JOB_ID,
+    SUPPLY_PIPELINE_JOB_NAME,
+)
 from supply_infra.scheduler.jobs.run_supply_pipeline import run_supply_pipeline
 
 if TYPE_CHECKING:
@@ -19,9 +24,6 @@ logger = logging.getLogger(__name__)
 
 _scheduler: BackgroundScheduler | None = None
 
-_PIPELINE_CRON_HOUR = 14
-_PIPELINE_CRON_MINUTE = 0
-
 
 def create_scheduler() -> BackgroundScheduler:
     """Create and configure the scheduler with the chained supply pipeline job."""
@@ -29,12 +31,12 @@ def create_scheduler() -> BackgroundScheduler:
     tz = settings.scheduler_timezone
     scheduler = BackgroundScheduler(timezone=tz)
 
-    # 每天 14:00(上海时区)串行执行:全局树 → 需求池 → 分级 → 视频点位拓展
+    # 每天 14:30(上海时区)串行执行:全局树 → 需求池 → 分级 → 拓展 → find_agent → 发布
     scheduler.add_job(
         run_supply_pipeline,
         trigger=CronTrigger(
-            hour=_PIPELINE_CRON_HOUR,
-            minute=_PIPELINE_CRON_MINUTE,
+            hour=PIPELINE_CRON_HOUR,
+            minute=PIPELINE_CRON_MINUTE,
             timezone=tz,
         ),
         id=SUPPLY_PIPELINE_JOB_ID,
@@ -49,8 +51,8 @@ def create_scheduler() -> BackgroundScheduler:
         "Scheduler configured with %d job(s) | timezone=%s | pipeline_cron=%02d:%02d",
         len(scheduler.get_jobs()),
         tz,
-        _PIPELINE_CRON_HOUR,
-        _PIPELINE_CRON_MINUTE,
+        PIPELINE_CRON_HOUR,
+        PIPELINE_CRON_MINUTE,
     )
     return scheduler
 

+ 8 - 0
supply_infra/scheduler/constants.py

@@ -2,3 +2,11 @@
 
 SUPPLY_PIPELINE_JOB_ID = "run_supply_pipeline"
 SUPPLY_PIPELINE_JOB_NAME = "供给数据流水线"
+
+# find_agent:每日按 score 取 top N 需求,2 线程并行找片
+PIPELINE_FIND_AGENT_TOP_DEMANDS = 200
+PIPELINE_FIND_AGENT_WORKERS = 2
+
+# 供给流水线每日触发时间(Asia/Shanghai)
+PIPELINE_CRON_HOUR = 14
+PIPELINE_CRON_MINUTE = 30

+ 9 - 2
supply_infra/scheduler/jobs/discover_videos_from_demands.py

@@ -102,6 +102,7 @@ def discover_videos_from_demands(
     workers: int = _DEFAULT_WORKERS,
     offset: int = 0,
     limit: int | None = None,
+    top_limit: int | None = None,
     skip_finished: bool = True,
     force: bool = False,
 ) -> dict[str, Any]:
@@ -110,12 +111,16 @@ def discover_videos_from_demands(
 
     每条记录对应一个 demand_grade + 其下全部视频与全部拓展点位。
     执行前会预写 video_discovery_run,并按 biz_dt + demand_grade_id 跳过已执行记录。
+    top_limit 按 score 取当日 top N 需求(S 优先于 A)。
     """
     started_at = datetime.now()
     batch_run_id = uuid.uuid4().hex
 
     try:
-        resolved_biz_dt, contexts = list_find_demand_contexts(biz_dt)
+        resolved_biz_dt, contexts = list_find_demand_contexts(
+            biz_dt,
+            top_limit=top_limit,
+        )
     except Exception as exc:
         logger.exception("discover_videos_from_demands preflight failed")
         return {
@@ -140,10 +145,11 @@ def discover_videos_from_demands(
         contexts = contexts[: int(limit)]
 
     logger.info(
-        "discover_videos_from_demands start: biz_dt=%s batch_run_id=%s workers=%s offset=%s pending=%s skipped=%s",
+        "discover_videos_from_demands start: biz_dt=%s batch_run_id=%s workers=%s top_limit=%s offset=%s pending=%s skipped=%s",
         resolved_biz_dt,
         batch_run_id,
         workers,
+        top_limit,
         offset,
         len(contexts),
         preload_stats.get("skipped_already_done", 0),
@@ -155,6 +161,7 @@ def discover_videos_from_demands(
         "biz_dt": resolved_biz_dt,
         "started_at": started_at.isoformat(),
         "offset": int(offset),
+        "top_limit": top_limit,
         **preload_stats,
         "total_records": len(contexts),
         "workers": 0,

+ 20 - 9
supply_infra/scheduler/jobs/publish_videos_from_discovery.py

@@ -6,7 +6,7 @@ from datetime import datetime
 from typing import Any
 from zoneinfo import ZoneInfo
 
-from supply_infra.aigc.client import AigcClient
+from supply_infra.aigc.client import AigcClient, MAX_VIDEOS_PER_CRAWLER_PLAN
 from supply_infra.aigc.plan_map import (
     AigcPlanPair,
     assignment_summary,
@@ -24,7 +24,14 @@ logger = logging.getLogger(__name__)
 
 # 仅发布主推荐与人工备选,不包含 rejected / pending_evaluation。
 _PUBLISHABLE_BUCKETS = ("primary", "backup")
-_MAX_VIDEOS_PER_CRAWLER_PLAN = 100
+
+
+@dataclass(frozen=True)
+class PublishCandidate:
+    """脱离 ORM Session 后仍可安全使用的发布候选快照。"""
+
+    id: int
+    aweme_id: str
 
 
 @dataclass
@@ -71,11 +78,11 @@ def _resolve_biz_dt(biz_dt: str | None) -> str | None:
 
 
 def _group_candidates_by_plan(
-    candidates: list[VideoDiscoveryCandidate],
+    candidates: list[PublishCandidate],
     plan_pairs: list[AigcPlanPair],
-) -> list[tuple[AigcPlanPair, list[VideoDiscoveryCandidate]]]:
+) -> list[tuple[AigcPlanPair, list[PublishCandidate]]]:
     buckets = distribute_evenly(candidates, len(plan_pairs))
-    grouped: list[tuple[AigcPlanPair, list[VideoDiscoveryCandidate]]] = []
+    grouped: list[tuple[AigcPlanPair, list[PublishCandidate]]] = []
     for plan_pair, bucket in zip(plan_pairs, buckets, strict=False):
         if bucket:
             grouped.append((plan_pair, bucket))
@@ -93,7 +100,7 @@ def publish_videos_from_discovery(
     """
     从 video_discovery_candidate 读取视频,按轮询均匀分配到各 AIGC 计划对。
 
-    仅处理 decision_bucket 为 primary / backup 的候选,不区分品类。
+    仅处理 decision_bucket 为 primary / backup 且 aweme_id 非空的候选,不区分品类。
     """
     resolved_biz_dt = _resolve_biz_dt(biz_dt)
     plan_pairs = list_unique_plan_pairs()
@@ -102,12 +109,16 @@ def publish_videos_from_discovery(
 
     with get_session() as session:
         repo = VideoDiscoveryRepository(session)
-        candidates = repo.list_publishable_candidates(
+        orm_candidates = repo.list_publishable_candidates(
             run_id=run_id,
             biz_dt=resolved_biz_dt,
             skip_published=skip_published,
             limit=limit,
         )
+        candidates = [
+            PublishCandidate(id=int(item.id), aweme_id=str(item.aweme_id))
+            for item in orm_candidates
+        ]
 
     if not candidates:
         return {
@@ -133,7 +144,7 @@ def publish_videos_from_discovery(
     batch_results: list[PublishBatchResult] = []
 
     for plan_pair, plan_candidates in grouped:
-        candidate_batches = chunk_list(plan_candidates, _MAX_VIDEOS_PER_CRAWLER_PLAN)
+        candidate_batches = chunk_list(plan_candidates, MAX_VIDEOS_PER_CRAWLER_PLAN)
         for batch_index, candidate_batch in enumerate(candidate_batches, start=1):
             aweme_batch = [str(item.aweme_id) for item in candidate_batch]
             id_batch = [int(item.id) for item in candidate_batch]
@@ -173,7 +184,7 @@ def publish_videos_from_discovery(
             if not batch.bind_success:
                 batch.bind_error = str(bind_result.get("error") or "绑定生成计划失败")
 
-            if not dry_run and crawler_plan_id:
+            if not dry_run and crawler_plan_id and batch.bind_success:
                 with get_session() as session:
                     updated = VideoDiscoveryRepository(session).mark_candidates_aigc_plans(
                         id_batch,

+ 26 - 1
supply_infra/scheduler/jobs/run_supply_pipeline.py

@@ -5,6 +5,8 @@
 2. ODPS → MySQL 策略需求池同步(当天 biz_dt)
 3. 需求池分级评估(同上 biz_dt)
 4. S/A 需求视频点位拓展判断
+5. find_agent 视频发现(top 200 需求,2 线程并行)
+6. AIGC 发布(find_agent 完成后发布全部符合条件的视频)
 
 各子步骤内部已做去重(INSERT IGNORE、diff 同步、跳过已分级词等);
 本文件额外用进程内锁防止同一轮次并发重入,并隔离各步骤异常:前一步失败时记录
@@ -20,14 +22,22 @@ from zoneinfo import ZoneInfo
 
 from supply_infra.config import get_infra_settings
 from supply_infra.scheduler.constants import (
+    PIPELINE_FIND_AGENT_TOP_DEMANDS,
+    PIPELINE_FIND_AGENT_WORKERS,
     SUPPLY_PIPELINE_JOB_ID,
     SUPPLY_PIPELINE_JOB_NAME,
 )
 from supply_infra.scheduler.job_execution import JobExecutionRecorder, record_skipped
+from supply_infra.scheduler.jobs.discover_videos_from_demands import (
+    discover_videos_from_demands,
+)
 from supply_infra.scheduler.jobs.expand_demand_from_video_points import (
     expand_demand_from_video_points,
 )
 from supply_infra.scheduler.jobs.grade_demand_pool import grade_demand_pool
+from supply_infra.scheduler.jobs.publish_videos_from_discovery import (
+    publish_videos_from_discovery,
+)
 from supply_infra.scheduler.jobs.sync_global_tree_odps_to_mysql import sync_global_tree_odps_to_mysql
 from supply_infra.scheduler.jobs.sync_multi_demand_pool_odps_to_mysql import (
     sync_multi_demand_pool_odps_to_mysql,
@@ -94,7 +104,8 @@ def _preflight_failure_result(biz_dt: str | None, exc: Exception) -> dict[str, A
 
 def run_supply_pipeline(biz_dt: str | None = None) -> dict[str, Any]:
     """
-    按顺序执行全局树同步 → 需求池同步 → 需求分级 → 视频点位拓展。
+    按顺序执行全局树同步 → 需求池同步 → 需求分级 → 视频点位拓展
+    → find_agent 找片 → AIGC 发布。
 
     Args:
         biz_dt: 业务日 YYYYMMDD;省略则取当天。
@@ -170,6 +181,20 @@ def run_supply_pipeline(biz_dt: str | None = None) -> dict[str, Any]:
                     workers=5,
                 ),
             ),
+            (
+                "discover_videos",
+                lambda: discover_videos_from_demands(
+                    biz_dt=resolved_biz_dt,
+                    workers=PIPELINE_FIND_AGENT_WORKERS,
+                    top_limit=PIPELINE_FIND_AGENT_TOP_DEMANDS,
+                ),
+            ),
+            (
+                "publish_videos",
+                lambda: publish_videos_from_discovery(
+                    biz_dt=resolved_biz_dt,
+                ),
+            ),
         ]
         for step_name, action in steps:
             payload, step_success, step_error = _run_step(step_name, action)