From 133c125e9c94d224fef59a8d8e74fe7fbe093f06 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Mon, 20 Apr 2026 17:33:51 +0300 Subject: [PATCH] fix: eliminate greenlet crash + memory leak protection - docker-compose: switch worker to --pool=solo --concurrency=1 --max-tasks-per-child=1 - tasks.py: remove ThreadPoolExecutor wrapper from _run_browser_job (greenlet crash) - scraper.py: browser fallback runs sequentially instead of ThreadPoolExecutor - scraper.py: add gc.collect() after scraper close to prevent memory leaks --- docker-compose.yml | 4 ++-- iaai_scraper/scraper.py | 22 +++++++++------------- iaai_scraper/worker/tasks.py | 35 +++++------------------------------ 3 files changed, 16 insertions(+), 45 deletions(-) diff --git a/docker-compose.yml b/docker-compose.yml index b7f1585..b49cb53 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -144,9 +144,9 @@ services: stop_grace_period: 60s command: > celery -A iaai_scraper.worker.celery_app worker - --loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-2} --pool=prefork + --loglevel=info --concurrency=1 --pool=solo --pidfile=/tmp/celery-worker.pid - -Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-3} + -Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-1} healthcheck: test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid)"] interval: 60s diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index ba9caa5..69a8977 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -1,3 +1,4 @@ +import gc import json import logging import os @@ -201,6 +202,7 @@ class IAAIScraper: pass finally: self._http_pool = None + gc.collect() def _new_context(self) -> BrowserContext: if self.browser is None: @@ -1623,23 +1625,17 @@ class IAAIScraper: "protection_events": local_protection, } - # Одна страница — в главном потоке. - if n_pages == 1: - res = _process_slice(slices[0] if slices else []) + # Playwright sync API использует greenlets, привязанные к потоку. + # Запуск в ThreadPoolExecutor вызывает greenlet crash. + # Обрабатываем все слайсы последовательно в текущем потоке. + for s in slices: + if not s: + continue + res = _process_slice(s) records.extend(res["records"]) failures.extend(res["failures"]) cars_failed += res["cars_failed"] protection_events += res["protection_events"] - else: - # Несколько страниц — параллельно, каждый поток со своей Page. - with ThreadPoolExecutor(max_workers=n_pages) as executor: - futures = [executor.submit(_process_slice, s) for s in slices if s] - for fut in as_completed(futures): - res = fut.result() - records.extend(res["records"]) - failures.extend(res["failures"]) - cars_failed += res["cars_failed"] - protection_events += res["protection_events"] return { "records": records, diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 581e4f5..d67185e 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -1,6 +1,5 @@ # Задачи Celery для синхронизации автомобилей и листинга IAAI. -from concurrent.futures import ThreadPoolExecutor import json import logging import os @@ -61,35 +60,11 @@ def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0): def _run_browser_job(func, *args, **kwargs): - # Браузерный код запускаем в отдельном потоке без активного loop. - # НЕ используем `with` — иначе shutdown(wait=True) заблокирует main thread - # если SoftTimeLimitExceeded прервёт future.result(), а browser thread ещё работает. - settings = Settings() - soft = settings.celery.task_soft_time_limit - hard = settings.celery.task_time_limit - # Таймаут для future.result(): берём hard limit (или soft + 120), чтобы не висеть вечно. - wait_timeout = min(hard, soft + 120) if soft and hard else None - - executor = ThreadPoolExecutor(max_workers=1, thread_name_prefix="iaai-browser") - future = executor.submit(func, *args, **kwargs) - try: - result = future.result(timeout=wait_timeout) - except SoftTimeLimitExceeded: - # Отпускаем executor без ожидания — thread умрёт когда Celery убьёт процесс (hard limit). - future.cancel() - executor.shutdown(wait=False, cancel_futures=True) - raise - except TimeoutError: - # future.result(timeout=...) вышел по таймауту — browser thread завис. - future.cancel() - executor.shutdown(wait=False, cancel_futures=True) - raise SoftTimeLimitExceeded("Browser thread did not finish within time limit") - except Exception: - executor.shutdown(wait=False, cancel_futures=True) - raise - else: - executor.shutdown(wait=True) - return result + # С pool=solo Celery worker работает в одном процессе/потоке. + # Playwright sync API использует greenlets, которые привязаны к потоку. + # Запуск в отдельном потоке вызывает greenlet.error: cannot switch to a different thread. + # Поэтому запускаем напрямую в текущем потоке. + return func(*args, **kwargs) def _sync_listing_lock_ttl_seconds() -> int: