diff --git a/core/task_manager.py b/core/task_manager.py index d8c0391..76cc08e 100644 --- a/core/task_manager.py +++ b/core/task_manager.py @@ -40,6 +40,10 @@ from tasks import list_task_types, get_task_class _log = get_logger("core.tm") +# 错峰启动间隔(秒):批量触发时设备逐个开始占用/连接,避免 adb 全局锁串行堆积 +# 和 STF occupy 并发风暴。100 台 × 0.2s = 20s 全部开始启动。 +_START_STAGGER_SEC = 0.2 + # ================== 设备分组 ================== class DeviceGroup: @@ -533,7 +537,11 @@ class TaskManager: return {"ok": True, "msg": f"任务 {job.name} 已触发"} def _run_job(self, job): - """执行任务:为每个目标设备起 worker(含重试循环)。""" + """执行任务:为每个目标设备起 worker(含重试循环)。 + + 大量设备同时启动会触发 adb connect 全局锁串行 + STF occupy 并发风暴, + 这里按 _START_STAGGER_SEC 间隔逐个启动,分摊连接压力(100 台 × 0.2s = 20s)。 + """ serials = job.resolve_serials(self) if not serials: _log.warning(f"任务 {job.name} 无可用设备") @@ -546,13 +554,14 @@ class TaskManager: max_attempts = max(1, job.retry.get("max_attempts", 1)) delay = job.retry.get("delay", 60) - for serial in serials: - # 每台设备一个重试循环线程,互不影响 + for idx, serial in enumerate(serials): + # 每台设备一个重试循环线程,互不影响;错峰延迟在各自线程内等待 t = threading.Thread(target=self._run_with_retry, - args=(task, serial, job, max_attempts, delay), daemon=True) + args=(task, serial, job, max_attempts, delay, + idx * _START_STAGGER_SEC), daemon=True) t.start() - def _run_with_retry(self, task, serial, job, max_attempts, delay): + def _run_with_retry(self, task, serial, job, max_attempts, delay, start_delay=0): """单设备任务执行 + 重试。 异常分类: @@ -560,6 +569,9 @@ class TaskManager: 其他异常 — 按 max_attempts 重试 用户停止(stop_device)会加入 _stop_requested,阻止任何后续重试。 """ + # 错峰启动:在各自线程内等待,分摊批量启动的连接/占用压力 + if start_delay > 0: + time.sleep(start_delay) try: for attempt in range(1, max_attempts + 1): # 用户已请求停止 → 不再启动新 attempt