diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py index 943bd84..d4009c4 100644 --- a/iaai_scraper/browser/listing.py +++ b/iaai_scraper/browser/listing.py @@ -81,16 +81,20 @@ class ListingCollector: def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult: # Считываем ссылки одним проходом по DOM. - raw_items = page.eval_on_selector_all( - "a[href*='/VehicleDetail/']", - """ - (nodes) => nodes.map((a) => ({ - href: a.getAttribute('href') || '', - title: a.getAttribute('title') || '', - text: (a.textContent || '').trim(), - })) - """, - ) + try: + raw_items = page.eval_on_selector_all( + "a[href*='/VehicleDetail/']", + """ + (nodes) => nodes.map((a) => ({ + href: a.getAttribute('href') || '', + 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) links: list[ListingVehicleLink] = [] seen: set[str] = set() @@ -127,13 +131,19 @@ class ListingCollector: locator = page.locator(selector).first if locator.count() == 0: continue - disabled = (locator.get_attribute("disabled") or "").lower() - aria_disabled = (locator.get_attribute("aria-disabled") or "").lower() - classes = (locator.get_attribute("class") or "").lower() + try: + disabled = (locator.get_attribute("disabled", timeout=1500) 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: continue - self.pacer.move_mouse_to(page, locator) - locator.click() + try: + self.pacer.move_mouse_to(page, locator) + locator.click(timeout=8000) + except Exception: + continue # Ждём смены контента (AJAX пагинация): первая VehicleDetail-ссылка должна измениться. if old_first_href: @@ -168,15 +178,29 @@ class ListingCollector: make: str | None = None, model: str | None = None, known_origin_ids: set[str] | None = None, + max_duration_seconds: float | None = None, ) -> dict[str, Any]: self.open_cars_listing(page) 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]] = [] all_links: list[str] = [] early_stopped = False threshold = self.settings.listing.early_stop_threshold 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) pages.append({ "page_number": page_result.page_number, @@ -236,6 +260,7 @@ class ListingCollector: "vehicles_collected": len(all_links), "vehicle_urls": all_links, "early_stopped": early_stopped, + "truncated_by_time_budget": truncated_by_time_budget, "pages": pages, "strategy": { "sequential": True, diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 21c3778..1f20fc0 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -263,6 +263,7 @@ class IAAIScraper: make: str | None = None, model: str | None = None, limit: int, + max_duration_seconds: float | None = None, ) -> tuple[list[str], list[str], dict, int]: # Поэтапный сбор листинга до нужного лимита новых URL. new_urls: list[str] = [] @@ -271,6 +272,8 @@ class IAAIScraper: skipped_existing = 0 pages_info: list[dict] = [] 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() try: @@ -278,6 +281,18 @@ class IAAIScraper: self.listing_collector.apply_filters(page, make=make, model=model) 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) pages_info.append({ "page_number": page_result.page_number, @@ -349,6 +364,7 @@ class IAAIScraper: "vehicles_collected": len(all_raw_urls), "vehicle_urls": all_raw_urls, "early_stopped": False, + "truncated_by_time_budget": truncated_by_time_budget, "pages": pages_info, } return new_urls, all_raw_urls, listing, skipped_existing @@ -358,6 +374,7 @@ class IAAIScraper: make: str | None = None, model: str | None = None, known_origin_ids: set[str] | None = None, + max_duration_seconds: float | None = None, ): page = self._get_page_with_warmup() @@ -367,6 +384,7 @@ class IAAIScraper: make=make, model=model, known_origin_ids=known_origin_ids, + max_duration_seconds=max_duration_seconds, ) finally: page.close() @@ -1157,6 +1175,15 @@ class IAAIScraper: try: 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 используем итеративный подход: @@ -1168,16 +1195,38 @@ class IAAIScraper: original_max_pages_per_run = self.settings.listing.max_pages_per_run 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( - make=make, model=model, limit=limit, + make=make, model=model, limit=effective_limit, max_duration_seconds=listing_time_budget_s, ) else: # Обычный сбор (без only_new или без limit). prefetch_cap: int | None = None if limit is not None and limit > 0: 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 # Если включён режим только новых — заранее загружаем все известные origin_id, @@ -1198,6 +1247,7 @@ class IAAIScraper: make=make, model=model, known_origin_ids=known_origin_ids, + max_duration_seconds=listing_time_budget_s, ) finally: self.settings.listing.max_vehicles_per_run = original_listing_cap @@ -1215,6 +1265,11 @@ class IAAIScraper: vehicle_urls = vehicle_urls[:max(0, limit)] 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) # Собираем нормализованные origin_url для последующей пометки проданных. @@ -1227,18 +1282,55 @@ class IAAIScraper: all_listing_origin_urls.add(normalized_url) # Обрабатываем пакетами. + # Каждый batch сохраняется в БД сразу — partial progress не теряется при crash. batch_size = self.settings.celery.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] logger.info( "Processing batch %d-%d of %d", batch_start + 1, min(batch_start + batch_size, total), total, ) - batch_result = self.sync_batch(batch_urls, lane=lane) - cars_upserted += batch_result.get("cars_upserted", 0) - cars_failed += batch_result.get("cars_failed", 0) - images_upserted += batch_result.get("images_upserted", 0) - failures.extend(batch_result.get("failures", [])) + try: + batch_result = self.sync_batch(batch_urls, lane=lane) + cars_upserted += batch_result.get("cars_upserted", 0) + cars_failed += batch_result.get("cars_failed", 0) + 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). @@ -1260,7 +1352,8 @@ class IAAIScraper: except Exception as exc: if not failures: 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: 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