Compare commits
3 Commits
c729f971e3
...
6a334d3c64
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6a334d3c64 | ||
|
|
d51216ddb7 | ||
|
|
9337344182 |
@@ -141,9 +141,9 @@ services:
|
|||||||
stop_grace_period: 60s
|
stop_grace_period: 60s
|
||||||
command: >
|
command: >
|
||||||
celery -A iaai_scraper.worker.celery_app worker
|
celery -A iaai_scraper.worker.celery_app worker
|
||||||
--loglevel=info --concurrency=4 --pool=prefork
|
--loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-4} --pool=prefork
|
||||||
--pidfile=/tmp/celery-worker.pid
|
--pidfile=/tmp/celery-worker.pid
|
||||||
-Q scraping --max-tasks-per-child=5
|
-Q scraping --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5}
|
||||||
healthcheck:
|
healthcheck:
|
||||||
test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid)"]
|
test: ["CMD-SHELL", "test -f /tmp/celery-worker.pid && kill -0 $(cat /tmp/celery-worker.pid)"]
|
||||||
interval: 60s
|
interval: 60s
|
||||||
|
|||||||
@@ -81,16 +81,20 @@ class ListingCollector:
|
|||||||
|
|
||||||
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
|
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
|
||||||
# Считываем ссылки одним проходом по DOM.
|
# Считываем ссылки одним проходом по DOM.
|
||||||
raw_items = page.eval_on_selector_all(
|
try:
|
||||||
"a[href*='/VehicleDetail/']",
|
raw_items = page.eval_on_selector_all(
|
||||||
"""
|
"a[href*='/VehicleDetail/']",
|
||||||
(nodes) => nodes.map((a) => ({
|
"""
|
||||||
href: a.getAttribute('href') || '',
|
(nodes) => nodes.map((a) => ({
|
||||||
title: a.getAttribute('title') || '',
|
href: a.getAttribute('href') || '',
|
||||||
text: (a.textContent || '').trim(),
|
title: a.getAttribute('title') || '',
|
||||||
}))
|
text: (a.textContent || '').trim(),
|
||||||
""",
|
}))
|
||||||
)
|
""",
|
||||||
|
)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("collect_current_page failed on page %d: %s", page_number, exc)
|
||||||
|
raw_items = []
|
||||||
total = min(len(raw_items), self.settings.listing.page_link_limit)
|
total = min(len(raw_items), self.settings.listing.page_link_limit)
|
||||||
links: list[ListingVehicleLink] = []
|
links: list[ListingVehicleLink] = []
|
||||||
seen: set[str] = set()
|
seen: set[str] = set()
|
||||||
@@ -127,13 +131,19 @@ class ListingCollector:
|
|||||||
locator = page.locator(selector).first
|
locator = page.locator(selector).first
|
||||||
if locator.count() == 0:
|
if locator.count() == 0:
|
||||||
continue
|
continue
|
||||||
disabled = (locator.get_attribute("disabled") or "").lower()
|
try:
|
||||||
aria_disabled = (locator.get_attribute("aria-disabled") or "").lower()
|
disabled = (locator.get_attribute("disabled", timeout=1500) or "").lower()
|
||||||
classes = (locator.get_attribute("class") or "").lower()
|
aria_disabled = (locator.get_attribute("aria-disabled", timeout=1500) or "").lower()
|
||||||
|
classes = (locator.get_attribute("class", timeout=1500) or "").lower()
|
||||||
|
except Exception:
|
||||||
|
continue
|
||||||
if disabled or aria_disabled == "true" or "disabled" in classes:
|
if disabled or aria_disabled == "true" or "disabled" in classes:
|
||||||
continue
|
continue
|
||||||
self.pacer.move_mouse_to(page, locator)
|
try:
|
||||||
locator.click()
|
self.pacer.move_mouse_to(page, locator)
|
||||||
|
locator.click(timeout=8000)
|
||||||
|
except Exception:
|
||||||
|
continue
|
||||||
|
|
||||||
# Ждём смены контента (AJAX пагинация): первая VehicleDetail-ссылка должна измениться.
|
# Ждём смены контента (AJAX пагинация): первая VehicleDetail-ссылка должна измениться.
|
||||||
if old_first_href:
|
if old_first_href:
|
||||||
@@ -168,15 +178,29 @@ class ListingCollector:
|
|||||||
make: str | None = None,
|
make: str | None = None,
|
||||||
model: str | None = None,
|
model: str | None = None,
|
||||||
known_origin_ids: set[str] | None = None,
|
known_origin_ids: set[str] | None = None,
|
||||||
|
max_duration_seconds: float | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
self.open_cars_listing(page)
|
self.open_cars_listing(page)
|
||||||
applied_filters = self.apply_filters(page, make=make, model=model)
|
applied_filters = self.apply_filters(page, make=make, model=model)
|
||||||
|
started_at = time.perf_counter()
|
||||||
|
truncated_by_time_budget = False
|
||||||
pages: list[dict[str, object]] = []
|
pages: list[dict[str, object]] = []
|
||||||
all_links: list[str] = []
|
all_links: list[str] = []
|
||||||
early_stopped = False
|
early_stopped = False
|
||||||
threshold = self.settings.listing.early_stop_threshold
|
threshold = self.settings.listing.early_stop_threshold
|
||||||
|
|
||||||
for page_number in range(1, max(1, self.settings.listing.max_pages_per_run) + 1):
|
for page_number in range(1, max(1, self.settings.listing.max_pages_per_run) + 1):
|
||||||
|
if max_duration_seconds is not None and max_duration_seconds > 0:
|
||||||
|
elapsed = time.perf_counter() - started_at
|
||||||
|
if elapsed >= max_duration_seconds:
|
||||||
|
truncated_by_time_budget = True
|
||||||
|
logger.warning(
|
||||||
|
"Listing collection stopped by time budget: page=%d elapsed=%.1fs budget=%.1fs",
|
||||||
|
page_number,
|
||||||
|
elapsed,
|
||||||
|
max_duration_seconds,
|
||||||
|
)
|
||||||
|
break
|
||||||
page_result = self.collect_current_page(page, page_number=page_number)
|
page_result = self.collect_current_page(page, page_number=page_number)
|
||||||
pages.append({
|
pages.append({
|
||||||
"page_number": page_result.page_number,
|
"page_number": page_result.page_number,
|
||||||
@@ -236,6 +260,7 @@ class ListingCollector:
|
|||||||
"vehicles_collected": len(all_links),
|
"vehicles_collected": len(all_links),
|
||||||
"vehicle_urls": all_links,
|
"vehicle_urls": all_links,
|
||||||
"early_stopped": early_stopped,
|
"early_stopped": early_stopped,
|
||||||
|
"truncated_by_time_budget": truncated_by_time_budget,
|
||||||
"pages": pages,
|
"pages": pages,
|
||||||
"strategy": {
|
"strategy": {
|
||||||
"sequential": True,
|
"sequential": True,
|
||||||
|
|||||||
@@ -263,6 +263,7 @@ class IAAIScraper:
|
|||||||
make: str | None = None,
|
make: str | None = None,
|
||||||
model: str | None = None,
|
model: str | None = None,
|
||||||
limit: int,
|
limit: int,
|
||||||
|
max_duration_seconds: float | None = None,
|
||||||
) -> tuple[list[str], list[str], dict, int]:
|
) -> tuple[list[str], list[str], dict, int]:
|
||||||
# Поэтапный сбор листинга до нужного лимита новых URL.
|
# Поэтапный сбор листинга до нужного лимита новых URL.
|
||||||
new_urls: list[str] = []
|
new_urls: list[str] = []
|
||||||
@@ -271,6 +272,8 @@ class IAAIScraper:
|
|||||||
skipped_existing = 0
|
skipped_existing = 0
|
||||||
pages_info: list[dict] = []
|
pages_info: list[dict] = []
|
||||||
max_pages = self.settings.listing.max_pages_per_run
|
max_pages = self.settings.listing.max_pages_per_run
|
||||||
|
started_at = time.perf_counter()
|
||||||
|
truncated_by_time_budget = False
|
||||||
|
|
||||||
page = self._get_page_with_warmup()
|
page = self._get_page_with_warmup()
|
||||||
try:
|
try:
|
||||||
@@ -278,6 +281,18 @@ class IAAIScraper:
|
|||||||
self.listing_collector.apply_filters(page, make=make, model=model)
|
self.listing_collector.apply_filters(page, make=make, model=model)
|
||||||
|
|
||||||
for page_number in range(1, max_pages + 1):
|
for page_number in range(1, max_pages + 1):
|
||||||
|
if max_duration_seconds is not None and max_duration_seconds > 0:
|
||||||
|
elapsed = time.perf_counter() - started_at
|
||||||
|
if elapsed >= max_duration_seconds:
|
||||||
|
truncated_by_time_budget = True
|
||||||
|
logger.warning(
|
||||||
|
"Iterative listing stopped by time budget: page=%d elapsed=%.1fs budget=%.1fs",
|
||||||
|
page_number,
|
||||||
|
elapsed,
|
||||||
|
max_duration_seconds,
|
||||||
|
)
|
||||||
|
break
|
||||||
|
|
||||||
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
|
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
|
||||||
pages_info.append({
|
pages_info.append({
|
||||||
"page_number": page_result.page_number,
|
"page_number": page_result.page_number,
|
||||||
@@ -349,6 +364,7 @@ class IAAIScraper:
|
|||||||
"vehicles_collected": len(all_raw_urls),
|
"vehicles_collected": len(all_raw_urls),
|
||||||
"vehicle_urls": all_raw_urls,
|
"vehicle_urls": all_raw_urls,
|
||||||
"early_stopped": False,
|
"early_stopped": False,
|
||||||
|
"truncated_by_time_budget": truncated_by_time_budget,
|
||||||
"pages": pages_info,
|
"pages": pages_info,
|
||||||
}
|
}
|
||||||
return new_urls, all_raw_urls, listing, skipped_existing
|
return new_urls, all_raw_urls, listing, skipped_existing
|
||||||
@@ -358,6 +374,7 @@ class IAAIScraper:
|
|||||||
make: str | None = None,
|
make: str | None = None,
|
||||||
model: str | None = None,
|
model: str | None = None,
|
||||||
known_origin_ids: set[str] | None = None,
|
known_origin_ids: set[str] | None = None,
|
||||||
|
max_duration_seconds: float | None = None,
|
||||||
):
|
):
|
||||||
page = self._get_page_with_warmup()
|
page = self._get_page_with_warmup()
|
||||||
|
|
||||||
@@ -367,6 +384,7 @@ class IAAIScraper:
|
|||||||
make=make,
|
make=make,
|
||||||
model=model,
|
model=model,
|
||||||
known_origin_ids=known_origin_ids,
|
known_origin_ids=known_origin_ids,
|
||||||
|
max_duration_seconds=max_duration_seconds,
|
||||||
)
|
)
|
||||||
finally:
|
finally:
|
||||||
page.close()
|
page.close()
|
||||||
@@ -1157,6 +1175,15 @@ class IAAIScraper:
|
|||||||
|
|
||||||
try:
|
try:
|
||||||
effective_only_new = self.settings.sync_only_new if only_new is None else only_new
|
effective_only_new = self.settings.sync_only_new if only_new is None else only_new
|
||||||
|
soft_time_limit = self.settings.celery.task_soft_time_limit
|
||||||
|
listing_time_budget_s: float | None = None
|
||||||
|
if soft_time_limit and soft_time_limit > 0:
|
||||||
|
# Делим время: ~40% на листинг, ~60% на парсинг + запись.
|
||||||
|
# При soft=120 → listing=48s, processing=72s (достаточно для ~7 батчей).
|
||||||
|
# При soft=600 → listing=240s, processing=360s (достаточно для ~30+ батчей).
|
||||||
|
# Минимум 30s на листинг, оставляем 30s запас до hard limit.
|
||||||
|
soft_f = float(soft_time_limit)
|
||||||
|
listing_time_budget_s = max(30.0, soft_f * 0.4)
|
||||||
|
|
||||||
# ── Сбор листинга ──
|
# ── Сбор листинга ──
|
||||||
# При only_new + limit используем итеративный подход:
|
# При only_new + limit используем итеративный подход:
|
||||||
@@ -1168,16 +1195,38 @@ class IAAIScraper:
|
|||||||
original_max_pages_per_run = self.settings.listing.max_pages_per_run
|
original_max_pages_per_run = self.settings.listing.max_pages_per_run
|
||||||
|
|
||||||
if effective_only_new and limit is not None and limit > 0:
|
if effective_only_new and limit is not None and limit > 0:
|
||||||
|
# Ограничиваем кол-во URL тем, что реально успеем обработать.
|
||||||
|
# При soft=600: processing=360s, ~3s на батч из 10 → ~120 батчей → ~1200 машин.
|
||||||
|
effective_limit = limit
|
||||||
|
if soft_time_limit and soft_time_limit > 0 and listing_time_budget_s is not None:
|
||||||
|
processing_budget_s = max(60.0, float(soft_time_limit) - listing_time_budget_s)
|
||||||
|
seconds_per_batch = 4.0 # ~3-4s на batch из 10 через HTTP fast-path
|
||||||
|
batches_fit = max(1, int(processing_budget_s / seconds_per_batch))
|
||||||
|
max_processable = batches_fit * self.settings.celery.batch_size
|
||||||
|
if max_processable < effective_limit:
|
||||||
|
logger.info(
|
||||||
|
"Capping iterative limit %d → %d (processing_budget=%.0fs, batch_size=%d)",
|
||||||
|
effective_limit, max_processable, processing_budget_s, self.settings.celery.batch_size,
|
||||||
|
)
|
||||||
|
effective_limit = max_processable
|
||||||
|
|
||||||
# Итеративный сбор: страница→фильтр→проверка→следующая страница.
|
# Итеративный сбор: страница→фильтр→проверка→следующая страница.
|
||||||
vehicle_urls, raw_urls, listing, skipped_existing = self._collect_listing_iterative(
|
vehicle_urls, raw_urls, listing, skipped_existing = self._collect_listing_iterative(
|
||||||
make=make, model=model, limit=limit,
|
make=make, model=model, limit=effective_limit, max_duration_seconds=listing_time_budget_s,
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
# Обычный сбор (без only_new или без limit).
|
# Обычный сбор (без only_new или без limit).
|
||||||
prefetch_cap: int | None = None
|
prefetch_cap: int | None = None
|
||||||
if limit is not None and limit > 0:
|
if limit is not None and limit > 0:
|
||||||
prefetch_cap = limit
|
prefetch_cap = limit
|
||||||
if prefetch_cap < original_listing_cap:
|
elif soft_time_limit and soft_time_limit > 0 and listing_time_budget_s is not None:
|
||||||
|
# При ограниченном soft-time-limit не раздуваем prefetch,
|
||||||
|
# чтобы быстрее перейти к парсингу/записи.
|
||||||
|
remaining_processing_s = max(30.0, float(soft_time_limit) - listing_time_budget_s)
|
||||||
|
batches_fit = max(2, int(remaining_processing_s // 60))
|
||||||
|
prefetch_cap = max(self.settings.celery.batch_size, self.settings.celery.batch_size * batches_fit)
|
||||||
|
|
||||||
|
if prefetch_cap is not None and prefetch_cap < original_listing_cap:
|
||||||
self.settings.listing.max_vehicles_per_run = prefetch_cap
|
self.settings.listing.max_vehicles_per_run = prefetch_cap
|
||||||
|
|
||||||
# Если включён режим только новых — заранее загружаем все известные origin_id,
|
# Если включён режим только новых — заранее загружаем все известные origin_id,
|
||||||
@@ -1198,6 +1247,7 @@ class IAAIScraper:
|
|||||||
make=make,
|
make=make,
|
||||||
model=model,
|
model=model,
|
||||||
known_origin_ids=known_origin_ids,
|
known_origin_ids=known_origin_ids,
|
||||||
|
max_duration_seconds=listing_time_budget_s,
|
||||||
)
|
)
|
||||||
finally:
|
finally:
|
||||||
self.settings.listing.max_vehicles_per_run = original_listing_cap
|
self.settings.listing.max_vehicles_per_run = original_listing_cap
|
||||||
@@ -1215,6 +1265,11 @@ class IAAIScraper:
|
|||||||
vehicle_urls = vehicle_urls[:max(0, limit)]
|
vehicle_urls = vehicle_urls[:max(0, limit)]
|
||||||
|
|
||||||
total = len(vehicle_urls)
|
total = len(vehicle_urls)
|
||||||
|
if listing.get("truncated_by_time_budget"):
|
||||||
|
logger.warning(
|
||||||
|
"Listing was truncated by time budget; processing collected subset (%d URLs)",
|
||||||
|
total,
|
||||||
|
)
|
||||||
logger.info("Starting sync: %d vehicles to process (batch mode)", total)
|
logger.info("Starting sync: %d vehicles to process (batch mode)", total)
|
||||||
|
|
||||||
# Собираем нормализованные origin_url для последующей пометки проданных.
|
# Собираем нормализованные origin_url для последующей пометки проданных.
|
||||||
@@ -1227,18 +1282,55 @@ class IAAIScraper:
|
|||||||
all_listing_origin_urls.add(normalized_url)
|
all_listing_origin_urls.add(normalized_url)
|
||||||
|
|
||||||
# Обрабатываем пакетами.
|
# Обрабатываем пакетами.
|
||||||
|
# Каждый batch сохраняется в БД сразу — partial progress не теряется при crash.
|
||||||
batch_size = self.settings.celery.batch_size
|
batch_size = self.settings.celery.batch_size
|
||||||
for batch_start in range(0, total, batch_size):
|
for batch_start in range(0, total, batch_size):
|
||||||
|
# Проверяем оставшееся время перед каждым batch.
|
||||||
|
if soft_time_limit and soft_time_limit > 0:
|
||||||
|
elapsed = time.perf_counter() - started_at
|
||||||
|
remaining = float(soft_time_limit) - elapsed
|
||||||
|
if remaining < 30:
|
||||||
|
logger.warning(
|
||||||
|
"Time budget exhausted before batch %d-%d (elapsed=%.1fs, soft=%ds). "
|
||||||
|
"Saving partial progress: %d upserted so far.",
|
||||||
|
batch_start + 1, min(batch_start + batch_size, total),
|
||||||
|
elapsed, soft_time_limit, cars_upserted,
|
||||||
|
)
|
||||||
|
break
|
||||||
|
|
||||||
batch_urls = vehicle_urls[batch_start:batch_start + batch_size]
|
batch_urls = vehicle_urls[batch_start:batch_start + batch_size]
|
||||||
logger.info(
|
logger.info(
|
||||||
"Processing batch %d-%d of %d",
|
"Processing batch %d-%d of %d",
|
||||||
batch_start + 1, min(batch_start + batch_size, total), total,
|
batch_start + 1, min(batch_start + batch_size, total), total,
|
||||||
)
|
)
|
||||||
batch_result = self.sync_batch(batch_urls, lane=lane)
|
try:
|
||||||
cars_upserted += batch_result.get("cars_upserted", 0)
|
batch_result = self.sync_batch(batch_urls, lane=lane)
|
||||||
cars_failed += batch_result.get("cars_failed", 0)
|
cars_upserted += batch_result.get("cars_upserted", 0)
|
||||||
images_upserted += batch_result.get("images_upserted", 0)
|
cars_failed += batch_result.get("cars_failed", 0)
|
||||||
failures.extend(batch_result.get("failures", []))
|
images_upserted += batch_result.get("images_upserted", 0)
|
||||||
|
failures.extend(batch_result.get("failures", []))
|
||||||
|
except PlaywrightError as pw_exc:
|
||||||
|
# Browser crashed — пробуем пересоздать контекст и продолжить.
|
||||||
|
logger.warning(
|
||||||
|
"Batch %d-%d browser crashed: %s. Reinitializing browser context.",
|
||||||
|
batch_start + 1, min(batch_start + batch_size, total), pw_exc,
|
||||||
|
)
|
||||||
|
cars_failed += len(batch_urls)
|
||||||
|
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(pw_exc)})
|
||||||
|
try:
|
||||||
|
self._new_context()
|
||||||
|
logger.info("Browser context reinitialized successfully after crash")
|
||||||
|
except Exception as reinit_exc:
|
||||||
|
logger.error("Failed to reinitialize browser: %s. Aborting remaining batches.", reinit_exc)
|
||||||
|
break
|
||||||
|
except Exception as batch_exc:
|
||||||
|
# Batch failed — логируем и не теряем уже накопленный прогресс.
|
||||||
|
logger.error(
|
||||||
|
"Batch %d-%d failed: %s. Continuing with next batch.",
|
||||||
|
batch_start + 1, min(batch_start + batch_size, total), batch_exc,
|
||||||
|
)
|
||||||
|
cars_failed += len(batch_urls)
|
||||||
|
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
|
||||||
|
|
||||||
# Помечаем авто как проданные, если они исчезли из листинга.
|
# Помечаем авто как проданные, если они исчезли из листинга.
|
||||||
# Только если сканирование было полным (не ограниченным limit/only_new/early_stop).
|
# Только если сканирование было полным (не ограниченным limit/only_new/early_stop).
|
||||||
@@ -1260,7 +1352,8 @@ class IAAIScraper:
|
|||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
if not failures:
|
if not failures:
|
||||||
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
|
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
|
||||||
logger.error("sync_listing failed: %s", exc)
|
logger.error("sync_listing failed: %s (partial progress: %d upserted)", exc, cars_upserted)
|
||||||
|
# Не пробрасываем — partial progress уже записан в БД через finish_sync_run ниже.
|
||||||
finally:
|
finally:
|
||||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||||
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
|
error_summary = "; ".join(item["error"] for item in failures[:10]) if failures else None
|
||||||
|
|||||||
@@ -1,11 +1,15 @@
|
|||||||
# Инициализация Celery-приложения и периодических задач.
|
# Инициализация Celery-приложения и периодических задач.
|
||||||
|
|
||||||
|
import logging
|
||||||
|
|
||||||
from celery import Celery
|
from celery import Celery
|
||||||
from celery.signals import worker_process_init, setup_logging as celery_setup_logging
|
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
|
||||||
|
|
||||||
from ..core.config import settings
|
from ..core.config import settings
|
||||||
from ..core.logs import setup_logging
|
from ..core.logs import setup_logging
|
||||||
|
|
||||||
|
logger = logging.getLogger("iaai_scraper.worker.celery_app")
|
||||||
|
|
||||||
|
|
||||||
@celery_setup_logging.connect
|
@celery_setup_logging.connect
|
||||||
def _configure_logging(loglevel=None, **kwargs):
|
def _configure_logging(loglevel=None, **kwargs):
|
||||||
@@ -36,14 +40,27 @@ celery_app = Celery(
|
|||||||
backend=_result_backend(),
|
backend=_result_backend(),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Auto-clamp: если hard limit слишком далёк от soft (> soft + 120),
|
||||||
|
# ограничиваем, чтобы зависший worker не жил вечно.
|
||||||
|
_soft = settings.celery.task_soft_time_limit
|
||||||
|
_hard = settings.celery.task_time_limit
|
||||||
|
_max_hard = _soft + 120 if _soft else _hard
|
||||||
|
if _hard > _max_hard:
|
||||||
|
logger.warning(
|
||||||
|
"CELERY_TASK_TIME_LIMIT=%d too far from CELERY_TASK_SOFT_TIME_LIMIT=%d; "
|
||||||
|
"clamping hard limit to %d",
|
||||||
|
_hard, _soft, _max_hard,
|
||||||
|
)
|
||||||
|
_hard = _max_hard
|
||||||
|
|
||||||
celery_app.conf.update(
|
celery_app.conf.update(
|
||||||
task_serializer="json",
|
task_serializer="json",
|
||||||
accept_content=["json"],
|
accept_content=["json"],
|
||||||
result_serializer="json",
|
result_serializer="json",
|
||||||
timezone="UTC",
|
timezone="UTC",
|
||||||
enable_utc=True,
|
enable_utc=True,
|
||||||
task_soft_time_limit=settings.celery.task_soft_time_limit,
|
task_soft_time_limit=_soft,
|
||||||
task_time_limit=settings.celery.task_time_limit,
|
task_time_limit=_hard,
|
||||||
task_acks_late=True,
|
task_acks_late=True,
|
||||||
task_reject_on_worker_lost=True,
|
task_reject_on_worker_lost=True,
|
||||||
task_track_started=True,
|
task_track_started=True,
|
||||||
@@ -73,3 +90,15 @@ celery_app.conf.update(
|
|||||||
)
|
)
|
||||||
|
|
||||||
celery_app.autodiscover_tasks(["iaai_scraper.worker"])
|
celery_app.autodiscover_tasks(["iaai_scraper.worker"])
|
||||||
|
|
||||||
|
|
||||||
|
@worker_ready.connect
|
||||||
|
def _on_worker_ready(**kwargs):
|
||||||
|
"""Сразу при старте worker отправляем первую задачу sync_listing,
|
||||||
|
чтобы не ждать час до первого beat-цикла."""
|
||||||
|
logger.info("Worker ready — dispatching initial sync_listing task")
|
||||||
|
celery_app.send_task(
|
||||||
|
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||||
|
kwargs={"limit": settings.celery.beat_sync_limit},
|
||||||
|
queue="scraping",
|
||||||
|
)
|
||||||
@@ -6,6 +6,7 @@ from threading import Event, Thread
|
|||||||
import time
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
|
|
||||||
|
from billiard.exceptions import SoftTimeLimitExceeded
|
||||||
from celery import shared_task
|
from celery import shared_task
|
||||||
from redis import Redis
|
from redis import Redis
|
||||||
|
|
||||||
@@ -42,15 +43,44 @@ def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
|||||||
|
|
||||||
def _run_browser_job(func, *args, **kwargs):
|
def _run_browser_job(func, *args, **kwargs):
|
||||||
# Браузерный код запускаем в отдельном потоке без активного loop.
|
# Браузерный код запускаем в отдельном потоке без активного loop.
|
||||||
with ThreadPoolExecutor(max_workers=1, thread_name_prefix="iaai-browser") as executor:
|
# НЕ используем `with` — иначе shutdown(wait=True) заблокирует main thread
|
||||||
future = executor.submit(func, *args, **kwargs)
|
# если SoftTimeLimitExceeded прервёт future.result(), а browser thread ещё работает.
|
||||||
return future.result()
|
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
|
||||||
|
|
||||||
|
|
||||||
def _sync_listing_lock_ttl_seconds() -> int:
|
def _sync_listing_lock_ttl_seconds() -> int:
|
||||||
settings = Settings()
|
settings = Settings()
|
||||||
# Небольшой запас к лимиту времени, чтобы lock снимался после сбоев.
|
soft = settings.celery.task_soft_time_limit
|
||||||
return max(settings.celery.task_time_limit + 120, 300)
|
hard = settings.celery.task_time_limit
|
||||||
|
# Используем clamped hard limit (soft + 120), а не сырой task_time_limit,
|
||||||
|
# чтобы lock не висел 11 дней при CELERY_TASK_TIME_LIMIT=999999.
|
||||||
|
effective_hard = min(hard, soft + 120) if soft else hard
|
||||||
|
return max(effective_hard + 120, 300)
|
||||||
|
|
||||||
|
|
||||||
def _get_persistence() -> PersistenceService:
|
def _get_persistence() -> PersistenceService:
|
||||||
@@ -187,7 +217,7 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
|
|||||||
@shared_task(
|
@shared_task(
|
||||||
name="iaai_scraper.worker.tasks.sync_listing_task",
|
name="iaai_scraper.worker.tasks.sync_listing_task",
|
||||||
bind=True,
|
bind=True,
|
||||||
max_retries=1,
|
max_retries=3,
|
||||||
default_retry_delay=120,
|
default_retry_delay=120,
|
||||||
acks_late=True,
|
acks_late=True,
|
||||||
)
|
)
|
||||||
@@ -244,21 +274,48 @@ def sync_listing_task(
|
|||||||
summary = {
|
summary = {
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
"run_id": result.get("run_id"),
|
"run_id": result.get("run_id"),
|
||||||
|
"status": result.get("status", "success"),
|
||||||
"cars_upserted": result.get("cars_upserted", 0),
|
"cars_upserted": result.get("cars_upserted", 0),
|
||||||
"cars_failed": result.get("cars_failed", 0),
|
"cars_failed": result.get("cars_failed", 0),
|
||||||
"images_upserted": result.get("images_upserted", 0),
|
"images_upserted": result.get("images_upserted", 0),
|
||||||
"skipped_existing": result.get("skipped_existing", 0),
|
"skipped_existing": result.get("skipped_existing", 0),
|
||||||
"elapsed_seconds": result.get("elapsed_seconds"),
|
"elapsed_seconds": result.get("elapsed_seconds"),
|
||||||
|
"failures_count": len(result.get("failures") or []),
|
||||||
}
|
}
|
||||||
logger.info(
|
logger.info(
|
||||||
"sync_listing_task completed: %d upserted, %d failed",
|
"sync_listing_task completed: status=%s, %d upserted, %d failed, failures=%d",
|
||||||
summary["cars_upserted"], summary["cars_failed"],
|
summary["status"],
|
||||||
|
summary["cars_upserted"],
|
||||||
|
summary["cars_failed"],
|
||||||
|
summary["failures_count"],
|
||||||
)
|
)
|
||||||
return {"status": "success", **summary}
|
return summary
|
||||||
|
|
||||||
|
except SoftTimeLimitExceeded:
|
||||||
|
logger.warning(
|
||||||
|
"sync_listing_task soft timeout exceeded — partial progress already saved to DB"
|
||||||
|
)
|
||||||
|
# Partial progress уже записан в БД через finish_sync_run.
|
||||||
|
# Не retry — следующий beat подхватит новые URL автоматически.
|
||||||
|
return {
|
||||||
|
"status": "timed_out",
|
||||||
|
"task_id": task_id,
|
||||||
|
"reason": "soft_time_limit_exceeded",
|
||||||
|
"note": "partial progress saved to DB; next beat cycle will continue",
|
||||||
|
}
|
||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
|
logger.error("sync_listing_task failed: %s", exc, exc_info=True)
|
||||||
raise self.retry(exc=exc)
|
# Retry только на не-таймаутные ошибки (сеть, БД, браузер).
|
||||||
|
try:
|
||||||
|
raise self.retry(exc=exc)
|
||||||
|
except self.MaxRetriesExceededError:
|
||||||
|
logger.error("sync_listing_task max retries exceeded, giving up")
|
||||||
|
return {
|
||||||
|
"status": "failed",
|
||||||
|
"task_id": task_id,
|
||||||
|
"error": str(exc),
|
||||||
|
}
|
||||||
finally:
|
finally:
|
||||||
if heartbeat_stop is not None:
|
if heartbeat_stop is not None:
|
||||||
heartbeat_stop.set()
|
heartbeat_stop.set()
|
||||||
|
|||||||
Reference in New Issue
Block a user