xueyiming 8 órája
szülő
commit
501883f0b2

+ 2 - 2
supply_infra/scheduler/constants.py

@@ -3,5 +3,5 @@
 SUPPLY_PIPELINE_JOB_ID = "run_supply_pipeline"
 SUPPLY_PIPELINE_JOB_ID = "run_supply_pipeline"
 SUPPLY_PIPELINE_JOB_NAME = "供给数据流水线"
 SUPPLY_PIPELINE_JOB_NAME = "供给数据流水线"
 
 
-# find_agent:当日全部 S/A 需求(有拓展点位),单线程串行找视频;有效视频满 200 提前结束
-PIPELINE_FIND_AGENT_WORKERS = 1
+# find_agent:当日全部 S/A 需求(有拓展点位),默认 2 线程并发找视频
+PIPELINE_FIND_AGENT_WORKERS = 2

+ 8 - 13
supply_infra/scheduler/jobs/discover_videos_from_demands.py

@@ -1,5 +1,5 @@
 """
 """
-从全部 S/A 级需求及其拓展点位触发 find_agent 视频发现;有效视频满 400 提前结束。
+从全部 S/A 级需求及其拓展点位触发 find_agent 视频发现;有效视频满 200 提前结束。
 
 
 任务层负责查库与组装上下文;Agent 负责搜索、画像与分池落库。
 任务层负责查库与组装上下文;Agent 负责搜索、画像与分池落库。
 """
 """
@@ -27,8 +27,8 @@ from supply_infra.services.video_discovery_service import get_video_discovery_se
 
 
 logger = logging.getLogger(__name__)
 logger = logging.getLogger(__name__)
 
 
-_DEFAULT_WORKERS = 1
-_DAILY_PASSED_VIDEO_LIMIT = 400
+_DEFAULT_WORKERS = 2
+_DAILY_PASSED_VIDEO_LIMIT = 200
 
 
 
 
 def _resolve_biz_dt(biz_dt: str | None) -> str:
 def _resolve_biz_dt(biz_dt: str | None) -> str:
@@ -140,17 +140,11 @@ def discover_videos_from_demands(
     每条记录对应一个 demand_grade + 其下全部视频与全部拓展点位。
     每条记录对应一个 demand_grade + 其下全部视频与全部拓展点位。
     执行前会预写 video_discovery_run,并按 biz_dt + demand_grade_id 跳过已执行记录。
     执行前会预写 video_discovery_run,并按 biz_dt + demand_grade_id 跳过已执行记录。
     默认处理全部 S/A(S 优先于 A、再按 score 排序);仅 CLI --top-limit 可人为截断。
     默认处理全部 S/A(S 优先于 A、再按 score 排序);仅 CLI --top-limit 可人为截断。
-    当日 primary 去重视频达到 200 时提前结束,否则跑完待处理队列。
-    固定单线程串行执行 find_agent(workers 参数已废弃,始终为 1)。
+    当日 primary 去重视频达到上限时提前结束,否则跑完待处理队列。
+    默认 2 线程并发执行 find_agent(每条 S/A 需求一个 worker)。
     """
     """
     started_at = datetime.now()
     started_at = datetime.now()
     batch_run_id = uuid.uuid4().hex
     batch_run_id = uuid.uuid4().hex
-    if int(workers) != 1:
-        logger.warning(
-            "discover_videos_from_demands ignores workers=%s; running single-threaded",
-            workers,
-        )
-    ensure_mysql_pool_capacity(1)
 
 
     try:
     try:
         resolved_biz_dt, contexts = list_find_demand_contexts(
         resolved_biz_dt, contexts = list_find_demand_contexts(
@@ -223,7 +217,8 @@ def discover_videos_from_demands(
         logger.info("discover_videos_from_demands finished: %s", result)
         logger.info("discover_videos_from_demands finished: %s", result)
         return result
         return result
 
 
-    worker_count = 1
+    worker_count = max(1, min(int(workers), len(contexts)))
+    ensure_mysql_pool_capacity(worker_count)
     result["workers"] = worker_count
     result["workers"] = worker_count
 
 
     with ThreadPoolExecutor(max_workers=worker_count) as executor:
     with ThreadPoolExecutor(max_workers=worker_count) as executor:
@@ -377,7 +372,7 @@ if __name__ == "__main__":
     run_cli(
     run_cli(
         lambda: discover_videos_from_demands(
         lambda: discover_videos_from_demands(
             _biz_dt,
             _biz_dt,
-            workers=1,
+            workers=_read_flag_value(_args, "--workers") or _DEFAULT_WORKERS,
             offset=_read_flag_value(_args, "--offset") or 0,
             offset=_read_flag_value(_args, "--offset") or 0,
             limit=_read_flag_value(_args, "--limit"),
             limit=_read_flag_value(_args, "--limit"),
             top_limit=_read_flag_value(_args, "--top-limit"),
             top_limit=_read_flag_value(_args, "--top-limit"),

+ 3 - 3
tests/supply_infra/scheduler/test_discover_videos_from_demands.py

@@ -819,7 +819,7 @@ def test_evaluate_find_agent_run_fails_without_candidates(
 @patch(
 @patch(
     "supply_infra.scheduler.jobs.discover_videos_from_demands.list_find_demand_contexts"
     "supply_infra.scheduler.jobs.discover_videos_from_demands.list_find_demand_contexts"
 )
 )
-def test_stops_discovery_after_400_passed_videos(
+def test_stops_discovery_after_200_passed_videos(
     mock_list_contexts,
     mock_list_contexts,
     mock_filter_contexts,
     mock_filter_contexts,
     mock_count_passed,
     mock_count_passed,
@@ -844,14 +844,14 @@ def test_stops_discovery_after_400_passed_videos(
         contexts,
         contexts,
         {"total_loaded": 2, "skipped_already_done": 0},
         {"total_loaded": 2, "skipped_already_done": 0},
     )
     )
-    mock_count_passed.side_effect = [399, 400]
+    mock_count_passed.side_effect = [199, 200]
     mock_process.return_value = {"success": True, "skipped": False}
     mock_process.return_value = {"success": True, "skipped": False}
 
 
     result = discover_videos_from_demands("20260727", workers=1)
     result = discover_videos_from_demands("20260727", workers=1)
 
 
     assert mock_process.call_count == 1
     assert mock_process.call_count == 1
     assert result["processed"] == 1
     assert result["processed"] == 1
-    assert result["passed_videos"] == 400
+    assert result["passed_videos"] == 200
     assert result["stopped_by_passed_video_limit"] is True
     assert result["stopped_by_passed_video_limit"] is True