Sfoglia il codice sorgente

修复一些问题

xueyiming 4 giorni fa
parent
commit
953d1a6c3c

+ 2 - 3
supply_infra/scheduler/constants.py

@@ -3,9 +3,8 @@
 SUPPLY_PIPELINE_JOB_ID = "run_supply_pipeline"
 SUPPLY_PIPELINE_JOB_NAME = "供给数据流水线"
 
-# find_agent:当日全部 S/A 需求(有拓展点位),2 线程并行;有效视频满 200 提前结束
-PIPELINE_FIND_AGENT_WORKERS = 2
-# step 子进程内 MYSQL_POOL_SIZE_CONTROL 默认须 >= 并行 worker 数
+# find_agent:当日全部 S/A 需求(有拓展点位),单线程串行找片;有效视频满 200 提前结束
+PIPELINE_FIND_AGENT_WORKERS = 1
 
 # 供给流水线每日触发时间(Asia/Shanghai)
 PIPELINE_CRON_HOUR = 14

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

@@ -128,10 +128,16 @@ def discover_videos_from_demands(
     执行前会预写 video_discovery_run,并按 biz_dt + demand_grade_id 跳过已执行记录。
     默认处理全部 S/A(S 优先于 A、再按 score 排序);仅 CLI --top-limit 可人为截断。
     当日 primary+backup 去重视频达到 200 时提前结束,否则跑完待处理队列。
+    固定单线程串行执行 find_agent(workers 参数已废弃,始终为 1)。
     """
     started_at = datetime.now()
     batch_run_id = uuid.uuid4().hex
-    ensure_mysql_pool_capacity(max(1, int(workers)))
+    if int(workers) != 1:
+        logger.warning(
+            "discover_videos_from_demands ignores workers=%s; running single-threaded",
+            workers,
+        )
+    ensure_mysql_pool_capacity(1)
 
     try:
         resolved_biz_dt, contexts = list_find_demand_contexts(
@@ -204,7 +210,7 @@ def discover_videos_from_demands(
         logger.info("discover_videos_from_demands finished: %s", result)
         return result
 
-    worker_count = max(1, min(int(workers), len(contexts)))
+    worker_count = 1
     result["workers"] = worker_count
 
     with ThreadPoolExecutor(max_workers=worker_count) as executor:
@@ -354,14 +360,11 @@ if __name__ == "__main__":
 
     _args = sys.argv[1:]
     _biz_dt = _args[0] if _args and not _args[0].startswith("-") else None
-    _workers_arg = None
-    if _biz_dt and len(_args) > 1 and not _args[1].startswith("-"):
-        _workers_arg = _args[1]
 
     run_cli(
         lambda: discover_videos_from_demands(
             _biz_dt,
-            workers=int(_workers_arg) if _workers_arg else 1,
+            workers=1,
             offset=_read_flag_value(_args, "--offset") or 0,
             limit=_read_flag_value(_args, "--limit"),
             top_limit=_read_flag_value(_args, "--top-limit"),

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

@@ -47,7 +47,7 @@ def test_stops_discovery_after_200_passed_videos(
     mock_count_passed.side_effect = [199, 200]
     mock_process.return_value = {"success": True, "skipped": False}
 
-    result = discover_videos_from_demands("20260727", workers=2)
+    result = discover_videos_from_demands("20260727", workers=1)
 
     assert mock_process.call_count == 1
     assert result["processed"] == 1