From be837bab0b1657d4c1586ec4542bbfa34d2389b2 Mon Sep 17 00:00:00 2001 From: Vladimir Bryzgalov Date: Sat, 15 Aug 2026 22:33:54 +0500 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B0=D1=80=D0=B0=D0=BB=D0=BB=D0=B5?= =?UTF-8?q?=D0=BB=D1=8C=D0=BD=D0=B0=D1=8F=20=D0=BE=D0=B1=D1=80=D0=B0=D0=B1?= =?UTF-8?q?=D0=BE=D1=82=D0=BA=D0=B0=20=D0=B2=D0=BC=D0=B5=D1=81=D1=82=D0=BE?= =?UTF-8?q?=20=D1=88=D0=B8=D1=80=D0=BE=D0=BA=D0=B8=D1=85=20=D0=BF=D0=BE?= =?UTF-8?q?=D1=82=D0=BE=D0=BA=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Замеры: распознавание даёт x52 на 4 потоках против x13 на 16, разделение говорящих x33 против x8. Дальше четырёх потоков синхронизация съедает весь выигрыш, поэтому ядра занимаются несколькими задачами сразу. По умолчанию 4 потока на задачу и до 4 задач параллельно. Захват задачи из очереди стал атомарным - без этого два воркера брали одну и ту же. Co-Authored-By: Claude Opus 5 (1M context) --- app/config.py | 27 ++++++++--- app/main.py | 61 +++++++++++++++--------- app/store.py | 15 ++++-- app/version.py | 2 +- tests/test_concurrency.py | 97 +++++++++++++++++++++++++++++++++++++++ tests/test_store.py | 12 ++--- 6 files changed, 176 insertions(+), 38 deletions(-) create mode 100644 tests/test_concurrency.py diff --git a/app/config.py b/app/config.py index 0aa2e5f..c6c4d71 100644 --- a/app/config.py +++ b/app/config.py @@ -34,8 +34,12 @@ allow_ips = "127.0.0.1, ::1" docs = false [processing] -# Сколько потоков отдать распознаванию. 0 = половина ядер. +# Потоков на одну задачу. 0 = 4, это оптимум по замерам: и распознавание, +# и разделение говорящих на 16 потоках работают вчетверо медленнее, чем на 4. threads = 0 +# Сколько записей обрабатывать одновременно. 0 = по числу ядер, но не больше 4. +# Каждый воркер держит свою копию моделей, это около 1 ГБ памяти на воркера. +workers = 0 # Ожидаемое число говорящих в записи. 0 = определять автоматически # (на реальных звонках работает плохо, для диалога ставьте 2). speakers = 2 @@ -63,6 +67,7 @@ class Settings: allow_ips: str = "" docs: bool = False threads: int = 0 + workers: int = 0 speakers: int = 2 max_upload_mb: int = 500 keep_results_hours: float = 72.0 @@ -88,12 +93,21 @@ class Settings: def max_upload_bytes(self) -> int: return self.max_upload_mb * 1024 * 1024 + # Замеры на 16-ядерной машине: распознавание x52 на 4 потоках против x13 + # на 16, разделение говорящих x33 против x8. Дальше 4 потоков накладные + # расходы на синхронизацию съедают весь выигрыш, поэтому ядра занимаем + # несколькими задачами сразу, а не шириной одной. + OPTIMAL_THREADS = 4 + MAX_WORKERS = 4 + def effective_threads(self) -> int: - if self.threads > 0: - return self.threads - # Половина логических ядер: оставляем запас, чтобы машина не вставала колом - # во время обработки. На Ryzen 9 9950X это 16 потоков. - return max(1, (os.cpu_count() or 4) // 2) + return self.threads if self.threads > 0 else self.OPTIMAL_THREADS + + def effective_workers(self) -> int: + if self.workers > 0: + return self.workers + cores = os.cpu_count() or 4 + return max(1, min(self.MAX_WORKERS, cores // self.effective_threads())) def load_settings(config_path: Path | None = None) -> Settings: @@ -116,6 +130,7 @@ def load_settings(config_path: Path | None = None) -> Settings: allow_ips=str(security.get("allow_ips", "")), docs=bool(security.get("docs", False)), threads=int(proc.get("threads", 0)), + workers=int(proc.get("workers", 0)), speakers=int(proc.get("speakers", 2)), max_upload_mb=int(proc.get("max_upload_mb", 500)), keep_results_hours=float(proc.get("keep_results_hours", 72)), diff --git a/app/main.py b/app/main.py index f02745b..6bf9234 100644 --- a/app/main.py +++ b/app/main.py @@ -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" diff --git a/app/store.py b/app/store.py index a0f6b62..c47a563 100644 --- a/app/store.py +++ b/app/store.py @@ -82,13 +82,20 @@ class JobStore: job["options"] = json.loads(job["options"]) if job["options"] else {} return job - def take_next(self) -> str | None: - """Возвращает id самой старой задачи в очереди, не меняя её статус.""" + def claim_next(self) -> str | None: + """Забирает самую старую задачу из очереди и сразу помечает её в работе. + + Выборка и пометка выполняются одним оператором под общим замком: иначе + два воркера успевают увидеть одну и ту же задачу и берут её оба. + """ with self._lock: row = self._conn.execute( - "SELECT id FROM jobs WHERE status=? ORDER BY created_at LIMIT 1", - (JobStatus.QUEUED,), + "UPDATE jobs SET status=?, started_at=? WHERE id = (" + " SELECT id FROM jobs WHERE status=? ORDER BY created_at LIMIT 1" + ") RETURNING id", + (JobStatus.RUNNING, time.time(), JobStatus.QUEUED), ).fetchone() + self._conn.commit() return row["id"] if row else None def queue_position(self, job_id: str) -> int: diff --git a/app/version.py b/app/version.py index f1380ee..d3ec452 100644 --- a/app/version.py +++ b/app/version.py @@ -1 +1 @@ -__version__ = "0.1.7" +__version__ = "0.2.0" diff --git a/tests/test_concurrency.py b/tests/test_concurrency.py new file mode 100644 index 0000000..45bc05c --- /dev/null +++ b/tests/test_concurrency.py @@ -0,0 +1,97 @@ +"""Тесты параллельной обработки. + +Замеры показали: обе стадии упираются в 4 потока, а на 16 работают вчетверо +медленнее. Значит ядра нужно занимать не шириной одной задачи, а несколькими +задачами сразу - и тогда очередь обязана быть устойчивой к гонкам. +""" +import threading + +import pytest + +from app.store import JobStatus, JobStore + + +@pytest.fixture +def store(tmp_path): + return JobStore(tmp_path / "jobs.db") + + +class TestClaimIsAtomic: + def test_claim_marks_running(self, store): + job_id = store.create(filename="a.wav", duration_sec=1.0) + assert store.claim_next() == job_id + assert store.get(job_id)["status"] == JobStatus.RUNNING + + def test_second_claim_gets_nothing(self, store): + store.create(filename="a.wav", duration_sec=1.0) + store.claim_next() + assert store.claim_next() is None + + def test_each_job_claimed_once_under_load(self, store): + """Главное требование: два воркера не должны взять одну задачу.""" + ids = {store.create(filename=f"{i}.wav", duration_sec=1.0) for i in range(50)} + claimed: list[str] = [] + lock = threading.Lock() + + def worker(): + while True: + job_id = store.claim_next() + if job_id is None: + return + with lock: + claimed.append(job_id) + + threads = [threading.Thread(target=worker) for _ in range(8)] + for t in threads: + t.start() + for t in threads: + t.join() + + assert len(claimed) == len(set(claimed)) == 50 + assert set(claimed) == ids + + def test_claims_oldest_first(self, store): + first = store.create(filename="1.wav", duration_sec=1.0) + store.create(filename="2.wav", duration_sec=1.0) + assert store.claim_next() == first + + +class TestWorkerSettings: + def test_default_threads_is_four_not_all_cores(self, tmp_path): + """Широкие потоки замедляют обе стадии, поэтому по умолчанию их немного.""" + from app.config import load_settings + + config = tmp_path / "config.toml" + config.write_text('[processing]\nthreads=0\n', encoding="utf-8") + assert load_settings(config).effective_threads() == 4 + + def test_explicit_threads_respected(self, tmp_path): + from app.config import load_settings + + config = tmp_path / "config.toml" + config.write_text('[processing]\nthreads=6\n', encoding="utf-8") + assert load_settings(config).effective_threads() == 6 + + def test_workers_derived_from_cores(self, tmp_path, monkeypatch): + from app.config import load_settings + + monkeypatch.setattr("os.cpu_count", lambda: 32) + config = tmp_path / "config.toml" + config.write_text('[processing]\nthreads=4\nworkers=0\n', encoding="utf-8") + # 32 логических ядра при 4 потоках на задачу - но не больше разумного предела + assert 2 <= load_settings(config).effective_workers() <= 4 + + def test_workers_never_below_one(self, tmp_path, monkeypatch): + from app.config import load_settings + + monkeypatch.setattr("os.cpu_count", lambda: 1) + config = tmp_path / "config.toml" + config.write_text('[processing]\nthreads=4\nworkers=0\n', encoding="utf-8") + assert load_settings(config).effective_workers() == 1 + + def test_explicit_workers_respected(self, tmp_path): + from app.config import load_settings + + config = tmp_path / "config.toml" + config.write_text('[processing]\nworkers=3\n', encoding="utf-8") + assert load_settings(config).effective_workers() == 3 diff --git a/tests/test_store.py b/tests/test_store.py index 40364d2..d8e087b 100644 --- a/tests/test_store.py +++ b/tests/test_store.py @@ -51,18 +51,18 @@ class TestJobLifecycle: class TestQueue: - def test_take_next_returns_oldest_queued(self, store): + def test_claim_returns_oldest_queued(self, store): first = store.create(filename="1.wav", duration_sec=1.0) store.create(filename="2.wav", duration_sec=1.0) - assert store.take_next() == first + assert store.claim_next() == first - def test_take_next_skips_running(self, store): + def test_claim_skips_running(self, store): job_id = store.create(filename="1.wav", duration_sec=1.0) store.mark_running(job_id) - assert store.take_next() is None + assert store.claim_next() is None - def test_take_next_on_empty_queue(self, store): - assert store.take_next() is None + def test_claim_on_empty_queue(self, store): + assert store.claim_next() is None def test_queue_position_counts_only_waiting(self, store): a = store.create(filename="1.wav", duration_sec=1.0)