feat: 任务批量触发错峰启动,分摊 adb 锁/STF occupy 风暴

大量设备(如 100 台)同时触发时,按 0.2s/台间隔逐个启动 worker,
避免 adb connect 全局锁串行堆积和 STF occupy 并发风暴。
100 台 × 0.2s = 20s 全部开始启动,连接压力分摊。
This commit is contained in:
2026-08-08 21:09:03 +08:00
parent 47262abdcd
commit 6d3cc10bca
+17 -5
View File
@@ -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