Параллельные задачи считаются в отдельных процессах
Замер подтвердил догадку: sherpa-onnx и onnxruntime держат GIL. Две задачи в двух потоках идут ровно столько же, сколько подряд (1.04x у диаризации, 1.17x у распознавания) - отсюда и незагруженный процессор при работе. В двух процессах те же задачи дают 1.59x даже с загрузкой моделей в каждом. Пул процессов включается при workers > 1. Функции воркера вынесены в модуль, не тянущий app.main: на Windows дочерний процесс поднимается через spawn и иначе стартовал бы ещё один веб-сервер. Если процессы не запустятся, сервис откатывается на однопроцессный режим. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
e2db164f0b
commit
41f3651cc8
+52
-14
@@ -42,18 +42,25 @@ _worker_stop = threading.Event()
|
||||
_state: dict = {"ready": False, "error": None}
|
||||
|
||||
|
||||
def _process(job_id: str, worker: Pipeline) -> None:
|
||||
def _process(job_id: str, worker: Pipeline, pool) -> None:
|
||||
"""Задача уже помечена в работе тем, кто её забрал."""
|
||||
job = store.get(job_id)
|
||||
if job is None:
|
||||
return
|
||||
upload = settings.data_dir / "uploads" / job_id
|
||||
try:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
wav = Path(tmp) / "audio.wav"
|
||||
to_wav16k(upload, wav, worker.ffmpeg)
|
||||
speakers = int(job["options"].get("speakers", settings.speakers))
|
||||
result = worker.transcribe(wav, num_speakers=speakers)
|
||||
speakers = int(job["options"].get("speakers", settings.speakers))
|
||||
if pool is not None:
|
||||
# Считаем в отдельном процессе: библиотеки держат GIL, и в потоках
|
||||
# задачи выстраиваются в очередь вместо параллельной работы.
|
||||
from app.worker import run_job
|
||||
|
||||
result = pool.submit(run_job, str(upload), speakers, worker.ffmpeg).result()
|
||||
else:
|
||||
with tempfile.TemporaryDirectory() as tmp:
|
||||
wav = Path(tmp) / "audio.wav"
|
||||
to_wav16k(upload, wav, worker.ffmpeg)
|
||||
result = worker.transcribe(wav, num_speakers=speakers)
|
||||
result["filename"] = job["filename"]
|
||||
store.mark_done(job_id, result)
|
||||
log.info("задача %s готова: %.1f с аудио, x%s", job_id,
|
||||
@@ -65,7 +72,7 @@ def _process(job_id: str, worker: Pipeline) -> None:
|
||||
upload.unlink(missing_ok=True)
|
||||
|
||||
|
||||
def _worker_loop(index: int, worker: Pipeline) -> None:
|
||||
def _worker_loop(index: int, worker: Pipeline, pool=None) -> None:
|
||||
"""Разбирает очередь. Задача захватывается атомарно, поэтому воркеров может быть много."""
|
||||
last_cleanup = 0.0
|
||||
while not _worker_stop.is_set():
|
||||
@@ -73,7 +80,7 @@ def _worker_loop(index: int, worker: Pipeline) -> None:
|
||||
job_id = store.claim_next()
|
||||
if job_id:
|
||||
log.info("воркер %d взял задачу %s", index, job_id)
|
||||
_process(job_id, worker)
|
||||
_process(job_id, worker, pool)
|
||||
continue
|
||||
# Уборкой занимается только первый воркер, чтобы не делать её хором.
|
||||
if index == 0 and time.time() - last_cleanup > 3600:
|
||||
@@ -84,6 +91,31 @@ def _worker_loop(index: int, worker: Pipeline) -> None:
|
||||
_worker_stop.wait(1.0)
|
||||
|
||||
|
||||
def _make_pool(count: int):
|
||||
"""Пул процессов для расчётов. При неудаче работаем в одном процессе.
|
||||
|
||||
На Windows процессы поднимаются через spawn, и сбои тут возможны, поэтому
|
||||
отказ не должен ронять сервис - он просто станет однопроцессным.
|
||||
"""
|
||||
from concurrent.futures import ProcessPoolExecutor
|
||||
|
||||
from app.worker import init_worker
|
||||
|
||||
try:
|
||||
pool = ProcessPoolExecutor(
|
||||
max_workers=count,
|
||||
initializer=init_worker,
|
||||
initargs=(str(settings.models_dir), settings.effective_threads(),
|
||||
str(settings.replacements_path), str(settings.base_dir)),
|
||||
)
|
||||
# Пустая задача проверяет, что процессы действительно поднялись
|
||||
pool.submit(str, "ok").result(timeout=300)
|
||||
return pool
|
||||
except Exception as exc: # noqa: BLE001 - причин может быть много, важен откат
|
||||
log.warning("не удалось запустить процессы (%s), работаю в одном", exc)
|
||||
return None
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
(settings.data_dir / "uploads").mkdir(parents=True, exist_ok=True)
|
||||
@@ -93,22 +125,28 @@ async def lifespan(app: FastAPI):
|
||||
except ModelsMissing as exc:
|
||||
_state["error"] = str(exc)
|
||||
log.error("сервис запущен без моделей: %s", exc)
|
||||
count = settings.effective_workers()
|
||||
pool = None
|
||||
if count > 1 and _state["ready"]:
|
||||
pool = _make_pool(count)
|
||||
|
||||
workers = []
|
||||
for i in range(settings.effective_workers()):
|
||||
# Первый воркер использует уже прогретый конвейер, остальные греются сами
|
||||
# при первой задаче: держать копии моделей впустую незачем.
|
||||
for i in range(count):
|
||||
worker = pipeline if i == 0 else make_pipeline()
|
||||
thread = threading.Thread(target=_worker_loop, args=(i, worker),
|
||||
thread = threading.Thread(target=_worker_loop, args=(i, worker, pool),
|
||||
name=f"asr-worker-{i}", daemon=True)
|
||||
thread.start()
|
||||
workers.append(thread)
|
||||
|
||||
log.info("сервис слушает %s:%s, воркеров %d по %d потоков", settings.host,
|
||||
settings.port, settings.effective_workers(), settings.effective_threads())
|
||||
log.info("сервис слушает %s:%s, воркеров %d по %d потоков, режим %s",
|
||||
settings.host, settings.port, count, settings.effective_threads(),
|
||||
"процессы" if pool else "один процесс")
|
||||
yield
|
||||
_worker_stop.set()
|
||||
for thread in workers:
|
||||
thread.join(timeout=5)
|
||||
if pool is not None:
|
||||
pool.shutdown(wait=False, cancel_futures=True)
|
||||
|
||||
|
||||
# Штатные /docs и /openapi.json отключены: они не требуют токена. Вместо них
|
||||
|
||||
Reference in New Issue
Block a user