improve listing sync

This commit is contained in:
qananasikq
2026-04-14 14:03:16 +03:00
parent 9337344182
commit d51216ddb7
2 changed files with 141 additions and 23 deletions

View File

@@ -81,6 +81,7 @@ 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.
try:
raw_items = page.eval_on_selector_all( raw_items = page.eval_on_selector_all(
"a[href*='/VehicleDetail/']", "a[href*='/VehicleDetail/']",
""" """
@@ -91,6 +92,9 @@ class ListingCollector:
})) }))
""", """,
) )
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
try:
self.pacer.move_mouse_to(page, locator) self.pacer.move_mouse_to(page, locator)
locator.click() 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,

View File

@@ -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,
) )
try:
batch_result = self.sync_batch(batch_urls, lane=lane) batch_result = self.sync_batch(batch_urls, lane=lane)
cars_upserted += batch_result.get("cars_upserted", 0) cars_upserted += batch_result.get("cars_upserted", 0)
cars_failed += batch_result.get("cars_failed", 0) cars_failed += batch_result.get("cars_failed", 0)
images_upserted += batch_result.get("images_upserted", 0) images_upserted += batch_result.get("images_upserted", 0)
failures.extend(batch_result.get("failures", [])) 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