fix: 修复自适应线程池限流扩容及缩容滞后

This commit is contained in:
2026-09-11 15:59:46 +08:00
parent 13f72178e9
commit f3faad0391
4 changed files with 262 additions and 179 deletions
+120 -158
View File
@@ -1,27 +1,23 @@
"""自适应线程池。
用于对耗时的独立子任务(如逐帧 VLM OCR)做弹性并发加速
- 以滚动时间窗口统计已完成任务的平均响应时间
- 窗口内平均响应 < fast_threshold(默认 0.3s)→ 增加 1 个工作线程(上限 max_workers
- 窗口内平均响应 > slow_threshold(默认 1.0s)→ 减少 1 个工作线程(下限 min_workers
用于逐帧 VLM OCR、逐条 LLM 判定等独立 I/O 子任务
- 从 min_workers 起步,每个时间窗口按平均单任务耗时增减目标并发
- 快响应增加 1 个在途任务额度,慢响应减少 1 个,受 min/max_workers 限制
- 限流降低有效上限,发生错误的窗口禁止扩容,干净窗口逐步恢复上限
线程数从 min_workers(默认 1)起步,按实测负载自适应:服务端空闲(响应快)
就加大并发,服务端变慢就退避,避免盲目并发压垮上游(如本地 Ollama)
线程安全说明:worker 会在多个线程中并发调用,调用方需保证 worker 无共享
可变状态(registry 处理器是纯函数,符合要求);结果按输入顺序返回。
执行器按需创建线程并复用;map 只提交目标额度内的任务,不把整批输入压入
执行器队列。缩容立即限制后续提交,已发出的请求允许完成,不强制中断
结果与进度由 map 所在线程统一收集,worker 只负责处理输入;返回结果保持
输入顺序,worker 异常作为结果交给调用方决定是否重试。
"""
from __future__ import annotations
import queue
import threading
import time
from concurrent.futures import FIRST_COMPLETED, ThreadPoolExecutor, wait
from typing import Callable
# 停止哨兵:压入队列让空闲工作线程退出(用于缩容)。
_POISON = object()
def decide(
current: int,
@@ -31,11 +27,7 @@ def decide(
fast_threshold: float,
slow_threshold: float,
) -> int:
"""根据窗口平均响应时间返回调整后的目标线程数(纯决策函数)。
响应快(avg < fast_threshold)且未达上限 → 加 1;响应慢
avg > slow_threshold)且未达下限 → 减 1;其余情况保持不变。
"""
"""按平均耗时返回目标并发:快则 +1、慢则 -1,达到上下界后保持。"""
if avg < fast_threshold and current < max_workers:
return current + 1
if avg > slow_threshold and current > min_workers:
@@ -44,7 +36,7 @@ def decide(
class AdaptiveThreadPool:
"""自适应线程池:单次 map 按输入顺序返回全部结果"""
"""有界自适应执行器;允许顺序重复 map,不允许同一实例并行调用 map"""
def __init__(
self,
@@ -57,7 +49,7 @@ class AdaptiveThreadPool:
clock=time.monotonic,
on_progress: Callable[[int, int, float, float, int], None] | None = None,
) -> None:
"""初始化;clock 可注入便于测试;on_progress(done,total,rate) 每次完成回调"""
"""保存 worker、窗口策略及进度回调;clock 可在测试中注入"""
self._worker = worker
self.min_workers = max(1, min_workers)
self.max_workers = max(self.min_workers, max_workers)
@@ -65,181 +57,151 @@ class AdaptiveThreadPool:
self.fast_threshold = fast_threshold
self.slow_threshold = slow_threshold
self._clock = clock
self._queue: queue.Queue = queue.Queue()
# 并发目标线程数:决策/缩容的权威依据(线程退出是异步的,不能用
# len(_threads) 判断,否则并发缩容会重复放哨兵把全部线程毒死)。
# 并发目标是提交额度,不能用执行器已创建的线程数判断缩容。
self._target_workers = 0
self._threads: list[threading.Thread] = []
self._results: list = []
self._lock = threading.Lock()
self._stop = threading.Event()
# 取消标记:worker 检测到取消(如暂停信号)后设置,后续完成的任务
# 不再触发进度回调——暂停时队列中剩余大量任务会快速退出,若仍逐项
# 打印进度会在数秒内打出上万行日志。
self._map_lock = threading.Lock()
# 保留 cancel 的调用约定:抑制回调,worker 自行检测暂停并返回异常。
# OCR 因此仍能为每个输入得到结果,同时不会产生上万条暂停进度日志。
self._cancel_event = threading.Event()
# 有效最大线程数:初始等于 max_workers;消费错误(如 API 限流)时
# report_failure 临时收紧,连续无错误窗口后逐步回升——并发自适应配额。
self._effective_max_workers = max_workers
# 当前窗口内消费错误计数:窗口评估时无错误才允许恢复有效上限。
# 有效上限跨 map 保留:LLM 失败条目重试时继续遵守已收紧的配额。
self._effective_max_workers = self.max_workers
self._window_failures = 0
# 滚动窗口起点与已记录的单次耗时。
self._window_start = clock()
# 观测到的最大并发线程数(供测试与监控)
# 记录实际提交时的最大在途数量,供监控与测试检查
self.max_concurrency = 0
# 进度回调与计数:on_progress(已完成数, 总数, 平均速度/秒)。
self._on_progress = on_progress
self._completed = 0
self._total = 0
self._started_at = 0.0
self._elapsed_sum = 0.0
self._window_times: list[float] = []
# 最近一次窗口评估的平均单任务耗时(秒):供进度回调诊断使用,
# 与扩缩容决策共用同一依据;窗口尚未评估时为 None(回退累计平均)。
self._window_avg_time: float | None = None
def _run(self) -> None:
"""工作线程主循环:取任务 → 执行 → 记录耗时并自适应评估"""
def _run(self, item) -> tuple[object, float]:
"""执行一次 worker,保留异常对象并记录真实单任务耗时"""
start = self._clock()
try:
while not self._stop.is_set():
try:
seq, item = self._queue.get(timeout=0.2)
except queue.Empty:
continue
if item is _POISON:
# 缩容哨兵:处理完即可退出(队列计数照常)。
self._queue.task_done()
break
start = self._clock()
try:
result = self._worker(item)
except Exception as exc:
# 单任务异常不拖垮整体:以异常对象作为结果,由调用方判定。
result = exc
finally:
elapsed = self._clock() - start
self._results.append((seq, result))
# 进度回调:已完成数、总数与平均处理速度(条/秒)。
self._completed += 1
if self._on_progress is not None and not self._cancel_event.is_set():
elapsed_total = max(self._clock() - self._started_at, 1e-9)
with self._lock:
workers = self._target_workers
self._on_progress(
self._completed,
self._total,
self._completed / elapsed_total,
self._current_avg_time(elapsed_total),
workers,
)
self._tick(elapsed)
self._queue.task_done()
finally:
# 无论何种退出路径都从线程列表移除,保证线程数统计准确。
with self._lock:
if threading.current_thread() in self._threads:
self._threads.remove(threading.current_thread())
result = self._worker(item)
except Exception as exc:
result = exc
return result, self._clock() - start
def _tick(self, elapsed: float) -> None:
"""记录一次完成耗时;窗口满时按平均响应时间调整线程数"""
self._window_times.append(elapsed)
if self._clock() - self._window_start < self.window_seconds:
return
avg = sum(self._window_times) / len(self._window_times)
self._window_start = self._clock()
self._window_times.clear()
# 记录本次窗口平均耗时:进度回调据此展示"当前扩缩容依据"。
self._window_avg_time = avg
"""收集窗口耗时并调整额度;与 worker 报告限流共用锁,避免决策竞态"""
with self._lock:
current = self._target_workers
# 窗口内无消费错误 → 有效上限逐步回升(错误降下来的并发慢慢恢复)。
if (
self._window_failures == 0
and self._effective_max_workers < self.max_workers
):
self._window_times.append(elapsed)
now = self._clock()
if now - self._window_start < self.window_seconds:
return
avg = sum(self._window_times) / len(self._window_times)
self._window_start = now
self._window_times.clear()
self._window_avg_time = avg
had_failures = self._window_failures > 0
# 有错误的窗口禁止恢复上限或增加并发;干净窗口每次只恢复 1。
if not had_failures and self._effective_max_workers < self.max_workers:
self._effective_max_workers += 1
# 重置窗口错误计数,进入下一窗口。
self._window_failures = 0
# 扩容上限用有效最大线程数:错误窗口内即使响应快也不超过收紧后的上限。
self._resize(
decide(
current, avg, self.min_workers, self._effective_max_workers,
self.fast_threshold, self.slow_threshold,
target = decide(
self._target_workers, avg, self.min_workers,
self._effective_max_workers, self.fast_threshold, self.slow_threshold,
)
if had_failures:
target = min(target, self._target_workers)
self._target_workers = max(
self.min_workers, min(target, self._effective_max_workers)
)
)
def _current_avg_time(self, elapsed_total: float) -> float:
"""返回供进度回调展示的平均单任务耗时(秒)
"""返回最近窗口均值;窗口未满时返回实际单任务耗时均值
优先使用最近一次窗口评估的平均耗时(与扩缩容决策同一依据);
窗口尚未评估过时回退为启动至今的累计平均,避免无数据可看
elapsed_total 保留旧调用签名;墙钟时间除以任务数会受并发倍数影响,
因此均值改用 worker 耗时总和计算
"""
if self._window_avg_time is not None:
return self._window_avg_time
return elapsed_total / max(self._completed, 1)
return self._elapsed_sum / max(self._completed, 1)
def _resize(self, target: int) -> None:
"""调整并发目标:扩容启动新线程;缩容压入等量停止哨兵(幂等)。
以 _target_workers 为当前值:重复调用同一 target 不会重复放哨兵,
避免并发缩容把所有线程毒死导致队列任务无人处理而挂起。
"""
"""幂等调整提交额度;不创建退出哨兵,不等待积压队列消费完再缩容。"""
with self._lock:
current = self._target_workers
if target > current:
self.max_concurrency = max(self.max_concurrency, target)
for _ in range(target - current):
thread = threading.Thread(target=self._run, daemon=True)
thread.start()
self._threads.append(thread)
self._target_workers = target
elif target < current:
for _ in range(current - target):
self._queue.put((None, _POISON))
self._target_workers = target
self._target_workers = max(
self.min_workers, min(target, self._effective_max_workers)
)
def cancel(self) -> None:
"""请求取消本批任务:后续完成的任务不再触发进度回调
"""抑制本批后续进度回调;暂停与输入结果处理仍交给 worker
供调用方在工作线程内检测到外部信号(如暂停)时调用,抑制暂停后
队列中剩余任务快速退出导致的进度日志井喷;下一批 map 自动重置。
下一批 map 重置此标记,保持 OCR 暂停后继续及 LLM 重试的调用约定。
"""
self._cancel_event.set()
def report_failure(self) -> None:
"""通知一次消费错误(如 API 限流 429):临时降低有效最大线程数并缩容。
供工作线程捕获可退避错误(限流/服务端 5xx)后调用:并发立即收紧到
新上限,后续请求减少从而避开持续限流;连续无错误窗口后有效上限
逐步回升到 max_workers(见 _tick 的恢复逻辑)。
"""
"""限流/服务端错误收紧有效上限,仅允许保持或减少当前提交额度。"""
with self._lock:
self._window_failures += 1
if self._effective_max_workers > self.min_workers:
self._effective_max_workers -= 1
# 缩容到新上限(幂等:目标低于当前才放停止哨兵)。
self._resize(self._effective_max_workers)
self._effective_max_workers = max(
self.min_workers, self._effective_max_workers - 1
)
# 上限 20 -> 19 不意味着当前 1 个任务应扩到 19 个。
self._target_workers = min(self._target_workers, self._effective_max_workers)
def map(self, items) -> list:
"""按输入顺序返回每个 item 经 worker 处理后的结果列表"""
self._results = []
self._completed = 0
self._total = len(items)
self._started_at = self._clock()
self._stop.clear()
# 每批任务开始时重置取消状态:上一批的取消不延续到下一批。
self._cancel_event.clear()
# 上一批任务结束后工作线程已全部退出(_stop 停止)但 _target_workers
# 仍记旧值,_resize 不会重新启动线程——实际无线程时归零后重建。
with self._lock:
if not self._threads:
self._target_workers = 0
self._resize(self.min_workers)
for seq, item in enumerate(items):
self._queue.put((seq, item))
self._queue.join()
self._stop.set()
with self._lock:
threads = list(self._threads)
for thread in threads:
thread.join(1.0)
self._results.sort(key=lambda pair: pair[0])
return [result for _, result in self._results]
"""有界提交并按输入顺序返回结果;上下文退出时回收执行器线程"""
if not self._map_lock.acquire(blocking=False):
raise RuntimeError("同一自适应线程池不能同时执行多个 map")
try:
self._completed = 0
self._total = len(items)
self._started_at = self._clock()
self._elapsed_sum = 0.0
self._cancel_event.clear()
with self._lock:
# 新批次重新计时,空闲时间不构成快响应窗口;错误上限仍保留。
self._window_start = self._started_at
self._window_times.clear()
self._window_avg_time = None
self._target_workers = self.min_workers
results = [None] * self._total
pending = {}
iterator = iter(enumerate(items))
exhausted = False
with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
while True:
# 与 report_failure 共用锁:每次提交都依据最新额度。
# pending 包含尚未收集的完成任务,限制只会更保守,不会超额。
with self._lock:
while not exhausted and len(pending) < self._target_workers:
try:
seq, item = next(iterator)
except StopIteration:
exhausted = True
break
future = executor.submit(self._run, item)
pending[future] = seq
self.max_concurrency = max(self.max_concurrency, len(pending))
if not pending:
break
done, _ = wait(pending, return_when=FIRST_COMPLETED)
for future in sorted(done, key=pending.__getitem__):
seq = pending.pop(future)
result, elapsed = future.result()
results[seq] = result
self._completed += 1
self._elapsed_sum += elapsed
# 单一收集线程串行回调和统计,不再发生 queue.task_done
# 因回调异常未执行而使整批永久挂起的问题。
if self._on_progress is not None and not self._cancel_event.is_set():
elapsed_total = max(self._clock() - self._started_at, 1e-9)
with self._lock:
workers = self._target_workers
self._on_progress(
self._completed, self._total,
self._completed / elapsed_total,
self._current_avg_time(elapsed_total), workers,
)
self._tick(elapsed)
return results
finally:
self._map_lock.release()