From cac68bda4c84ab9d72c63863cb6df179006eda29 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Tue, 14 Apr 2026 19:32:34 +0300 Subject: [PATCH] add full scan mode --- iaai_scraper/core/config.py | 4 +- iaai_scraper/scraper.py | 470 ++++++++++++++++++++++++------------ runtime_config.json | 2 +- 3 files changed, 318 insertions(+), 158 deletions(-) diff --git a/iaai_scraper/core/config.py b/iaai_scraper/core/config.py index 4690527..34c5a77 100644 --- a/iaai_scraper/core/config.py +++ b/iaai_scraper/core/config.py @@ -173,7 +173,7 @@ class ProxyConfig: return result -# --- Главный объект настроек: собирает все блоки конфигурации --- +# Главный объект настроек: собирает все блоки конфигурации @dataclass(slots=True) class Settings: @@ -192,7 +192,7 @@ class Settings: log_level: str = _env_str("IAAI_LOG_LEVEL", "INFO") log_file: str | None = _env_optional_str("IAAI_LOG_FILE") enable_trace_id_logs: bool = _env_bool("IAAI_ENABLE_TRACE_ID_LOGS", True) - sync_only_new: bool = _env_bool("IAAI_SYNC_ONLY_NEW", True) + sync_only_new: bool = _env_bool("IAAI_SYNC_ONLY_NEW", False) raw_output_json: str | None = _env_optional_str("IAAI_RAW_OUTPUT_JSON") tokens_file: str | None = _env_path_str("IAAI_TOKENS_FILE") runtime_config_file: str | None = _env_path_str("IAAI_RUNTIME_CONFIG_FILE") diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index 1f20fc0..fc22b35 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -9,6 +9,7 @@ import html as html_module from concurrent.futures import ThreadPoolExecutor, as_completed from datetime import datetime, timezone from pathlib import Path +from typing import Any from urllib.parse import urlsplit, urlunsplit from urllib.request import Request, urlopen @@ -394,6 +395,299 @@ class IAAIScraper: **listing, } + def _sync_listing_streaming( + self, + *, + make: str | None, + model: str | None, + lane: str, + limit: int | None, + effective_only_new: bool, + started_at: float, + ) -> dict[str, Any]: + batch_size = self.settings.celery.batch_size + pending_urls: list[str] = [] + seen_urls: set[str] = set() + all_raw_urls: list[str] = [] + all_listing_origin_urls: set[str] = set() + pages_info: list[dict[str, Any]] = [] + skipped_existing = 0 + total = 0 + cars_upserted = 0 + cars_failed = 0 + images_upserted = 0 + failures: list[dict[str, str]] = [] + early_stopped = False + truncated_by_time_budget = False + + known_origin_ids: set[str] | None = None + threshold = self.settings.listing.early_stop_threshold + if effective_only_new and threshold > 0.0: + try: + known_origin_ids = self.persistence.get_all_origin_ids_for_lane("iaai:") + logger.info( + "Loaded %d known origin_ids for early-stop listing", + len(known_origin_ids), + ) + 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): + listing_payload = self.collect_listing( + make=make, + model=model, + known_origin_ids=known_origin_ids, + max_duration_seconds=None, + ) + raw_urls = list(listing_payload.get("vehicle_urls", [])) + deduped_urls = self._dedupe_urls(raw_urls) + fresh_urls, page_skipped = self._filter_known_urls(deduped_urls) + skipped_existing += page_skipped + pending_urls.extend(fresh_urls) + total = len(fresh_urls) + + for raw in raw_urls: + normalized = self._normalize_vehicle_url(raw) + if normalized: + all_listing_origin_urls.add(normalized) + + logger.info( + "Starting sync: %d vehicles to process (batch mode, only_new legacy path)", + total, + ) + + while len(pending_urls) >= batch_size: + batch_urls = pending_urls[:batch_size] + 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 Exception as batch_exc: + logger.error("Legacy only_new batch failed: %s", batch_exc) + cars_failed += len(batch_urls) + failures.append({"vehicle_url": "legacy_only_new_batch", "error": str(batch_exc)}) + pending_urls = pending_urls[batch_size:] + + if pending_urls: + try: + batch_result = self.sync_batch(pending_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 Exception as batch_exc: + logger.error("Legacy only_new final batch failed: %s", batch_exc) + cars_failed += len(pending_urls) + failures.append({"vehicle_url": "legacy_only_new_final_batch", "error": str(batch_exc)}) + + return { + "listing": listing_payload, + "total": total, + "skipped_existing": skipped_existing, + "cars_upserted": cars_upserted, + "cars_failed": cars_failed, + "images_upserted": images_upserted, + "failures": failures, + "all_listing_origin_urls": all_listing_origin_urls, + } + + def _process_pending_batch(batch_urls: list[str], batch_start: int) -> bool: + nonlocal cars_upserted, cars_failed, images_upserted, failures + if not batch_urls: + return True + + logger.info( + "Processing batch %d-%d of %d", + batch_start + 1, + batch_start + len(batch_urls), + total, + ) + 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", [])) + return True + except PlaywrightError as pw_exc: + logger.warning( + "Batch %d-%d browser crashed: %s. Reinitializing browser context.", + batch_start + 1, + batch_start + len(batch_urls), + 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") + return True + except Exception as reinit_exc: + logger.error("Failed to reinitialize browser: %s. Aborting remaining batches.", reinit_exc) + return False + except Exception as batch_exc: + logger.error( + "Batch %d-%d failed: %s. Continuing with next batch.", + batch_start + 1, + batch_start + len(batch_urls), + batch_exc, + ) + cars_failed += len(batch_urls) + failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)}) + return True + + page = self._get_page_with_warmup() + try: + self.listing_collector.open_cars_listing(page) + applied_filters = self.listing_collector.apply_filters(page, make=make, model=model) + + for page_number in range(1, max(1, self.settings.listing.max_pages_per_run) + 1): + 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), + }) + + page_urls: list[str] = [] + for item in page_result.vehicle_links: + normalized = self._normalize_vehicle_url(item.href) + all_raw_urls.append(item.href) + if normalized: + all_listing_origin_urls.add(normalized) + if normalized and normalized not in seen_urls: + seen_urls.add(normalized) + page_urls.append(normalized) + + if not page_urls: + if page_number == 1: + logger.warning("Page 1 returned 0 links — retrying listing open...") + time.sleep(3) + self.listing_collector.open_cars_listing(page) + page_result = self.listing_collector.collect_current_page(page, page_number=page_number) + for item in page_result.vehicle_links: + normalized = self._normalize_vehicle_url(item.href) + all_raw_urls.append(item.href) + if normalized: + all_listing_origin_urls.add(normalized) + if normalized and normalized not in seen_urls: + seen_urls.add(normalized) + page_urls.append(normalized) + if not page_urls: + logger.info("Page %d: 0 new links, stopping pagination", page_number) + break + + if known_origin_ids is not None and threshold > 0.0 and page_result.vehicle_links: + page_known = sum( + 1 + for item in page_result.vehicle_links + if item.lot_number and f"iaai:{item.lot_number}" in known_origin_ids + ) + page_new = len(page_result.vehicle_links) - page_known + ratio = page_known / len(page_result.vehicle_links) + if ratio >= threshold and page_new == 0: + logger.info( + "Early stop on page %d: %.0f%% known (%d/%d), 0 new >= threshold %.0f%%", + page_number, + ratio * 100, + page_known, + len(page_result.vehicle_links), + threshold * 100, + ) + early_stopped = True + break + + fresh_urls = page_urls + page_skipped = 0 + if effective_only_new: + fresh_urls, page_skipped = self._filter_known_urls(page_urls) + skipped_existing += page_skipped + + if limit is not None and limit > 0: + remaining_limit = max(0, limit - total - len(pending_urls)) + if remaining_limit <= 0: + break + fresh_urls = fresh_urls[:remaining_limit] + + pending_urls.extend(fresh_urls) + total += len(fresh_urls) + + logger.info( + "Page %d: %d links, %d new, %d known (total new: %d/%s)", + page_number, + len(page_urls), + len(fresh_urls), + page_skipped, + total, + limit if limit is not None else "∞", + ) + + # Явный прогресс в логах даже при более высоком уровне логирования. + if page_number % 10 == 0: + logger.warning( + "sync_listing progress: pages=%d, discovered=%d, upserted=%d, failed=%d, pending=%d", + page_number, + total, + cars_upserted, + cars_failed, + len(pending_urls), + ) + + while len(pending_urls) >= batch_size: + batch_urls = pending_urls[:batch_size] + if not _process_pending_batch(batch_urls, cars_upserted + cars_failed): + pending_urls = [] + break + pending_urls = pending_urls[batch_size:] + + if limit is not None and limit > 0 and total >= limit: + break + if ( + self.settings.listing.collect_current_page_only + or not self.settings.listing.include_pagination + or not page_result.next_page_detected + ): + break + if not self.listing_collector.go_to_next_page(page): + break + + if pending_urls: + _process_pending_batch(pending_urls, cars_upserted + cars_failed) + + logger.warning( + "sync_listing final progress: pages=%d, discovered=%d, upserted=%d, failed=%d", + len(pages_info), + total, + cars_upserted, + cars_failed, + ) + finally: + page.close() + + return { + "listing": { + "status": "ok", + "listing_url": self.settings.listing.cars_url, + "applied_filters": applied_filters, + "pages_collected": len(pages_info), + "vehicles_collected": len(all_raw_urls), + "vehicle_urls": all_raw_urls, + "early_stopped": early_stopped, + "truncated_by_time_budget": truncated_by_time_budget, + "pages": pages_info, + }, + "total": total, + "skipped_existing": skipped_existing, + "cars_upserted": cars_upserted, + "cars_failed": cars_failed, + "images_upserted": images_upserted, + "failures": failures, + "all_listing_origin_urls": all_listing_origin_urls, + } + def _get_page(self) -> Page: if not self.context: self._new_context() @@ -1170,167 +1464,30 @@ class IAAIScraper: images_upserted = 0 total = 0 skipped_existing = 0 + is_partial_scan = True failures: list[dict[str, str]] = [] listing: dict = {} 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) + stream_result = self._sync_listing_streaming( + make=make, + model=model, + lane=lane, + limit=limit, + effective_only_new=effective_only_new, + started_at=started_at, + ) + listing = stream_result["listing"] + total = stream_result["total"] + skipped_existing = stream_result["skipped_existing"] + cars_upserted = stream_result["cars_upserted"] + cars_failed = stream_result["cars_failed"] + images_upserted = stream_result["images_upserted"] + failures.extend(stream_result["failures"]) + all_listing_origin_urls = stream_result["all_listing_origin_urls"] - # ── Сбор листинга ── - # При only_new + limit используем итеративный подход: - # листаем страницы одну за другой, фильтруем known на лету, - # останавливаемся когда набрали limit новых. - original_listing_cap = self.settings.listing.max_vehicles_per_run - original_include_pagination = self.settings.listing.include_pagination - original_collect_current_page_only = self.settings.listing.collect_current_page_only - 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=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 - 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, - # чтобы listing_collector мог остановиться при встрече старых страниц. - known_origin_ids: set[str] | None = None - if effective_only_new and self.settings.listing.early_stop_threshold > 0.0: - try: - known_origin_ids = self.persistence.get_all_origin_ids_for_lane("iaai:") - logger.info( - "Loaded %d known origin_ids for early-stop listing", - len(known_origin_ids), - ) - except Exception as exc: - logger.warning("Could not load known origin_ids for early-stop: %s", exc) - - try: - listing = self.collect_listing( - 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 - self.settings.listing.include_pagination = original_include_pagination - self.settings.listing.collect_current_page_only = original_collect_current_page_only - self.settings.listing.max_pages_per_run = original_max_pages_per_run - - raw_urls = list(listing.get("vehicle_urls", [])) - vehicle_urls = self._dedupe_urls(raw_urls) - - if effective_only_new: - vehicle_urls, skipped_existing = self._filter_known_urls(vehicle_urls) - - if limit is not None: - 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 для последующей пометки проданных. - # Используем URL (а не origin_id из URL), т.к. в БД origin_id берётся из parsed lot_number, - # который может отличаться от числа в URL листинга. - all_listing_origin_urls = set() - for raw_url in raw_urls: - normalized_url = self._normalize_vehicle_url(raw_url) - if normalized_url: - 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, - ) - 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)}) + logger.info("Streaming sync processed %d vehicles", total) # Помечаем авто как проданные, если они исчезли из листинга. # Только если сканирование было полным (не ограниченным limit/only_new/early_stop). @@ -1338,6 +1495,7 @@ class IAAIScraper: effective_only_new or (limit is not None and limit > 0) or listing.get("early_stopped", False) + or listing.get("truncated_by_time_budget", False) ) if all_listing_origin_urls and not is_partial_scan: try: @@ -1375,6 +1533,8 @@ class IAAIScraper: "trace_id": trace_id, "status": status, "run_id": run_id, + "full_scan_completed": (not is_partial_scan) and status == "success", + "only_new_effective": effective_only_new, "listing": listing, "cars_upserted": cars_upserted, "cars_failed": cars_failed, diff --git a/runtime_config.json b/runtime_config.json index a64bd1a..bb43bcb 100644 --- a/runtime_config.json +++ b/runtime_config.json @@ -6,7 +6,7 @@ "ids_max_pages": null, "condition_check_enabled": false, "lane": "iaai_cars", - "only_new": true, + "only_new": false, "limit": null }, "filters": {