Параллельная обработка вместо широких потоков

Замеры: распознавание даёт x52 на 4 потоках против x13 на 16, разделение
говорящих x33 против x8. Дальше четырёх потоков синхронизация съедает весь
выигрыш, поэтому ядра занимаются несколькими задачами сразу.

По умолчанию 4 потока на задачу и до 4 задач параллельно. Захват задачи
из очереди стал атомарным - без этого два воркера брали одну и ту же.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Vladimir Bryzgalov
2026-08-15 22:33:54 +05:00
co-authored by Claude Opus 5
parent e3b4e3c73f
commit be837bab0b
6 changed files with 176 additions and 38 deletions
+40 -21
View File
@@ -25,29 +25,34 @@ log = logging.getLogger("talkscore-asr")
settings: Settings = load_settings()
allowlist = parse_allowlist(settings.allow_ips)
store = JobStore(settings.data_dir / "jobs.db")
pipeline = Pipeline(
models_dir=settings.models_dir,
threads=settings.effective_threads(),
replacements_path=settings.replacements_path,
base_dir=settings.base_dir,
)
def make_pipeline() -> Pipeline:
"""Своя копия моделей на каждого воркера: они не рассчитаны на общий доступ."""
return Pipeline(
models_dir=settings.models_dir,
threads=settings.effective_threads(),
replacements_path=settings.replacements_path,
base_dir=settings.base_dir,
)
pipeline = make_pipeline()
_worker_stop = threading.Event()
_state: dict = {"ready": False, "error": None}
def _process(job_id: str) -> None:
def _process(job_id: str, worker: Pipeline) -> None:
"""Задача уже помечена в работе тем, кто её забрал."""
job = store.get(job_id)
if job is None:
return
store.mark_running(job_id)
upload = settings.data_dir / "uploads" / job_id
try:
with tempfile.TemporaryDirectory() as tmp:
wav = Path(tmp) / "audio.wav"
to_wav16k(upload, wav, pipeline.ffmpeg)
to_wav16k(upload, wav, worker.ffmpeg)
speakers = int(job["options"].get("speakers", settings.speakers))
result = pipeline.transcribe(wav, num_speakers=speakers)
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,
@@ -59,16 +64,18 @@ def _process(job_id: str) -> None:
upload.unlink(missing_ok=True)
def _worker_loop() -> None:
"""Один воркер: модели тяжёлые, параллельные задачи только мешали бы друг другу."""
def _worker_loop(index: int, worker: Pipeline) -> None:
"""Разбирает очередь. Задача захватывается атомарно, поэтому воркеров может быть много."""
last_cleanup = 0.0
while not _worker_stop.is_set():
if _state["ready"]:
job_id = store.take_next()
job_id = store.claim_next()
if job_id:
_process(job_id)
log.info("воркер %d взял задачу %s", index, job_id)
_process(job_id, worker)
continue
if time.time() - last_cleanup > 3600:
# Уборкой занимается только первый воркер, чтобы не делать её хором.
if index == 0 and time.time() - last_cleanup > 3600:
removed = store.cleanup(settings.keep_results_hours)
if removed:
log.info("удалено старых задач: %d", removed)
@@ -85,13 +92,22 @@ async def lifespan(app: FastAPI):
except ModelsMissing as exc:
_state["error"] = str(exc)
log.error("сервис запущен без моделей: %s", exc)
worker = threading.Thread(target=_worker_loop, name="asr-worker", daemon=True)
worker.start()
log.info("сервис слушает %s:%s, потоков %d", settings.host, settings.port,
settings.effective_threads())
workers = []
for i in range(settings.effective_workers()):
# Первый воркер использует уже прогретый конвейер, остальные греются сами
# при первой задаче: держать копии моделей впустую незачем.
worker = pipeline if i == 0 else make_pipeline()
thread = threading.Thread(target=_worker_loop, args=(i, worker),
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())
yield
_worker_stop.set()
worker.join(timeout=5)
for thread in workers:
thread.join(timeout=5)
# Штатные /docs и /openapi.json отключены: они не требуют токена. Вместо них
@@ -166,6 +182,7 @@ def health(request: Request) -> JSONResponse:
"error": _state["error"],
"queue": store.stats(),
"threads": settings.effective_threads(),
"workers": settings.effective_workers(),
"your_ip": client_ip,
"your_ip_allowed": ip_allowed(client_ip, allowlist),
"ip_filter_active": bool(allowlist),
@@ -270,7 +287,9 @@ def run() -> None:
access = f"только с {nets} адресов" if nets else "со ВСЕХ адресов"
print(f"\n Настройки: {settings.base_dir / 'config.toml'}"
f"\n Токен: {'задан' if settings.token else 'НЕ ЗАДАН, сервис никого не пустит'}"
f"\n Доступ: {access}")
f"\n Доступ: {access}"
f"\n Обработка: {settings.effective_workers()} задач одновременно"
f" по {settings.effective_threads()} потока")
print(f"\n talkscore-asr слушает http://{settings.host}:{settings.port}"
f"\n Потоков: {settings.effective_threads()}"
f"\n Проверка: curl http://localhost:{settings.port}/health"