add full scan mode

This commit is contained in:
qananasikq
2026-04-14 19:32:34 +03:00
parent 0add3913eb
commit cac68bda4c
3 changed files with 318 additions and 158 deletions

View File

@@ -173,7 +173,7 @@ class ProxyConfig:
return result return result
# --- Главный объект настроек: собирает все блоки конфигурации --- # Главный объект настроек: собирает все блоки конфигурации
@dataclass(slots=True) @dataclass(slots=True)
class Settings: class Settings:
@@ -192,7 +192,7 @@ class Settings:
log_level: str = _env_str("IAAI_LOG_LEVEL", "INFO") log_level: str = _env_str("IAAI_LOG_LEVEL", "INFO")
log_file: str | None = _env_optional_str("IAAI_LOG_FILE") log_file: str | None = _env_optional_str("IAAI_LOG_FILE")
enable_trace_id_logs: bool = _env_bool("IAAI_ENABLE_TRACE_ID_LOGS", True) 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") raw_output_json: str | None = _env_optional_str("IAAI_RAW_OUTPUT_JSON")
tokens_file: str | None = _env_path_str("IAAI_TOKENS_FILE") tokens_file: str | None = _env_path_str("IAAI_TOKENS_FILE")
runtime_config_file: str | None = _env_path_str("IAAI_RUNTIME_CONFIG_FILE") runtime_config_file: str | None = _env_path_str("IAAI_RUNTIME_CONFIG_FILE")

View File

@@ -9,6 +9,7 @@ import html as html_module
from concurrent.futures import ThreadPoolExecutor, as_completed from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from typing import Any
from urllib.parse import urlsplit, urlunsplit from urllib.parse import urlsplit, urlunsplit
from urllib.request import Request, urlopen from urllib.request import Request, urlopen
@@ -394,6 +395,299 @@ class IAAIScraper:
**listing, **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: def _get_page(self) -> Page:
if not self.context: if not self.context:
self._new_context() self._new_context()
@@ -1170,167 +1464,30 @@ class IAAIScraper:
images_upserted = 0 images_upserted = 0
total = 0 total = 0
skipped_existing = 0 skipped_existing = 0
is_partial_scan = True
failures: list[dict[str, str]] = [] failures: list[dict[str, str]] = []
listing: dict = {} listing: dict = {}
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 stream_result = self._sync_listing_streaming(
listing_time_budget_s: float | None = None make=make,
if soft_time_limit and soft_time_limit > 0: model=model,
# Делим время: ~40% на листинг, ~60% на парсинг + запись. lane=lane,
# При soft=120 → listing=48s, processing=72s (достаточно для ~7 батчей). limit=limit,
# При soft=600 → listing=240s, processing=360s (достаточно для ~30+ батчей). effective_only_new=effective_only_new,
# Минимум 30s на листинг, оставляем 30s запас до hard limit. started_at=started_at,
soft_f = float(soft_time_limit) )
listing_time_budget_s = max(30.0, soft_f * 0.4) 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"]
# ── Сбор листинга ── logger.info("Streaming sync processed %d vehicles", total)
# При 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)})
# Помечаем авто как проданные, если они исчезли из листинга. # Помечаем авто как проданные, если они исчезли из листинга.
# Только если сканирование было полным (не ограниченным limit/only_new/early_stop). # Только если сканирование было полным (не ограниченным limit/only_new/early_stop).
@@ -1338,6 +1495,7 @@ class IAAIScraper:
effective_only_new effective_only_new
or (limit is not None and limit > 0) or (limit is not None and limit > 0)
or listing.get("early_stopped", False) or listing.get("early_stopped", False)
or listing.get("truncated_by_time_budget", False)
) )
if all_listing_origin_urls and not is_partial_scan: if all_listing_origin_urls and not is_partial_scan:
try: try:
@@ -1375,6 +1533,8 @@ class IAAIScraper:
"trace_id": trace_id, "trace_id": trace_id,
"status": status, "status": status,
"run_id": run_id, "run_id": run_id,
"full_scan_completed": (not is_partial_scan) and status == "success",
"only_new_effective": effective_only_new,
"listing": listing, "listing": listing,
"cars_upserted": cars_upserted, "cars_upserted": cars_upserted,
"cars_failed": cars_failed, "cars_failed": cars_failed,

View File

@@ -6,7 +6,7 @@
"ids_max_pages": null, "ids_max_pages": null,
"condition_check_enabled": false, "condition_check_enabled": false,
"lane": "iaai_cars", "lane": "iaai_cars",
"only_new": true, "only_new": false,
"limit": null "limit": null
}, },
"filters": { "filters": {