checkpoint: pipeline + 27k/h baseline before sitemap discovery

This commit is contained in:
qananasikq
2026-04-18 16:34:09 +03:00
parent a8c5f8aa41
commit 5999c2e42d
30 changed files with 4993 additions and 61 deletions

View File

@@ -17,8 +17,7 @@ def first_non_empty(values: Iterable[Any]) -> Any | None:
return None
# Регулярные выражения для VIN, lot и price.
VIN_RE = re.compile(r"\b([A-HJ-NPR-Z0-9]{17})\b", re.IGNORECASE)
# Регулярные выражения для lot и price.
LOT_RE = re.compile(r"\b(\d{7,10})\b")
PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)")

View File

@@ -1,12 +1,13 @@
import json
import logging
import os
import random
import re
import signal
import time
import uuid
import html as html_module
from concurrent.futures import ThreadPoolExecutor, as_completed
from concurrent.futures import Future, ThreadPoolExecutor, as_completed
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Callable
@@ -15,6 +16,13 @@ from urllib.request import Request, urlopen
import urllib3
try:
import orjson as _orjson # ~3-5× быстрее stdlib json на парсинге
_HAS_ORJSON = True
except ImportError:
_orjson = None
_HAS_ORJSON = False
from playwright.sync_api import Error as PlaywrightError
from playwright.sync_api import BrowserContext, Page, sync_playwright
from playwright.sync_api import TimeoutError as PlaywrightTimeoutError
@@ -37,6 +45,69 @@ VEHICLE_ID_RE = re.compile(r"/VehicleDetail/(\d+)(?:~[A-Z]{2})?", re.IGNORECASE)
HTML_TAG_RE = re.compile(r"<[^>]+>")
SCRIPT_STYLE_RE = re.compile(r"<(script|style)[^>]*>.*?</\1>", re.IGNORECASE | re.DOTALL)
# ── Anti-stall watchdog ──
# Жёсткие верхние границы на per-page операции браузера. Если Playwright "тихо" завис
# (частая реакция IAAI на затяжной скрапинг — page hangs без raise), SIGALRM вернёт
# управление и существующая recover-логика (_reopen_listing_and_resume) стартует
# новый context с нуля.
_PAGE_COLLECT_TIMEOUT_S = 90
_PAGE_NEXT_TIMEOUT_S = 60
# Проактивная ротация контекста отключена: reactive-восстановление
# (retry при пустой странице + _reopen_listing_and_resume) уже закрывает
# anti-bot соскоки, а проактивная ротация требовала долистывать кликами
# после reopen и конфликтовала с лимитом max_nav_pages в recovery-пути.
# Включить можно переменной окружения IAAI_CONTEXT_ROTATE_EVERY_PAGES.
_CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0)
class PageOperationTimeoutError(Exception):
"""Browser page operation exceeded per-op watchdog timeout."""
class _PageOpWatchdog:
"""SIGALRM-based watchdog.
Работает только в главном потоке процесса (Celery prefork child — это ок).
В других контекстах — no-op.
"""
def __init__(self, seconds: int, label: str) -> None:
self.seconds = max(1, int(seconds))
self.label = label
self._old_handler = None
self._active = False
def _handler(self, signum, frame): # noqa: ARG002
raise PageOperationTimeoutError(
f"Page operation '{self.label}' exceeded {self.seconds}s watchdog"
)
def __enter__(self):
import threading as _threading
if _threading.current_thread() is not _threading.main_thread():
return self
if not hasattr(signal, "SIGALRM"):
return self
try:
self._old_handler = signal.signal(signal.SIGALRM, self._handler)
signal.alarm(self.seconds)
self._active = True
except (ValueError, OSError):
# signal можно выставлять только из main thread; в остальном пропускаем.
self._active = False
return self
def __exit__(self, exc_type, exc, tb):
if not self._active:
return False
try:
signal.alarm(0)
if self._old_handler is not None:
signal.signal(signal.SIGALRM, self._old_handler)
except (ValueError, OSError):
pass
return False
class IAAIScraper:
@@ -120,13 +191,24 @@ class IAAIScraper:
self.car_mapper = CarMapper()
self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False
# HTTP pool: при parallel_tabs=24 и concurrency=4 пик ≈ 96 конкурентных сокетов;
# даём запас до 256, чтобы не упираться в PoolError под всплесками PX/retry.
# TCP_NODELAY выключает Nagle: для коротких HTTPS-запросов с keep-alive
# это снижает tail-latency на десятки ms.
import socket as _socket
_socket_options = [
(_socket.IPPROTO_TCP, _socket.TCP_NODELAY, 1),
(_socket.SOL_SOCKET, _socket.SO_KEEPALIVE, 1),
]
proxy_url = self.settings.proxy.server
if proxy_url:
_proxy_kwargs: dict = {
"num_pools": 4,
"maxsize": 64,
"num_pools": 8,
"maxsize": 256,
"block": False,
"retries": False,
"timeout": urllib3.Timeout(connect=5, read=10),
"timeout": urllib3.Timeout(connect=5, read=20),
"socket_options": _socket_options,
}
if self.settings.proxy.username:
_proxy_kwargs["proxy_headers"] = urllib3.make_headers(
@@ -136,10 +218,12 @@ class IAAIScraper:
logger.info("HTTP fast-path using proxy: %s", proxy_url)
else:
self._http_pool = urllib3.PoolManager(
num_pools=4,
maxsize=64,
num_pools=8,
maxsize=256,
block=False,
retries=False,
timeout=urllib3.Timeout(connect=5, read=10),
timeout=urllib3.Timeout(connect=5, read=20),
socket_options=_socket_options,
)
def _new_trace_id(self, prefix: str) -> str:
@@ -515,32 +599,44 @@ class IAAIScraper:
"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
# ── Pipeline-режим: батчи (HTTP + DB upsert) выполняются в фоне,
# пока главный поток продолжает собирать listing-страницы через
# Playwright. Это убирает простой Playwright во время HTTP-фазы
# (~50-65 сек на 400 машин) и даёт ~25-30% к throughput.
# max_workers=1: батчи исполняются последовательно, но ВНЕ
# главного потока, чтобы не блокировать сбор listing-а.
# Browser fallback откладывается (skip_browser_fallback=True),
# потому что Playwright не thread-safe — обработаем после loop'а.
batch_executor = ThreadPoolExecutor(max_workers=1, thread_name_prefix="batch-pipeline")
pending_future: Future | None = None
pending_future_meta: tuple[int, int] | None = None # (batch_start, batch_size)
deferred_fallbacks: list[tuple[str, int]] = []
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", []))
def _drain_pending_future() -> bool:
"""Waits for in-flight batch and accumulates its result.
Returns False if a non-recoverable error occurred and remaining
batches should be aborted (mirrors the legacy contract).
"""
nonlocal pending_future, pending_future_meta
nonlocal cars_upserted, cars_failed, images_upserted, failures
if pending_future is None:
return True
future = pending_future
meta = pending_future_meta or (0, 0)
pending_future = None
pending_future_meta = None
batch_start, batch_size_actual = meta
try:
batch_result = future.result()
except PlaywrightError as pw_exc:
logger.warning(
"Batch %d-%d browser crashed: %s. Reinitializing browser context.",
batch_start + 1,
batch_start + len(batch_urls),
batch_start + batch_size_actual,
pw_exc,
)
cars_failed += len(batch_urls)
cars_failed += batch_size_actual
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(pw_exc)})
try:
self._new_context()
@@ -553,12 +649,51 @@ class IAAIScraper:
logger.error(
"Batch %d-%d failed: %s. Continuing with next batch.",
batch_start + 1,
batch_start + len(batch_urls),
batch_start + batch_size_actual,
batch_exc,
)
cars_failed += len(batch_urls)
cars_failed += batch_size_actual
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
return True
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", []))
deferred_fallbacks.extend(batch_result.get("deferred_fallback_urls", []))
return True
def _process_pending_batch(batch_urls: list[str], batch_start: int) -> bool:
nonlocal pending_future, pending_future_meta
# Drain previous batch first (if still running).
if not _drain_pending_future():
return False
if not batch_urls:
return True
logger.info(
"Processing batch %d-%d of %d (pipelined)",
batch_start + 1,
batch_start + len(batch_urls),
total,
)
# Cookie header обязательно вычисляем в главном потоке — Playwright
# context.cookies() использует greenlet, который нельзя дёргать из background thread'а.
cookie_header = self._cookie_header_for_context()
pending_future = batch_executor.submit(
self.sync_batch,
batch_urls,
lane=lane,
skip_browser_fallback=True,
cookie_header_override=cookie_headerkground thread'а.
cookie_header = self._cookie_header_for_context()
pending_future = batch_executor.submit(
self.sync_batch,
batch_urls,
lane=lane,
skip_browser_fallback=True,
cookie_header_override=cookie_header,
)
pending_future_meta = (batch_start, len(batch_urls))
return True
page = None
try:
@@ -570,12 +705,77 @@ class IAAIScraper:
year_max=year_max,
)
pages_since_rotation = 0
for page_number in range(1, self.settings.listing.max_pages_per_run + 1):
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
# Проактивная ротация контекста (опционально, по умолчанию выключена).
# Если включена — требует долистывания до page_number после reopen,
# поэтому max_nav_pages поднят под реальную глубину.
if _CONTEXT_ROTATE_EVERY_PAGES > 0 and pages_since_rotation >= _CONTEXT_ROTATE_EVERY_PAGES:
logger.warning(
"Rotating browser context at page %d (every %d pages) to reset anti-bot state",
page_number, _CONTEXT_ROTATE_EVERY_PAGES,
)
try:
page.close()
except Exception:
pass
page = None
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,
max_nav_pages=max(page_number, 500),
)
pages_since_rotation = 0
except ListingResumeError as resume_exc:
logger.warning(
"Context rotation could not resume at page %d (%s) — stopping segment",
page_number, resume_exc,
)
break
try:
with _PageOpWatchdog(_PAGE_COLLECT_TIMEOUT_S, f"collect page {page_number}"):
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
except (PageOperationTimeoutError, PlaywrightError) as op_exc:
logger.warning(
"Page %d collect stalled/failed (%s) — reopening listing and resuming",
page_number, op_exc,
)
try:
page.close()
except Exception:
pass
page = None
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,
max_nav_pages=max(page_number, 500),
)
pages_since_rotation = 0
with _PageOpWatchdog(_PAGE_COLLECT_TIMEOUT_S, f"collect page {page_number} (after reopen)"):
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
except (ListingResumeError, PageOperationTimeoutError, PlaywrightError) as resume_exc:
logger.warning(
"Cannot recover at page %d (%s) — stopping pagination for this segment",
page_number, resume_exc,
)
break
pages_info.append({
"page_number": page_result.page_number,
"links_found": len(page_result.vehicle_links),
})
pages_since_rotation += 1
page_urls = self._extract_page_urls(
page_result,
@@ -604,6 +804,7 @@ class IAAIScraper:
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
max_nav_pages=max(page_number, 500),
)
page_result = self.listing_collector.collect_current_page(page, page_number=page_number)
page_urls = self._extract_page_urls(
@@ -694,9 +895,21 @@ class IAAIScraper:
or not page_result.next_page_detected
):
break
if not self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1):
try:
with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"):
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1)
except (PageOperationTimeoutError, PlaywrightError) as nav_exc:
logger.warning(
"go_to_next_page stalled/failed at page %d (%s) — will reopen",
page_number + 1, nav_exc,
)
next_ok = False
if not next_ok:
logger.warning("Failed to navigate to page %d — reopening listing and resuming", page_number + 1)
page.close()
try:
page.close()
except Exception:
pass
page = None
try:
page, _ = self._reopen_listing_and_resume(
@@ -706,7 +919,9 @@ class IAAIScraper:
listing_url=listing_url,
year_min=year_min,
year_max=year_max,
max_nav_pages=max(page_number + 1, 500),
)
pages_since_rotation = 0
except ListingResumeError as resume_exc:
logger.warning(
"Cannot resume at page %d (%s) — treating pagination as exhausted",
@@ -717,6 +932,35 @@ class IAAIScraper:
if pending_urls:
_process_pending_batch(pending_urls, cars_upserted + cars_failed)
# Финальный drain последнего in-flight батча.
_drain_pending_future()
# Обрабатываем отложенные browser fallback'и (Playwright теперь свободен).
if deferred_fallbacks:
fb_pages = min(len(deferred_fallbacks), self._FALLBACK_PARALLEL_PAGES)
logger.info(
"Processing %d deferred browser fallbacks (%d parallel pages)",
len(deferred_fallbacks), fb_pages,
)
try:
fb_results = self._browser_fallback_parallel(
deferred_fallbacks, len(deferred_fallbacks), fb_pages,
)
fb_records = fb_results.get("records", [])
if fb_records:
try:
upsert_result = self.persistence.upsert_cars_batch(fb_records)
cars_upserted += upsert_result["inserted"] + upsert_result["updated"]
images_upserted += upsert_result["images_upserted"]
except Exception as exc:
logger.error("Deferred fallback upsert failed: %s", exc)
cars_failed += len(fb_records)
cars_failed += fb_results.get("cars_failed", 0)
failures.extend(fb_results.get("failures", []))
except Exception as fb_exc:
logger.error("Deferred fallback batch failed: %s", fb_exc)
cars_failed += len(deferred_fallbacks)
logger.warning(
"sync_listing final progress: pages=%d, discovered=%d, upserted=%d, failed=%d",
len(pages_info),
@@ -725,6 +969,12 @@ class IAAIScraper:
cars_failed,
)
finally:
# Ensure no future is left dangling on exception path.
try:
_drain_pending_future()
except Exception:
pass
batch_executor.shutdown(wait=True)
if page is not None:
page.close()
@@ -979,7 +1229,8 @@ class IAAIScraper:
"""Извлекает IAAI inventoryView.attributes из HTML без браузера.
Использует json.JSONDecoder.raw_decode для быстрого поиска JSON
вместо посимвольного сканирования скобок.
вместо посимвольного сканирования скобок. Если получившийся срез
целиком отдаётся декодеру — пробуем orjson (3-5× быстрее stdlib).
"""
search = "inventoryView"
pos = html.find(search)
@@ -996,11 +1247,24 @@ class IAAIScraper:
if content_start < 0:
return None
decoder = json.JSONDecoder()
try:
obj, _ = decoder.raw_decode(html, content_start)
except (json.JSONDecodeError, ValueError):
return None
# Быстрый путь: вырезаем всё до </script>, пытаемся orjson; иначе stdlib.
obj: dict | None = None
if _HAS_ORJSON:
script_end = html.find("</script>", content_start)
if script_end > content_start:
# Подрезаем хвост до последней закрывающей скобки JSON.
snippet = html[content_start:script_end].rstrip().rstrip(";")
try:
obj = _orjson.loads(snippet)
except Exception:
obj = None
if obj is None:
decoder = json.JSONDecoder()
try:
obj, _ = decoder.raw_decode(html, content_start)
except (json.JSONDecodeError, ValueError):
return None
iv = obj.get("inventoryView")
if not iv or not iv.get("attributes"):
@@ -1118,11 +1382,19 @@ class IAAIScraper:
cookie_header: str,
user_agent: str,
) -> CarRecord:
"""Быстрый HTTP-запрос через urllib3 (connection pooling / keep-alive)."""
"""Быстрый HTTP-запрос через urllib3 (connection pooling / keep-alive).
Заголовок Accept-Encoding включает gzip/deflate/br — IAAI отдаёт HTML
со сжатием 5-10×, что снижает байты по сети и время декода.
urllib3 сам прозрачно распаковывает gzip и deflate; для br нужен
пакет `brotli` (он в зависимостях). Если brotli не установлен — он
просто не появится в Accept-Encoding и сервер ответит gzip.
"""
headers = {
"User-Agent": user_agent,
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Accept-Encoding": self._ACCEPT_ENCODING,
"Cache-Control": "no-cache",
"Pragma": "no-cache",
"Connection": "keep-alive",
@@ -1131,7 +1403,10 @@ class IAAIScraper:
if cookie_header:
headers["Cookie"] = cookie_header
response = self._http_pool.request("GET", vehicle_url, headers=headers)
# decode_content=True — urllib3 сам распакует gzip/deflate/br.
response = self._http_pool.request(
"GET", vehicle_url, headers=headers, decode_content=True,
)
if response.status >= 400:
raise RuntimeError(f"HTTP fetch failed for {vehicle_url}: {response.status}")
html = response.data.decode("utf-8", errors="ignore")
@@ -1205,6 +1480,10 @@ class IAAIScraper:
_HTTP_RETRY_DELAYS = (0.4, 1.0) # Задержки между повторами
_FALLBACK_PARALLEL_PAGES = 4 # Параллельных вкладок для fallback
_FALLBACK_NAV_TIMEOUT_MS = 8000 # Таймаут навигации fallback
# gzip+deflate: нативно в urllib3 (zlib, ~0 CPU). Brotli отключён намеренно —
# python-brotli декодирует на CPU и под 32 параллельными потоками становится
# узким местом (хуже, чем гонять чуть больше байт по сети).
_ACCEPT_ENCODING = "gzip, deflate"
def _browser_fallback_parallel(
self,
@@ -1359,11 +1638,15 @@ class IAAIScraper:
"protection_events": protection_events,
}
cookie_header_override: str | None = None,
def sync_batch(
self,
vehicle_urls: list[str],
lane: str = "iaai_cars",
parallel_tabs: int | None = None,
*,
skip_browser_fallback: bool = False,
cookie_header_override: str | None = None,
) -> dict:
trace_id = self._new_trace_id("sync-batch")
started_at = time.perf_counter()
@@ -1374,7 +1657,13 @@ class IAAIScraper:
cars_failed = 0
images_upserted = 0
records: list[CarRecord] = []
failures: list[dict[str, str]] = []
# cookie_header_override передаётся из главного потока при pipeline-вызове —
# вычисление cookie_header из background thread'а ломает Playwright sync greenlet.
cookie_header = (
cookie_header_override
if cookie_header_override is not None
else self._cookie_header_for_context()
http_successes = 0
http_fallbacks = 0
protection_events = 0
@@ -1382,17 +1671,25 @@ class IAAIScraper:
total = len(vehicle_urls)
logger.info("sync_batch: %d vehicles, %d parallel workers", total, num_workers)
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"
# cookie_header_override передаётся из главного потока при pipeline-вызове —
# вычисление cookie_header из background thread'а ломает Playwright sync greenlet.
cookie_header = (
cookie_header_override
if cookie_header_override is not None
else self._cookie_header_for_context()
)
# UA должен совпадать с UA Playwright-контекста, иначе Imperva видит
# «cookies от Chrome, ходим под Firefox» — типичный сигнал бота.
user_agent = self.settings.fingerprint.user_agent
# ── Фаза 1: HTTP fast-path ──
fallback_urls: list[tuple[str, int]] = []
def _fetch_one_with_retry(idx_url: tuple[int, str]) -> tuple[int, CarRecord | Exception]:
idx, url = idx_url
# Небольшой анти-бёрст джиттер: разносим старт запросов в пуле,
# чтобы N воркеров не били в одну TLS-/PX-волну.
time.sleep(random.uniform(0.0, 0.15))
last_exc: Exception | None = None
for attempt in range(1 + self._HTTP_MAX_RETRIES):
try:
@@ -1431,16 +1728,28 @@ class IAAIScraper:
)
# ── Фаза 2: браузерный fallback ──
# При skip_browser_fallback=True (pipeline-режим) возвращаем fallback_urls
# вызывающему — Playwright занят сбором листинга в главном потоке,
# его нельзя трогать из фонового executor'а.
deferred_fallback_urls: list[tuple[str, int]] = []
if fallback_urls:
logger.info(
"Fallback browser mode: %d/%d URLs, single page",
len(fallback_urls), total,
)
fb_results = self._browser_fallback_parallel(fallback_urls, total, 1)
records.extend(fb_results["records"])
cars_failed += fb_results["cars_failed"]
protection_events += fb_results["protection_events"]
failures.extend(fb_results["failures"])
if skip_browser_fallback:
deferred_fallback_urls = fallback_urls
logger.info(
"Deferring browser fallback: %d/%d URLs (pipeline mode)",
len(fallback_urls), total,
)
else:
fb_pages = min(len(fallback_urls), self._FALLBACK_PARALLEL_PAGES)
logger.info(
"Fallback browser mode: %d/%d URLs, %d parallel pages",
len(fallback_urls), total, fb_pages,
)
fb_results = self._browser_fallback_parallel(fallback_urls, total, fb_pages)
records.extend(fb_results["records"])
cars_failed += fb_results["cars_failed"]
protection_events += fb_results["protection_events"]
failures.extend(fb_results["failures"])
# ── Фаза 3: запись в БД ──
if records:
@@ -1475,6 +1784,7 @@ class IAAIScraper:
"protection_events": protection_events,
"elapsed_seconds": round(time.perf_counter() - started_at, 3),
"failures": failures,
"deferred_fallback_urls": deferred_fallback_urls,
}
def sync_listing(

View File

@@ -608,7 +608,10 @@ def sync_listing_task(
)
if force_bootstrap_full_scan:
_set_full_scan_done(redis_client, False)
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False)
# SoftTimeLimit считаем failure для circuit breaker:
# если сегмент систематически не укладывается в time limit,
# нельзя крутить followup бесконечно.
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=True)
# Partial progress уже записан в БД через finish_sync_run.
# Не retry — следующий запуск продолжит обработку по расписанию.
return {