diff --git a/docker-compose.yml b/docker-compose.yml index fe72c94..dbe9cc7 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -15,6 +15,7 @@ x-app-env: &app-env CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-200} IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-40} IAAI_BLOCK_RESOURCES: ${IAAI_BLOCK_RESOURCES:-true} + IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-auto} IAAI_MAX_PAGES_PER_RUN: ${IAAI_MAX_PAGES_PER_RUN:-9999} IAAI_MAX_VEHICLES_PER_RUN: ${IAAI_MAX_VEHICLES_PER_RUN:-50000} IAAI_HUMAN_PACE_ENABLED: ${IAAI_HUMAN_PACE_ENABLED:-true} diff --git a/iaai_scraper/browser/listing.py b/iaai_scraper/browser/listing.py index fe2d20b..f1c8b59 100644 --- a/iaai_scraper/browser/listing.py +++ b/iaai_scraper/browser/listing.py @@ -133,9 +133,9 @@ class ListingCollector: return not old_first_href - def open_cars_listing(self, page: Page) -> None: - logger.info("Opening cars listing page: %s", self.settings.listing.cars_url) - url = self.settings.listing.cars_url + def open_cars_listing(self, page: Page, *, url_override: str | None = None) -> None: + url = url_override or self.settings.listing.cars_url + logger.info("Opening cars listing page: %s", url) last_err = None for attempt in range(3): try: @@ -202,16 +202,82 @@ class ListingCollector: # Короткая пауза вместо длинного sleep. time.sleep(1.0) - def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]: - applied = {"make": None, "model": None} + def apply_filters( + self, + page: Page, + make: str | None = None, + model: str | None = None, + year_min: int | None = None, + year_max: int | None = None, + ) -> dict[str, str | int | None]: + applied: dict[str, str | int | None] = {"make": None, "model": None, "year_min": None, "year_max": None} if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make): applied["make"] = make self.pacer.after_filter_action() if model and self._try_fill_filter_input(page, ["input[placeholder*='Model']", "input[aria-label*='Model']"], model): applied["model"] = model self.pacer.after_filter_action() + if year_min is not None or year_max is not None: + if self._apply_year_range(page, year_min, year_max): + applied["year_min"] = year_min + applied["year_max"] = year_max + self.pacer.after_filter_action() return applied + def _apply_year_range(self, page: Page, year_min: int | None, year_max: int | None) -> bool: + """Заполняет поля фильтра Year и нажимает Apply Year.""" + if year_min is None and year_max is None: + return False + try: + success = page.evaluate( + """([yearMin, yearMax]) => { + const inputs = Array.from(document.querySelectorAll('input')); + const yearInputs = inputs.filter(inp => { + const v = parseInt(inp.value, 10); + return !isNaN(v) && v >= 1900 && v <= 2100; + }); + if (yearInputs.length < 2) return false; + yearInputs.sort((a, b) => parseInt(a.value) - parseInt(b.value)); + const setVal = (el, val) => { + const setter = Object.getOwnPropertyDescriptor( + HTMLInputElement.prototype, 'value' + ).set; + setter.call(el, String(val)); + el.dispatchEvent(new Event('input', {bubbles: true})); + el.dispatchEvent(new Event('change', {bubbles: true})); + }; + if (yearMin !== null) setVal(yearInputs[0], yearMin); + if (yearMax !== null) setVal(yearInputs[yearInputs.length - 1], yearMax); + const container = yearInputs[0].closest( + '[class*="filter"], [class*="year"], section, fieldset' + ) || yearInputs[0].parentElement.parentElement; + if (container) { + const btn = Array.from(container.querySelectorAll( + 'button, a, [role="button"], span[class*="apply"]' + )).find(el => /apply|\u043f\u0440\u0438\u043c\u0435\u043d/i.test(el.textContent)); + if (btn) { btn.click(); return true; } + } + yearInputs[yearInputs.length - 1].dispatchEvent( + new KeyboardEvent('keydown', { + key: 'Enter', code: 'Enter', keyCode: 13, bubbles: true + }) + ); + return true; + }""", + [year_min, year_max], + ) + if success: + try: + page.wait_for_load_state("domcontentloaded", timeout=15_000) + except Exception: + pass + self._wait_for_listing_content(page) + logger.info("Applied year range filter: %s — %s", year_min, year_max) + return True + except Exception as exc: + logger.warning("Failed to apply year range filter: %s", exc) + return False + def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult: # Считываем ссылки одним проходом по DOM. self._accept_cookie_banner(page) diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 34c5a77..45e0943 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -1,3 +1,4 @@ +import json import os from dataclasses import dataclass, field from pathlib import Path @@ -111,6 +112,79 @@ class ListingConfig: # Порог раннего останова: если доля уже известных машин на странице >= этого значения, # прекращаем листать — все новые машины уже найдены. 0 = отключено. early_stop_threshold: float = _env_float("IAAI_EARLY_STOP_THRESHOLD", 0.8) + # Сегментация листинга по брендам для обхода лимита пагинации IAAI (~22 600 машин). + # JSON-массив объектов: [{"make":"TOYOTA"},{"make":"FORD"},...] или "auto" для авто-списка. + # Пустая строка = без сегментации (backward compatible). + listing_segments_json: str = _env_str("IAAI_LISTING_SEGMENTS", "") + + +# Список брендов IAAI для автоматической сегментации. +# Покрывает >99% автомобилей на сайте. Порядок: от крупных к мелким. +IAAI_DEFAULT_MAKES: tuple[str, ...] = ( + "TOYOTA", "FORD", "CHEVROLET", "HONDA", "NISSAN", "HYUNDAI", + "KIA", "DODGE", "JEEP", "BMW", "MERCEDES-BENZ", "SUBARU", + "VOLKSWAGEN", "GMC", "MAZDA", "LEXUS", "CHRYSLER", "AUDI", + "RAM", "BUICK", "CADILLAC", "ACURA", "INFINITI", "LINCOLN", + "MITSUBISHI", "VOLVO", "JAGUAR", "LAND ROVER", "PORSCHE", + "MINI", "TESLA", "GENESIS", "FIAT", "ALFA ROMEO", "MASERATI", + "SCION", "PONTIAC", "SATURN", "MERCURY", "SAAB", "SUZUKI", + "OLDSMOBILE", "ISUZU", "HUMMER", "PLYMOUTH", "SMART", + "RIVIAN", "LUCID", "POLESTAR", "FERRARI", "LAMBORGHINI", + "BENTLEY", "ROLLS-ROYCE", "ASTON MARTIN", "MCLAREN", "LOTUS", + "MAYBACH", "FISKER", "GEO", "DAEWOO", "EAGLE", +) + +# Пагинационный потолок IAAI: ~226 страниц × 100 = 22 600 результатов. +IAAI_PAGINATION_CEILING = 22_600 + +# Бренды, потенциально превышающие потолок пагинации — разбиваем по годам. +_LARGE_MAKES: frozenset[str] = frozenset({ + "TOYOTA", "FORD", "CHEVROLET", "HONDA", "NISSAN", "HYUNDAI", + "KIA", "DODGE", "JEEP", +}) +_YEAR_SPLITS: tuple[tuple[int, int], ...] = ( + (1900, 2012), + (2013, 2019), + (2020, 2027), +) + + +def parse_listing_segments(raw: str) -> list[dict[str, str | int | None]]: + """Парсит IAAI_LISTING_SEGMENTS в список сегментов. + + Каждый сегмент — dict с ключами: make (str), year_min/year_max (int|None). + Специальное значение ``"auto"`` генерирует сегменты из IAAI_DEFAULT_MAKES. + Крупные бренды автоматически разбиваются по диапазонам годов. + """ + raw = raw.strip() + if not raw: + return [] + if raw.lower() == "auto": + segments: list[dict[str, str | int | None]] = [] + for m in IAAI_DEFAULT_MAKES: + if m in _LARGE_MAKES: + for yr_min, yr_max in _YEAR_SPLITS: + segments.append({"make": m, "year_min": yr_min, "year_max": yr_max}) + else: + segments.append({"make": m, "year_min": None, "year_max": None}) + return segments + try: + data = json.loads(raw) + except (json.JSONDecodeError, ValueError): + return [] + if not isinstance(data, list): + return [] + segments = [] + for item in data: + if isinstance(item, str): + segments.append({"make": item.upper(), "year_min": None, "year_max": None}) + elif isinstance(item, dict): + segments.append({ + "make": str(item.get("make") or "").upper() or None, + "year_min": int(item["year_min"]) if item.get("year_min") is not None else None, + "year_max": int(item["year_max"]) if item.get("year_max") is not None else None, + }) + return segments # - Конфиг PostgreSQL (URL, пул соединений, pool_recycle) diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index ddb7698..cf792de 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -20,7 +20,7 @@ from playwright.sync_api import BrowserContext, Page, sync_playwright from playwright.sync_api import TimeoutError as PlaywrightTimeoutError from .browser import BrowserFactory, HumanPacer, NetworkCapture -from .core.config import Settings, settings +from .core.config import Settings, settings, parse_listing_segments from .core.exceptions import AntiBotDetectedError, SiteStructureChangedError from .core.logs import set_trace_id, setup_logging from .core.retry import retryable @@ -120,8 +120,6 @@ class IAAIScraper: self.car_mapper = CarMapper() self.persistence = PersistenceService(self.settings) self._shutdown_requested = False - # HTTP-клиент для fast-path. - # При наличии прокси используем ProxyManager. proxy_url = self.settings.proxy.server if proxy_url: _proxy_kwargs: dict = { @@ -182,6 +180,13 @@ class IAAIScraper: pass finally: self.playwright = None + if getattr(self, "_http_pool", None) is not None: + try: + self._http_pool.clear() + except Exception: + pass + finally: + self._http_pool = None def _new_context(self) -> BrowserContext: if self.browser is None: @@ -199,21 +204,17 @@ class IAAIScraper: return {"status": "ok", "database_url": self.settings.database.url} def _warmup_visit(self, page: Page) -> None: - # Прогрев главной страницы. try: logger.info("Warmup: visiting homepage to pass anti-bot challenge...") page.goto(self.settings.home_url, wait_until="commit", timeout=30_000) - # Ждём domcontentloaded вместо networkidle. try: page.wait_for_load_state("domcontentloaded", timeout=6_000) except PlaywrightTimeoutError: pass - # Короткая проверка, что страница стабилизировалась. try: page.wait_for_function("() => document.title && document.title.length > 3", timeout=4_000) except PlaywrightTimeoutError: pass - # Короткая пауза для установки cookies. time.sleep(0.3) logger.info("Warmup done: %s (title=%s)", page.url, page.title()[:50]) except Exception as e: @@ -223,12 +224,10 @@ class IAAIScraper: def _get_page_with_warmup(self) -> Page: context = self._new_context() page = context.new_page() - # Прогрев с первой антибот-проверкой. self._warmup_visit(page) return page def _dedupe_urls(self, raw_urls: list[str]) -> list[str]: - # Нормализация и дедупликация URL. vehicle_urls: list[str] = [] seen: set[str] = set() for raw in raw_urls: @@ -239,7 +238,6 @@ class IAAIScraper: return vehicle_urls def _filter_known_urls(self, vehicle_urls: list[str]) -> tuple[list[str], int]: - # Отсев уже известных URL. url_to_origin_id = { url: self._extract_db_origin_id_from_url(url) for url in vehicle_urls @@ -276,6 +274,14 @@ class IAAIScraper: page_urls.append(normalized) return page_urls + @staticmethod + def _build_segment_listing_url(base_url: str, make: str | None) -> str: + """Построить URL листинга для сегмента. Make передаётся через query-параметр.""" + if not make: + return base_url + sep = "&" if "?" in base_url else "?" + return f"{base_url}{sep}Make={make.replace(' ', '%20')}" + def _recover_empty_listing_page( self, page: Page, @@ -284,11 +290,12 @@ class IAAIScraper: all_raw_urls: list[str], seen_urls: set[str], all_listing_origin_urls: set[str] | None = None, + listing_url: str | None = None, ) -> tuple[ListingPageResult, list[str]]: if page_number == 1: logger.warning("Page 1 returned 0 links — retrying listing open...") time.sleep(3) - self.listing_collector.open_cars_listing(page) + self.listing_collector.open_cars_listing(page, url_override=listing_url) page_result = self.listing_collector.collect_current_page(page, page_number=page_number) return page_result, self._extract_page_urls(page_result, all_raw_urls, seen_urls, all_listing_origin_urls) @@ -309,10 +316,24 @@ class IAAIScraper: target_page_number: int, make: str | None, model: str | None, + listing_url: str | None = None, + year_min: int | None = None, + year_max: int | None = None, + max_nav_pages: int = 10, ) -> Page: + # Ограничиваем глубину навигации: если до цели > max_nav_pages кликов — не пытаемся. + if target_page_number > max_nav_pages + 1: + raise RuntimeError( + f"Cannot resume at page {target_page_number}: " + f"exceeds max navigation depth ({max_nav_pages} pages)" + ) page = self._get_page_with_warmup() - self.listing_collector.open_cars_listing(page) - self.listing_collector.apply_filters(page, make=make, model=model) + try: + self.listing_collector.open_cars_listing(page, url_override=listing_url) + self.listing_collector.apply_filters(page, make=make, model=model, year_min=year_min, year_max=year_max) + except Exception: + page.close() + raise for expected_page in range(2, target_page_number + 1): if not self.listing_collector.go_to_next_page(page, expected_page_number=expected_page): @@ -322,120 +343,6 @@ class IAAIScraper: logger.warning("Listing resumed at page %d after recovery", target_page_number) return page - def _collect_listing_iterative( - self, - *, - 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] = [] - all_raw_urls: list[str] = [] - seen: set[str] = set() - 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: - self.listing_collector.open_cars_listing(page) - 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, - "links_found": len(page_result.vehicle_links), - }) - - # Собираем URL со страницы. - page_urls = self._extract_page_urls(page_result, all_raw_urls, seen) - - if not page_urls: - page_result, page_urls = self._recover_empty_listing_page( - page, - page_number=page_number, - all_raw_urls=all_raw_urls, - seen_urls=seen, - ) - if not page_urls and page_number > 1: - logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) - page.close() - page = self._reopen_listing_and_resume( - target_page_number=page_number, - make=make, - model=model, - ) - page_result = self.listing_collector.collect_current_page(page, page_number=page_number) - page_urls = self._extract_page_urls(page_result, all_raw_urls, seen) - if not page_urls: - logger.info("Page %d: 0 new links, stopping pagination", page_number) - break - - # Фильтруем known. - fresh, page_skipped = self._filter_known_urls(page_urls) - skipped_existing += page_skipped - new_urls.extend(fresh) - - logger.info( - "Page %d: %d links, %d new, %d known (total new: %d/%d)", - page_number, len(page_urls), len(fresh), page_skipped, - len(new_urls), limit, - ) - - # Стопаем только если набрали нужное количество И на этой странице уже нет новых. - # Если последняя страница дала новые — проверяем следующую (там могут быть ещё). - if len(new_urls) >= limit and len(fresh) == 0: - break - if len(new_urls) >= limit: - # Набрали достаточно, дальше не листаем - break - - if not page_result.next_page_detected: - logger.info("No next page detected, stopping") - break - if not self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1): - logger.warning("Failed to navigate to page %d — reopening listing and resuming", page_number + 1) - page.close() - page = self._reopen_listing_and_resume( - target_page_number=page_number + 1, - make=make, - model=model, - ) - finally: - page.close() - - # Обрезаем до limit. - new_urls = new_urls[:limit] - - listing = { - "status": "ok", - "listing_url": self.settings.listing.cars_url, - "pages_collected": len(pages_info), - "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 - def collect_listing( self, make: str | None = None, @@ -472,6 +379,9 @@ class IAAIScraper: started_at: float, start_page: int = 1, progress_callback: Callable[[int], None] | None = None, + listing_url: str | None = None, + year_min: int | None = None, + year_max: int | None = None, ) -> dict[str, Any]: batch_size = self.settings.celery.batch_size pending_urls: list[str] = [] @@ -500,9 +410,9 @@ class IAAIScraper: except Exception as exc: logger.warning("Could not load known origin_ids for early-stop: %s", exc) - # Совместимость: если only_new без явного limit, используем старую схему - # collect_listing -> filter_known -> batch upsert (это ожидают тесты и API-поток only_new). - if effective_only_new and (limit is None or limit <= 0): + # Совместимость: если only_new без limit и без сегментного URL — старая схема. + # При сегментированном скрапинге (listing_url/year фильтры) всегда streaming. + if effective_only_new and (limit is None or limit <= 0) and not listing_url and year_min is None and year_max is None: listing_payload = self.collect_listing( make=make, model=model, @@ -612,15 +522,20 @@ class IAAIScraper: try: if start_page <= 1: assert page is not None - self.listing_collector.open_cars_listing(page) - applied_filters = self.listing_collector.apply_filters(page, make=make, model=model) + self.listing_collector.open_cars_listing(page, url_override=listing_url) + applied_filters = self.listing_collector.apply_filters( + page, make=make, model=model, year_min=year_min, year_max=year_max, + ) else: page = self._reopen_listing_and_resume( target_page_number=start_page, make=make, model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, ) - applied_filters = {"make": make, "model": model} + applied_filters = {"make": make, "model": model, "year_min": year_min, "year_max": year_max} for page_number in range(start_page, max(start_page, self.settings.listing.max_pages_per_run) + 1): page_result = self.listing_collector.collect_current_page(page, page_number=page_number) @@ -643,22 +558,34 @@ class IAAIScraper: all_raw_urls=all_raw_urls, seen_urls=seen_urls, all_listing_origin_urls=all_listing_origin_urls, + listing_url=listing_url, ) if not page_urls and page_number > 1: logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number) page.close() - page = self._reopen_listing_and_resume( - target_page_number=page_number, - make=make, - model=model, - ) - page_result = self.listing_collector.collect_current_page(page, page_number=page_number) - page_urls = self._extract_page_urls( - page_result, - all_raw_urls, - seen_urls, - all_listing_origin_urls, - ) + try: + page = self._reopen_listing_and_resume( + target_page_number=page_number, + make=make, + model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) + page_result = self.listing_collector.collect_current_page(page, page_number=page_number) + page_urls = self._extract_page_urls( + page_result, + all_raw_urls, + seen_urls, + all_listing_origin_urls, + ) + except RuntimeError as resume_exc: + logger.warning( + "Cannot resume at page %d (%s) — stopping pagination for this segment", + page_number, resume_exc, + ) + page = self._get_page_with_warmup() + page_urls = [] if not page_urls: logger.info("Page %d: 0 new links, stopping pagination", page_number) break @@ -708,7 +635,7 @@ class IAAIScraper: limit if limit is not None else "∞", ) - # Явный прогресс в логах даже при более высоком уровне логирования. + # Прогресс каждые 10 страниц на уровне WARNING. if page_number % 10 == 0: logger.warning( "sync_listing progress: pages=%d, discovered=%d, upserted=%d, failed=%d, pending=%d", @@ -743,11 +670,22 @@ class IAAIScraper: if not self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1): logger.warning("Failed to navigate to page %d — reopening listing and resuming", page_number + 1) page.close() - page = self._reopen_listing_and_resume( - target_page_number=page_number + 1, - make=make, - model=model, - ) + try: + page = self._reopen_listing_and_resume( + target_page_number=page_number + 1, + make=make, + model=model, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, + ) + except RuntimeError as resume_exc: + logger.warning( + "Cannot resume at page %d (%s) — stopping pagination for this segment", + page_number + 1, resume_exc, + ) + page = self._get_page_with_warmup() + break if pending_urls: _process_pending_batch(pending_urls, cars_upserted + cars_failed) @@ -760,7 +698,8 @@ class IAAIScraper: cars_failed, ) finally: - page.close() + if page is not None: + page.close() return { "listing": { @@ -790,14 +729,12 @@ class IAAIScraper: @retryable(max_attempts=3, jitter_seconds=0.25) def scrape_vehicle_detail(self, vehicle_url: str): - # Открыть страницу, перехватить JSON, вернуть данные. page = self._get_page() try: return self._scrape_on_page(page, vehicle_url) finally: page.close() - # Быстрый извлекатель данных из inline JSON. _JS_EXTRACT = """ () => { try { @@ -867,7 +804,6 @@ class IAAIScraper: @staticmethod def _build_car_from_js(js_data: dict, vehicle_url: str) -> dict: - # Сбор vehicle_summary из JS-данных. if not js_data or not js_data.get("ok"): return {} @@ -909,7 +845,6 @@ class IAAIScraper: @staticmethod def _build_payload_insights(vehicle_summary: dict) -> dict: - # Сбор payload_insights из vehicle_summary. return { "vehicle_core": vehicle_summary, "pricing": { @@ -929,7 +864,6 @@ class IAAIScraper: trace_id = self._new_trace_id("scrape") started_at = time.perf_counter() - # Блокируем тяжёлые ресурсы страницы. if self.settings.block_resources: BrowserFactory.enable_resource_blocking(page) @@ -938,13 +872,11 @@ class IAAIScraper: page.goto(vehicle_url, wait_until="commit", timeout=60_000) - # На VPS достаточно domcontentloaded. try: page.wait_for_load_state("domcontentloaded", timeout=5_000) except PlaywrightTimeoutError: pass - # Быстрый путь: читаем данные из inline JSON. js_data: dict = {} try: js_data = page.evaluate(self._JS_EXTRACT) or {} @@ -952,9 +884,7 @@ class IAAIScraper: pass if js_data.get("ok"): - # Успешное извлечение структуры авто. vehicle_summary = self._build_car_from_js(js_data, vehicle_url) - # Быстрая проверка валидности данных. has_identity = bool(vehicle_summary.get("make") or vehicle_summary.get("lot_number")) if not has_identity: raise SiteStructureChangedError(f"JS extraction returned no vehicle identity for {vehicle_url}") @@ -982,7 +912,7 @@ class IAAIScraper: save_to_json(network_dump, self.settings.raw_output_json) return result - # Резервный путь: полный парсинг страницы. + # Резервный путь: полный HTML-парсинг. try: page.wait_for_selector("#VehicleDetailViewModel, .veh-details, .vehicle-details, [data-uname='vehicleDetailPage']", timeout=400) except PlaywrightTimeoutError: @@ -1029,7 +959,6 @@ class IAAIScraper: if pos < 0: return None - # Ищем начало внешнего JSON-объекта. script_start = html.rfind(" CarRecord: - # Быстрый запрос через context.request. if not self.context: self._new_context() assert self.context is not None @@ -1119,7 +1046,6 @@ class IAAIScraper: html = response.text() - # Пытаемся извлечь JSON напрямую из HTML. js_data = self._extract_iaai_json_from_html(html, vehicle_url) if not js_data: raise RuntimeError(f"HTTP fast-path missing embedded JSON for {vehicle_url}") @@ -1200,7 +1126,6 @@ class IAAIScraper: return CarRecord.model_validate(db_record.model_dump(mode="json")) def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"): - # Скрапинг и upsert одного авто. trace_id = self._new_trace_id("sync-vehicle") started_at = time.perf_counter() self.persistence.create_tables() @@ -1247,12 +1172,12 @@ class IAAIScraper: error_summary=error_summary, ) - # ── Tunables for sync_batch ── - _HTTP_MICRO_BATCH = 25 # URLs per micro-batch (avoid mass rate-limit) - _HTTP_MAX_RETRIES = 2 # Retries per URL before giving up to browser - _HTTP_RETRY_DELAYS = (0.4, 1.0) # Backoff between retries - _FALLBACK_PARALLEL_PAGES = 4 # Concurrent browser tabs for fallback - _FALLBACK_NAV_TIMEOUT_MS = 8000 # Reduced from 15 000 + # ── Настройки sync_batch ── + _HTTP_MICRO_BATCH = 25 # URL за микро-батч + _HTTP_MAX_RETRIES = 2 # Повторов на URL до перехода в браузер + _HTTP_RETRY_DELAYS = (0.4, 1.0) # Задержки между повторами + _FALLBACK_PARALLEL_PAGES = 4 # Параллельных вкладок для fallback + _FALLBACK_NAV_TIMEOUT_MS = 8000 # Таймаут навигации fallback def _browser_fallback_parallel( self, @@ -1260,13 +1185,11 @@ class IAAIScraper: total: int, n_pages: int, ) -> dict: - # Обработка fallback URL через браузер. records: list[CarRecord] = [] failures: list[dict[str, str]] = [] cars_failed = 0 protection_events = 0 - # Деление URL по страницам. slices: list[list[tuple[str, int]]] = [[] for _ in range(n_pages)] for i, item in enumerate(fallback_urls): slices[i % n_pages].append(item) @@ -1384,7 +1307,7 @@ class IAAIScraper: "protection_events": local_protection, } - # Один page обрабатываем в главном потоке. + # Одна страница — в главном потоке. if n_pages == 1: res = _process_slice(slices[0] if slices else []) records.extend(res["records"]) @@ -1392,7 +1315,7 @@ class IAAIScraper: cars_failed += res["cars_failed"] protection_events += res["protection_events"] else: - # Несколько страниц — параллельно через потоки (каждый поток со своей Page). + # Несколько страниц — параллельно, каждый поток со своей 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): @@ -1415,7 +1338,6 @@ class IAAIScraper: lane: str = "iaai_cars", parallel_tabs: int | None = None, ) -> dict: - # Пакетный скрапинг списка URL. trace_id = self._new_trace_id("sync-batch") started_at = time.perf_counter() num_workers = parallel_tabs or self.settings.parallel_tabs @@ -1433,14 +1355,13 @@ class IAAIScraper: total = len(vehicle_urls) logger.info("sync_batch: %d vehicles, %d parallel workers", total, num_workers) - # Подготовка для HTTP fast-path: cookies + user-agent из Playwright сессии. cookie_header = self._cookie_header_for_context() user_agent = ( "Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:128.0) " "Gecko/20100101 Firefox/128.0" ) - # ── Phase 1: HTTP fast-path — все URLs одним ThreadPoolExecutor ── + # ── Фаза 1: HTTP fast-path ── fallback_urls: list[tuple[str, int]] = [] def _fetch_one_with_retry(idx_url: tuple[int, str]) -> tuple[int, CarRecord | Exception]: @@ -1459,7 +1380,6 @@ class IAAIScraper: time.sleep(self._HTTP_RETRY_DELAYS[min(attempt, len(self._HTTP_RETRY_DELAYS) - 1)]) return idx, last_exc # type: ignore[return-value] - # Запускаем все URLs сразу — один пул потоков для минимального времени ожидания. indexed_urls = list(enumerate(vehicle_urls)) with ThreadPoolExecutor(max_workers=min(num_workers, total)) as executor: for idx, result in executor.map(_fetch_one_with_retry, indexed_urls): @@ -1483,7 +1403,7 @@ class IAAIScraper: http_successes, http_fallbacks, time.perf_counter() - started_at, ) - # ── Phase 2: browser fallback (single-threaded — Playwright sync API is not thread-safe) ── + # ── Фаза 2: браузерный fallback ── if fallback_urls: logger.info( "Fallback browser mode: %d/%d URLs, single page", @@ -1495,9 +1415,8 @@ class IAAIScraper: protection_events += fb_results["protection_events"] failures.extend(fb_results["failures"]) - # ── Phase 3: single DB flush ── + # ── Фаза 3: запись в БД ── if records: - # Дедупликация перед записью. seen_origins: set[str] = set() unique_records: list[CarRecord] = [] for rec in records: @@ -1540,9 +1459,11 @@ class IAAIScraper: only_new: bool | None = None, start_page: int = 1, progress_callback: Callable[[int], None] | None = None, + listing_url: str | None = None, + year_min: int | None = None, + year_max: int | None = None, ): - # Листинг + sync всех найденных машин. - # Применяем runtime_config как дефолты (CLI/API аргументы имеют приоритет). + # runtime_config — дефолты; CLI/API аргументы приоритетнее. rc = self.runtime_config.sync if limit is None and rc.limit is not None: limit = rc.limit @@ -1576,6 +1497,9 @@ class IAAIScraper: started_at=started_at, start_page=max(1, int(start_page)), progress_callback=progress_callback, + listing_url=listing_url, + year_min=year_min, + year_max=year_max, ) listing = stream_result["listing"] total = stream_result["total"] @@ -1588,8 +1512,7 @@ class IAAIScraper: logger.info("Streaming sync processed %d vehicles", total) - # Помечаем авто как проданные, если они исчезли из листинга. - # Только если сканирование было полным (не ограниченным limit/only_new/early_stop). + # Помечаем проданные авто, исчезнувшие из листинга (только при полном скане). is_partial_scan = ( effective_only_new or (limit is not None and limit > 0) @@ -1610,7 +1533,6 @@ class IAAIScraper: if not failures: failures.append({"vehicle_url": "collect_listing", "error": str(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 @@ -1632,7 +1554,9 @@ class IAAIScraper: "trace_id": trace_id, "status": status, "run_id": run_id, - "full_scan_completed": (not is_partial_scan) and status == "success", + # partial_success допустим — отдельные машины могли не спарситься, + # это не повод повторять весь bootstrap. + "full_scan_completed": (not is_partial_scan) and status in ("success", "partial_success"), "only_new_effective": effective_only_new, "listing": listing, "cars_upserted": cars_upserted, @@ -1644,9 +1568,132 @@ class IAAIScraper: "failures": failures, } + def sync_listing_segmented( + self, + segments: list[dict[str, Any]], + lane: str = "iaai_cars", + only_new: bool | None = None, + start_segment: int = 0, + start_page: int = 1, + progress_callback: Callable[[int, int], None] | None = None, + ) -> dict[str, Any]: + """Итеративный sync_listing по списку сегментов (бренд / бренд+годы). + + Args: + segments: список dict с ключами make, year_min, year_max. + start_segment: индекс сегмента для resume (0-based). + start_page: страница внутри start_segment для resume. + progress_callback: вызывается (segment_index, page_number) после каждой страницы. + """ + trace_id = self._new_trace_id("sync-segmented") + started_at = time.perf_counter() + base_url = self.settings.listing.cars_url + + total_cars_upserted = 0 + total_cars_failed = 0 + total_images_upserted = 0 + total_skipped = 0 + total_discovered = 0 + all_failures: list[dict[str, str]] = [] + segment_results: list[dict[str, Any]] = [] + completed_all = True + + logger.warning( + "Starting segmented sync: %d segments, resume from segment=%d page=%d", + len(segments), start_segment, start_page, + ) + + for seg_idx in range(start_segment, len(segments)): + seg = segments[seg_idx] + seg_make = seg.get("make") + seg_year_min = seg.get("year_min") + seg_year_max = seg.get("year_max") + seg_url = self._build_segment_listing_url(base_url, seg_make) if seg_make else None + + seg_start_page = start_page if seg_idx == start_segment else 1 + seg_label = f"{seg_make or 'ALL'}" + if seg_year_min is not None or seg_year_max is not None: + seg_label += f" ({seg_year_min}-{seg_year_max})" + + logger.warning( + "Segment %d/%d: %s (start_page=%d)", + seg_idx + 1, len(segments), seg_label, seg_start_page, + ) + + def _seg_progress(page_number: int, _si=seg_idx) -> None: + if progress_callback: + progress_callback(_si, page_number) + + try: + result = self.sync_listing( + make=None if seg_url else seg_make, + model=None, + lane=lane, + only_new=only_new, + start_page=seg_start_page, + progress_callback=_seg_progress, + listing_url=seg_url, + year_min=seg_year_min, + year_max=seg_year_max, + ) + total_cars_upserted += result.get("cars_upserted", 0) + total_cars_failed += result.get("cars_failed", 0) + total_images_upserted += result.get("images_upserted", 0) + total_skipped += result.get("skipped_existing", 0) + total_discovered += result.get("listing", {}).get("vehicles_collected", 0) + all_failures.extend(result.get("failures", [])) + segment_results.append({ + "segment": seg, + "segment_index": seg_idx, + "status": result.get("status"), + "cars_upserted": result.get("cars_upserted", 0), + "cars_failed": result.get("cars_failed", 0), + "vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0), + }) + + logger.warning( + "Segment %d/%d done: %s → upserted=%d, failed=%d, collected=%d", + seg_idx + 1, len(segments), seg_label, + result.get("cars_upserted", 0), + result.get("cars_failed", 0), + result.get("listing", {}).get("vehicles_collected", 0), + ) + except Exception as exc: + logger.error("Segment %d/%d failed: %s — %s", seg_idx + 1, len(segments), seg_label, exc) + all_failures.append({"vehicle_url": f"segment_{seg_idx}_{seg_label}", "error": str(exc)}) + completed_all = False + # Продолжаем оставшиеся сегменты — одна ошибка не должна убивать весь прогон + continue + + elapsed = round(time.perf_counter() - started_at, 3) + status = "success" if not all_failures else "partial_success" if total_cars_upserted else "failed" + + logger.warning( + "Segmented sync done: %d/%d segments, upserted=%d, failed=%d, discovered=%d, elapsed=%.1fs", + len(segment_results), len(segments), + total_cars_upserted, total_cars_failed, total_discovered, elapsed, + ) + + return { + "trace_id": trace_id, + "status": status, + # Bootstrap считается завершённым если все сегменты пройдены, + # даже если часть машин failed (они будут обновлены в следующих циклах). + "full_scan_completed": completed_all, + "segments_total": len(segments), + "segments_completed": len(segment_results), + "cars_upserted": total_cars_upserted, + "cars_failed": total_cars_failed, + "images_upserted": total_images_upserted, + "skipped_existing": total_skipped, + "total_discovered": total_discovered, + "elapsed_seconds": elapsed, + "failures": all_failures, + "segment_results": segment_results, + } + def run_scheduled(self) -> None: - # Резерв для standalone-режима. - # В текущей архитектуре планирование выполняется через Celery beat + worker/tasks.py. + """Standalone-планировщик (в продакшене используется Celery beat).""" def _handle_shutdown(signum, frame): logger.info("Received signal %s, shutting down gracefully...", signum) self._shutdown_requested = True diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index e0ec832..038a164 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -11,7 +11,7 @@ from billiard.exceptions import SoftTimeLimitExceeded from celery import shared_task from redis import Redis -from ..core.config import Settings +from ..core.config import Settings, parse_listing_segments from ..scraper import IAAIScraper from ..storage.db import PersistenceService @@ -166,7 +166,7 @@ def _release_lock_if_owner(redis_client: Redis, key: str, owner_token: str) -> N logger.warning("Failed to release lock %s", key, exc_info=True) -def _has_running_sync_listing_tasks(celery_app) -> bool: +def _has_running_sync_listing_tasks(celery_app, *, exclude_task_id: str | None = None) -> bool: try: inspector = celery_app.control.inspect(timeout=1.0) snapshots = [ @@ -182,12 +182,17 @@ def _has_running_sync_listing_tasks(celery_app) -> bool: for entries in snapshot.values(): for entry in entries or []: task_name = str(entry.get("name") or entry.get("request", {}).get("name") or "") - if task_name == SYNC_LISTING_TASK_NAME: - return True + if task_name != SYNC_LISTING_TASK_NAME: + continue + # Исключаем текущую задачу — она не считается "другой запущенной" + entry_id = str(entry.get("id") or entry.get("request", {}).get("id") or "") + if exclude_task_id and entry_id == exclude_task_id: + continue + return True return False -def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app) -> bool: +def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app, *, current_task_id: str | None = None) -> bool: try: owner_token = redis_client.get(SYNC_LISTING_LOCK_KEY) if not owner_token: @@ -196,7 +201,7 @@ def _clear_orphan_sync_listing_lock(redis_client: Redis, celery_app) -> bool: logger.warning("Failed to read sync listing lock before cleanup", exc_info=True) return False - if _has_running_sync_listing_tasks(celery_app): + if _has_running_sync_listing_tasks(celery_app, exclude_task_id=current_task_id): logger.info("sync_listing lock preserved: active task still detected") return False @@ -256,6 +261,7 @@ def _save_sync_checkpoint( make: str | None, model: str | None, lane: str, + segment_index: int | None = None, ) -> None: payload = { "status": "in_progress", @@ -264,6 +270,7 @@ def _save_sync_checkpoint( "make": make, "model": model, "lane": lane, + "segment_index": segment_index, "updated_at": int(time.time()), } try: @@ -384,7 +391,7 @@ def sync_listing_task( lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl) if not lock_acquired: - orphan_cleared = _clear_orphan_sync_listing_lock(redis_client, self.app) + orphan_cleared = _clear_orphan_sync_listing_lock(redis_client, self.app, current_task_id=task_id) if orphan_cleared: lock_acquired = _acquire_lock(redis_client, SYNC_LISTING_LOCK_KEY, owner_token, lock_ttl) @@ -409,16 +416,20 @@ def sync_listing_task( checkpoint_make = checkpoint.get("make") checkpoint_model = checkpoint.get("model") checkpoint_lane = checkpoint.get("lane") + checkpoint_segment = checkpoint.get("segment_index") same_scope = ( checkpoint_make == make and checkpoint_model == model and checkpoint_lane == lane ) - if checkpoint_page > 0 and same_scope: + # Для сегментированного режима: совпадение по lane + наличие segment_index + is_segmented_checkpoint = checkpoint_segment is not None and checkpoint_lane == lane + if checkpoint_page > 0 and (same_scope or is_segmented_checkpoint): resume_from_page = checkpoint_page + 1 logger.warning( - "Resuming sync_listing from page %d using checkpoint", + "Resuming sync_listing from page %d (segment=%s) using checkpoint", resume_from_page, + checkpoint_segment, ) elif checkpoint_page > 0: logger.info("Ignoring stale checkpoint due to different sync parameters") @@ -442,8 +453,38 @@ def sync_listing_task( lock_ttl, ) self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id}) + + # Определяем сегменты из конфига. + settings = Settings() + segments = parse_listing_segments(settings.listing.listing_segments_json) + use_segmented = bool(segments) and make is None and model is None + + resume_from_segment = 0 + if use_segmented and checkpoint and str(checkpoint.get("status") or "") == "in_progress": + cp_segment = checkpoint.get("segment_index") + if cp_segment is not None and int(cp_segment) >= 0: + resume_from_segment = int(cp_segment) + # resume_from_page уже вычислен выше + def _job(): with IAAIScraper() as scraper: + if use_segmented: + return scraper.sync_listing_segmented( + segments=segments, + lane=lane, + only_new=effective_only_new, + start_segment=resume_from_segment, + start_page=resume_from_page, + progress_callback=lambda seg_idx, page_number: _save_sync_checkpoint( + redis_client, + task_id=task_id, + page_number=page_number, + make=None, + model=None, + lane=lane, + segment_index=seg_idx, + ), + ) return scraper.sync_listing( make=make, model=model,