1472 lines
62 KiB
Python
1472 lines
62 KiB
Python
import json
|
||
import logging
|
||
import re
|
||
import time
|
||
from collections import Counter
|
||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||
from dataclasses import dataclass, field, replace
|
||
from pathlib import Path
|
||
from typing import Any
|
||
from urllib.error import HTTPError, URLError
|
||
from urllib.parse import urlencode, quote
|
||
from urllib.request import Request, urlopen
|
||
|
||
import urllib3
|
||
|
||
from .core.config import Settings
|
||
from .core.runtime_config import FiltersConfig
|
||
from .storage.db import PersistenceService
|
||
from .storage.schemas import CarRecord, ImageRecord
|
||
from .translations import (
|
||
BODY_TYPE_MAP,
|
||
BRAND_ENGLISH_ALIASES,
|
||
BRAND_TRANSLATIONS,
|
||
CITY_TO_COUNTRY,
|
||
COLOR_TRANSLATIONS,
|
||
FUEL_TYPE_MAP,
|
||
GEARBOX_MAP_KO,
|
||
MODEL_BODY_TYPE_MAP,
|
||
MODEL_TRANSLATIONS,
|
||
SELL_TYPE_MAP,
|
||
expand_allowed_brands,
|
||
)
|
||
|
||
logger = logging.getLogger("encar_scraper")
|
||
|
||
ENCAR_LISTING_API = "https://api.encar.com/search/car/list/general"
|
||
ENCAR_DETAIL_API_TEMPLATE = "https://api.encar.com/v1/readside/vehicle/{vehicle_id}"
|
||
ENCAR_PHOTO_API_TEMPLATE = "https://api.encar.com/legacy/usedcar/sale/car/photo/{vehicle_id}"
|
||
# Batch endpoint из фронтенд JS (fem.encar.com) — до 20 машин за запрос, все фото
|
||
ENCAR_BATCH_VEHICLES_API = "https://api.encar.com/v1/readside/vehicles"
|
||
BATCH_VEHICLES_CHUNK_SIZE = 20
|
||
ENCAR_DETAIL_URL_TEMPLATE = "https://www.encar.com/dc/dc_cardetailview.do?carid={vehicle_id}"
|
||
ENCAR_IMAGE_BASE = "https://ci.encar.com"
|
||
VEHICLE_ID_RE = re.compile(r"(?:carid|vehicleId)=?(\d+)")
|
||
PHOTO_VEHICLE_ID_RE = re.compile(r"/(\d+)_\d+\.(?:jpg|jpeg|png|webp)$", re.IGNORECASE)
|
||
|
||
# Максимальный номер фото для пробинга и допустимые промахи подряд
|
||
PHOTO_PROBE_MAX = 50
|
||
PHOTO_PROBE_MAX_CONSECUTIVE_MISS = 7
|
||
|
||
# Типы листинга
|
||
CAR_TYPE_ALL = "A" # все (отечественные + импорт)
|
||
CAR_TYPE_DOMESTIC = "Y" # только отечественные
|
||
CAR_TYPE_IMPORT = "N" # только импорт
|
||
CAR_TYPE_ELECTRIC = "G" # электро/гибрид
|
||
|
||
# Regex для парсинга объёма двигателя из Badge
|
||
ENGINE_VOLUME_RE = re.compile(r"(\d+\.\d+)\s*(?:T|터보)?")
|
||
# Regex для европейских кодов: "530i", "C200", "A220", "GLC300", "740d", "M340i", "S63", "120i", "E220d" и т.д.
|
||
# Ищем 3-значные числа в начале или после буквы: 530, 200, 220, 300, 740, 340 ...
|
||
EURO_ENGINE_CODE_RE = re.compile(
|
||
r"(?:^|[A-Za-z])(\d{3})(?:[a-zA-Z]|$|\s)" # C200, A220, 530i, 740d, etc.
|
||
)
|
||
# Маппинг европейских кодов → объём (cc). Первые 1-2 цифры = литраж
|
||
EURO_CODE_ENGINE_MAP: dict[str, int] = {
|
||
# BMW
|
||
"118": 1500, "120": 2000, "125": 2000, "128": 2000,
|
||
"218": 1500, "220": 2000, "225": 2000, "228": 2000, "230": 2000, "235": 2000,
|
||
"318": 1500, "320": 2000, "325": 2500, "328": 2000, "330": 2000, "335": 3000, "340": 3000,
|
||
"418": 2000, "420": 2000, "425": 2000, "428": 2000, "430": 2000, "435": 3000, "440": 3000,
|
||
"520": 2000, "525": 2500, "528": 2000, "530": 2000, "535": 3000, "540": 3000,
|
||
"620": 2000, "630": 3000, "640": 3000, "650": 4400,
|
||
"725": 2500, "730": 3000, "740": 3000, "745": 3000, "750": 4400, "760": 6600,
|
||
"X20": 2000, "X25": 2500, "X28": 2000, "X30": 2000, "X35": 3000, "X40": 3000, "X50": 4400,
|
||
# Mercedes-Benz
|
||
"160": 1600, "180": 1600, "200": 2000, "220": 2000, "250": 2000, "260": 2000,
|
||
"280": 3000, "300": 3000, "350": 3500, "400": 3000, "450": 3000,
|
||
"500": 4000, "550": 4700, "560": 5500, "600": 6000, "630": 6300,
|
||
"43": 3000, "53": 3000, "55": 5500, "63": 4000,
|
||
# Audi
|
||
"25": 1400, "30": 2000, "35": 2000, "40": 2000, "45": 2000, "50": 3000, "55": 3000,
|
||
}
|
||
|
||
# Маппинг для извлечения engine_volume из числовых кодов типа "M16 GDI"
|
||
BADGE_SHORTCODE_RE = re.compile(r"M(\d{2})\s") # M16 GDI → 1600
|
||
|
||
@dataclass
|
||
class EncarFilters:
|
||
"""Фильтры для API запросов."""
|
||
car_type: str = CAR_TYPE_ALL # A=все, Y=корейские, N=импорт
|
||
manufacturer: str | None = None # BMW, 현대, etc.
|
||
model: str | None = None
|
||
year_from: int | None = None
|
||
year_to: int | None = None
|
||
price_from: int | None = None # в 만원 (10k KRW)
|
||
price_to: int | None = None
|
||
mileage_max: int | None = None
|
||
|
||
def build_query(self) -> str:
|
||
"""Строит строку запроса для API."""
|
||
parts = ["Hidden.N"]
|
||
|
||
if self.car_type:
|
||
parts.append(f"CarType.{self.car_type}")
|
||
|
||
if self.manufacturer:
|
||
parts.append(f"Manufacturer.{self.manufacturer}")
|
||
|
||
if self.model:
|
||
parts.append(f"Model.{self.model}")
|
||
|
||
# Год: Year в формате YYYYMM (напр. 202001)
|
||
# year_from/year_to > 9999 → уже в формате YYYYMM, иначе YYYY → прибавляем 01/12
|
||
if self.year_from or self.year_to:
|
||
if self.year_from and self.year_from > 9999:
|
||
year_min = str(self.year_from)
|
||
else:
|
||
year_min = f"{self.year_from}01" if self.year_from else ""
|
||
if self.year_to and self.year_to > 9999:
|
||
year_max = str(self.year_to)
|
||
else:
|
||
year_max = f"{self.year_to}12" if self.year_to else ""
|
||
parts.append(f"Year.range({year_min}..{year_max})")
|
||
|
||
# Цена в 만원
|
||
if self.price_from or self.price_to:
|
||
p_min = str(self.price_from) if self.price_from else ""
|
||
p_max = str(self.price_to) if self.price_to else ""
|
||
parts.append(f"Price.range({p_min}..{p_max})")
|
||
|
||
# Пробег
|
||
if self.mileage_max:
|
||
parts.append(f"Mileage.range(..{self.mileage_max})")
|
||
|
||
return "(And." + "._." .join(parts) + ".)"
|
||
|
||
|
||
@dataclass(slots=True, frozen=True)
|
||
class RuntimeFilterSpec:
|
||
price_min: int | None = None
|
||
price_max: int | None = None
|
||
mileage_min: int | None = None
|
||
mileage_max: int | None = None
|
||
models: frozenset[str] = frozenset()
|
||
years: frozenset[int] = frozenset()
|
||
body_types: frozenset[str] = frozenset()
|
||
colors: frozenset[str] = frozenset()
|
||
drives: frozenset[str] = frozenset()
|
||
gearboxes: frozenset[str] = frozenset()
|
||
exclude_models: frozenset[str] = frozenset()
|
||
exclude_years: frozenset[int] = frozenset()
|
||
exclude_body_types: frozenset[str] = frozenset()
|
||
|
||
|
||
@dataclass(slots=True)
|
||
class EncarMapper:
|
||
def map_to_car_record(self, vehicle_url: str, payload: dict[str, Any], probe_all_photos: bool = False) -> CarRecord:
|
||
vehicle_id = self._extract_vehicle_id(vehicle_url) or str(
|
||
payload.get("Id") or payload.get("vehicleId") or ""
|
||
)
|
||
origin_id = f"encar:{vehicle_id}"
|
||
parser_id = f"encar:{vehicle_id}"
|
||
origin_url = vehicle_url or ENCAR_DETAIL_URL_TEMPLATE.format(vehicle_id=vehicle_id)
|
||
|
||
brand = self._as_str(payload.get("Manufacturer") or payload.get("Brand") or "UNKNOWN")
|
||
# Перевод бренда на русский
|
||
brand = BRAND_TRANSLATIONS.get(brand, brand)
|
||
|
||
# Model + Badge + BadgeDetail = полное название
|
||
model_raw = self._as_str(payload.get("Model") or "UNKNOWN")
|
||
badge = self._as_str(payload.get("Badge") or "")
|
||
badge_detail = self._as_str(payload.get("BadgeDetail") or "")
|
||
|
||
# Полное описание модели: "Model / Badge BadgeDetail"
|
||
model_parts = [model_raw]
|
||
if badge:
|
||
model_parts.append(badge)
|
||
if badge_detail and badge_detail not in badge:
|
||
model_parts.append(badge_detail)
|
||
model = " / ".join(model_parts) if len(model_parts) > 1 else model_parts[0]
|
||
model = self._translate_text(model)
|
||
|
||
year = self._to_year(payload.get("FormYear") or payload.get("Year"))
|
||
# Price в 만원 (10 000 KRW) → переводим в KRW
|
||
price_man_won = self._to_float(payload.get("Price"))
|
||
price = int(price_man_won * 10_000) if price_man_won is not None else None
|
||
mileage = self._to_int(payload.get("Mileage"))
|
||
|
||
# Цвет: листинг не содержит цвет, проверяем все варианты
|
||
color_raw = self._as_str(
|
||
payload.get("ExteriorColor") or payload.get("Color") or ""
|
||
)
|
||
color = self._translate_color(color_raw) if color_raw else "other"
|
||
|
||
currency = "KRW"
|
||
country = "KR"
|
||
|
||
# Топливо и коробка передач
|
||
fuel_raw = self._as_str(payload.get("FuelType") or "")
|
||
green_type = self._as_str(payload.get("GreenType") or "")
|
||
|
||
# Gearbox: Encar листинг не даёт Transmission, определяем по топливу
|
||
if fuel_raw == "전기":
|
||
# Чистый электромобиль
|
||
gearbox = "EV"
|
||
else:
|
||
# Гибрид (가솔린+전기, 디젤+전기) или ДВС → AT
|
||
gearbox_raw = self._as_str(payload.get("Transmission") or "")
|
||
if gearbox_raw:
|
||
gearbox = GEARBOX_MAP_KO.get(gearbox_raw) or GEARBOX_MAP_KO.get(gearbox_raw.lower()) or "AT"
|
||
else:
|
||
gearbox = "AT"
|
||
|
||
sell_type_raw = self._as_str(payload.get("SellType") or "")
|
||
selling_type = SELL_TYPE_MAP.get(sell_type_raw, "STOCK")
|
||
|
||
# Парсинг drive из Badge
|
||
drive = self._parse_drive(badge, model_raw)
|
||
|
||
# Парсинг body_type из Model (по ключевым словам в модели)
|
||
body_type_raw = self._as_str(payload.get("BodyType") or "")
|
||
if body_type_raw:
|
||
body_type = BODY_TYPE_MAP.get(body_type_raw, "OTHER")
|
||
else:
|
||
body_type = self._detect_body_type(model_raw)
|
||
|
||
# Парсинг engine_volume из Badge (напр. "가솔린 2.5 터보 AWD" → 2500)
|
||
engine_volume_raw = payload.get("Displacement") or payload.get("EngineCc")
|
||
if engine_volume_raw:
|
||
engine_volume = self._to_int(engine_volume_raw)
|
||
else:
|
||
engine_volume = self._parse_engine_volume(badge)
|
||
|
||
# "엔카진단" флаги
|
||
service_marks = payload.get("ServiceMark") or []
|
||
condition_list = payload.get("Condition") or []
|
||
has_inspection = bool(condition_list) or (
|
||
"EncarDiagnosisP1" in service_marks or "EncarDiagnosisP2" in service_marks
|
||
)
|
||
|
||
photo_prefix = self._as_str(payload.get("Photo") or "")
|
||
images = self._map_images(
|
||
payload.get("Photos") or [],
|
||
probe_all=probe_all_photos,
|
||
photo_prefix=photo_prefix or None,
|
||
)
|
||
|
||
return CarRecord(
|
||
parser_id=parser_id,
|
||
brand=brand,
|
||
model=model,
|
||
year=year,
|
||
price=price,
|
||
currency=currency,
|
||
mileage=mileage or 0,
|
||
country=country,
|
||
is_sold=False,
|
||
color=color or "other",
|
||
drive=drive,
|
||
gearbox=gearbox,
|
||
steering_wheel="LEFT", # Корея — левый руль
|
||
body_type=body_type,
|
||
engine_volume=engine_volume,
|
||
selling_type=selling_type,
|
||
one_owner=False,
|
||
new_car=False,
|
||
is_hidden=False,
|
||
origin="ENCAR",
|
||
origin_url=origin_url,
|
||
origin_id=origin_id,
|
||
is_damaged=False,
|
||
evaluation="diagnosis" if has_inspection else None,
|
||
non_smoking=False,
|
||
rental=False,
|
||
repair_history=False,
|
||
slug=self._build_slug(brand, model, vehicle_id),
|
||
images=images,
|
||
)
|
||
|
||
def _extract_vehicle_id(self, vehicle_url: str | None) -> str | None:
|
||
if not vehicle_url:
|
||
return None
|
||
match = VEHICLE_ID_RE.search(vehicle_url)
|
||
return match.group(1) if match else None
|
||
|
||
def _build_slug(self, brand: str, model: str, vehicle_id: str) -> str:
|
||
raw = re.sub(r"[^a-z0-9]+", "-", f"{brand}-{model}-{vehicle_id}".lower()).strip("-")
|
||
return raw or f"encar-{vehicle_id}"
|
||
|
||
def _as_str(self, value: Any) -> str:
|
||
return str(value).strip() if value is not None else ""
|
||
|
||
def _translate_text(self, text: str) -> str:
|
||
translated = self._as_str(text)
|
||
for source, target in sorted(MODEL_TRANSLATIONS.items(), key=lambda item: len(item[0]), reverse=True):
|
||
translated = translated.replace(source, target)
|
||
translated = re.sub(r"\s+", " ", translated).strip()
|
||
return translated
|
||
|
||
def _translate_color(self, color_raw: str) -> str:
|
||
value = self._as_str(color_raw)
|
||
if not value:
|
||
return "other"
|
||
return COLOR_TRANSLATIONS.get(value, self._translate_text(value).lower() or "other")
|
||
|
||
def _to_float(self, value: Any) -> float | None:
|
||
if value is None:
|
||
return None
|
||
try:
|
||
return float(value)
|
||
except (ValueError, TypeError):
|
||
digits = re.sub(r"[^0-9.]", "", str(value))
|
||
return float(digits) if digits else None
|
||
|
||
def _to_int(self, value: Any) -> int | None:
|
||
f = self._to_float(value)
|
||
return int(f) if f is not None else None
|
||
|
||
def _to_year(self, value: Any) -> int | None:
|
||
"""Парсинг года.
|
||
- FormYear: "2021" → 2021
|
||
- Year: 202102.0 → YYYYMM формат → берём YYYY часть
|
||
"""
|
||
if value is None:
|
||
return None
|
||
try:
|
||
s = str(value).strip().split(".")[0] # убираем .0
|
||
if len(s) >= 6:
|
||
# YYYYMM формат
|
||
return int(s[:4])
|
||
if len(s) == 4:
|
||
return int(s)
|
||
except (ValueError, TypeError):
|
||
pass
|
||
return None
|
||
|
||
def _map_images(self, photos: Any, probe_all: bool = False, photo_prefix: str | None = None) -> list[ImageRecord]:
|
||
if not isinstance(photos, list):
|
||
photos = []
|
||
images: list[ImageRecord] = []
|
||
known_orders: set[int] = set()
|
||
url_prefix: str | None = None
|
||
|
||
for raw in photos:
|
||
if not isinstance(raw, dict):
|
||
continue
|
||
location = raw.get("location") or raw.get("Location") or ""
|
||
if not location:
|
||
continue
|
||
url = self._build_image_url(str(location))
|
||
order = int(raw.get("ordering") or 0)
|
||
images.append(ImageRecord(fullres_image=url, preview_image=url, order_index=order))
|
||
known_orders.add(order)
|
||
if url_prefix is None and "_" in location:
|
||
url_prefix = self._build_image_url(location.rsplit("_", 1)[0])
|
||
|
||
# Если Photos не дал prefix, используем Photo (путь-префикс из premium API)
|
||
if url_prefix is None and photo_prefix:
|
||
url_prefix = self._build_image_url(photo_prefix.rstrip("_"))
|
||
|
||
if probe_all and url_prefix:
|
||
extra = self._probe_additional_photos(url_prefix, known_orders)
|
||
images.extend(extra)
|
||
|
||
images.sort(key=lambda item: item.order_index)
|
||
return images
|
||
|
||
@staticmethod
|
||
def _head_check(url: str) -> bool:
|
||
"""Проверка существования фото через HEAD-запрос."""
|
||
try:
|
||
req = Request(url, method="HEAD", headers={
|
||
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
|
||
"Referer": "https://www.encar.com/",
|
||
})
|
||
with urlopen(req, timeout=3) as resp:
|
||
return resp.status == 200
|
||
except (HTTPError, URLError, OSError):
|
||
return False
|
||
|
||
@staticmethod
|
||
def _head_check_pooled(pool: urllib3.HTTPSConnectionPool, path: str) -> bool:
|
||
"""HEAD-запрос через connection pool (быстрее за счёт reuse TCP/SSL)."""
|
||
try:
|
||
resp = pool.request(
|
||
"HEAD", path,
|
||
headers={
|
||
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
|
||
"Referer": "https://www.encar.com/",
|
||
},
|
||
timeout=2.0,
|
||
retries=False,
|
||
)
|
||
return resp.status == 200
|
||
except Exception:
|
||
return False
|
||
|
||
def _probe_additional_photos(self, url_prefix: str, known_orders: set[int]) -> list[ImageRecord]:
|
||
"""Пробинг дополнительных фото параллельно через HEAD-запросы."""
|
||
candidates = []
|
||
for i in range(1, PHOTO_PROBE_MAX + 1):
|
||
if i in known_orders:
|
||
continue
|
||
candidates.append((i, f"{url_prefix}_{i:03d}.jpg"))
|
||
|
||
if not candidates:
|
||
return []
|
||
|
||
found: list[ImageRecord] = []
|
||
with ThreadPoolExecutor(max_workers=10) as pool:
|
||
futures = {pool.submit(self._head_check, url): (order, url) for order, url in candidates}
|
||
for future in as_completed(futures):
|
||
order, url = futures[future]
|
||
if future.result():
|
||
found.append(ImageRecord(fullres_image=url, preview_image=url, order_index=order))
|
||
|
||
found.sort(key=lambda x: x.order_index)
|
||
|
||
# Отсекаем хвост после PHOTO_PROBE_MAX_CONSECUTIVE_MISS промахов подряд
|
||
all_found_orders = known_orders | {img.order_index for img in found}
|
||
cutoff = PHOTO_PROBE_MAX + 1
|
||
consecutive_miss = 0
|
||
for i in range(1, PHOTO_PROBE_MAX + 1):
|
||
if i in all_found_orders:
|
||
consecutive_miss = 0
|
||
else:
|
||
consecutive_miss += 1
|
||
if consecutive_miss >= PHOTO_PROBE_MAX_CONSECUTIVE_MISS:
|
||
cutoff = i - PHOTO_PROBE_MAX_CONSECUTIVE_MISS + 1
|
||
break
|
||
|
||
return [img for img in found if img.order_index < cutoff]
|
||
|
||
def _build_image_url(self, location: str) -> str:
|
||
if location.startswith("http"):
|
||
return location
|
||
if location.startswith("/"):
|
||
return f"{ENCAR_IMAGE_BASE}{location}"
|
||
return f"{ENCAR_IMAGE_BASE}/{location}"
|
||
|
||
def _parse_drive(self, badge: str, model: str) -> str | None:
|
||
"""Парсит тип привода из Badge и Model."""
|
||
text = f"{badge} {model}".lower()
|
||
if "awd" in text or "4wd" in text or "사륜" in text or "4륜" in text:
|
||
return "4WD"
|
||
elif "2wd" in text or "fwd" in text or "이륜" in text or "전륜" in text:
|
||
return "FWD"
|
||
elif "rwd" in text or "후륜" in text:
|
||
return "RWD"
|
||
return None
|
||
|
||
def _parse_engine_volume(self, badge: str) -> int | None:
|
||
"""Парсит объём двигателя из Badge.
|
||
|
||
Примеры:
|
||
"가솔린 2.5 터보 AWD" → 2500
|
||
"1.6 GDI 유니크" → 1600
|
||
"3.8 4WD" → 3800
|
||
"530i 럭셔리" → 2000 (BMW code)
|
||
"C200 AMG Line" → 2000 (Mercedes code)
|
||
"롱레인지 AWD" → None (электро, нет объёма)
|
||
"""
|
||
if not badge:
|
||
return None
|
||
|
||
# 1. Прямой литраж: "2.5", "1.6", "3.8"
|
||
m = ENGINE_VOLUME_RE.search(badge)
|
||
if m:
|
||
vol = float(m.group(1))
|
||
if 0.5 <= vol <= 8.0:
|
||
return int(vol * 1000)
|
||
|
||
# 2. Европейские коды: "530i", "C200", "GLC300", "M340i", "E220d"
|
||
euro_match = EURO_ENGINE_CODE_RE.search(badge)
|
||
if euro_match:
|
||
code = euro_match.group(1)
|
||
if code in EURO_CODE_ENGINE_MAP:
|
||
return EURO_CODE_ENGINE_MAP[code]
|
||
|
||
# 3. Корейские shortcode: "M16 GDI" → 1600
|
||
m_short = BADGE_SHORTCODE_RE.search(badge)
|
||
if m_short:
|
||
code = int(m_short.group(1))
|
||
if 10 <= code <= 60:
|
||
return code * 100
|
||
|
||
return None
|
||
|
||
def _detect_body_type(self, model_raw: str) -> str:
|
||
"""Определяет тип кузова из названия модели."""
|
||
if not model_raw:
|
||
return "OTHER"
|
||
# Проверяем каждый ключ из MODEL_BODY_TYPE_MAP
|
||
for keyword, body in MODEL_BODY_TYPE_MAP.items():
|
||
if keyword in model_raw:
|
||
return body
|
||
return "OTHER"
|
||
|
||
|
||
class EncarScraper:
|
||
DEFAULT_PAGE_SIZE = 1000 # API поддерживает до 1000
|
||
|
||
def __init__(self, settings: Settings | None = None, persistence: PersistenceService | None = None) -> None:
|
||
self.settings = settings or Settings()
|
||
self.persistence = persistence or PersistenceService(self.settings)
|
||
self.car_mapper = EncarMapper()
|
||
self._batch_pool: urllib3.HTTPSConnectionPool | None = None
|
||
|
||
def _get_batch_pool(self) -> urllib3.HTTPSConnectionPool:
|
||
"""Возвращает переиспользуемый connection pool для batch photo API."""
|
||
if self._batch_pool is None:
|
||
self._batch_pool = urllib3.HTTPSConnectionPool(
|
||
"api.encar.com",
|
||
port=443,
|
||
maxsize=30,
|
||
block=True,
|
||
headers={
|
||
"User-Agent": self.settings.fingerprint.user_agent,
|
||
"Accept": "application/json",
|
||
},
|
||
cert_reqs="CERT_REQUIRED",
|
||
)
|
||
return self._batch_pool
|
||
|
||
@staticmethod
|
||
def _compile_runtime_filters(runtime_filters: FiltersConfig | None) -> RuntimeFilterSpec | None:
|
||
if runtime_filters is None:
|
||
return None
|
||
|
||
spec = RuntimeFilterSpec(
|
||
price_min=runtime_filters.price_min,
|
||
price_max=runtime_filters.price_max,
|
||
mileage_min=runtime_filters.mileage_min,
|
||
mileage_max=runtime_filters.mileage_max,
|
||
models=frozenset(m for m in runtime_filters.models if m),
|
||
years=frozenset(int(y) for y in runtime_filters.years if y is not None),
|
||
body_types=frozenset(b for b in runtime_filters.body_types if b),
|
||
colors=frozenset(c for c in runtime_filters.colors if c),
|
||
drives=frozenset(d for d in runtime_filters.drives if d),
|
||
gearboxes=frozenset(g for g in runtime_filters.gearboxes if g),
|
||
exclude_models=frozenset(m for m in runtime_filters.exclude_models if m),
|
||
exclude_years=frozenset(int(y) for y in runtime_filters.exclude_years if y is not None),
|
||
exclude_body_types=frozenset(b for b in runtime_filters.exclude_body_types if b),
|
||
)
|
||
|
||
if not any((
|
||
spec.price_min is not None,
|
||
spec.price_max is not None,
|
||
spec.mileage_min is not None,
|
||
spec.mileage_max is not None,
|
||
spec.models,
|
||
spec.years,
|
||
spec.body_types,
|
||
spec.colors,
|
||
spec.drives,
|
||
spec.gearboxes,
|
||
spec.exclude_models,
|
||
spec.exclude_years,
|
||
spec.exclude_body_types,
|
||
)):
|
||
return None
|
||
|
||
return spec
|
||
|
||
@staticmethod
|
||
def _matches_any_token(value: str, tokens: frozenset[str]) -> bool:
|
||
if not tokens:
|
||
return True
|
||
return any(token in value for token in tokens if token)
|
||
|
||
def _record_passes_runtime_filters(
|
||
self,
|
||
record: CarRecord,
|
||
runtime_filters: RuntimeFilterSpec | None,
|
||
) -> bool:
|
||
if runtime_filters is None:
|
||
return True
|
||
|
||
model_lower = (record.model or "").lower()
|
||
body_type_lower = (record.body_type or "").lower()
|
||
color_lower = (record.color or "").lower()
|
||
drive_lower = (record.drive or "").lower()
|
||
gearbox_lower = (record.gearbox or "").lower()
|
||
|
||
if runtime_filters.price_min is not None and (record.price is None or record.price < runtime_filters.price_min):
|
||
return False
|
||
if runtime_filters.price_max is not None and (record.price is None or record.price > runtime_filters.price_max):
|
||
return False
|
||
if runtime_filters.mileage_min is not None and record.mileage < runtime_filters.mileage_min:
|
||
return False
|
||
if runtime_filters.mileage_max is not None and record.mileage > runtime_filters.mileage_max:
|
||
return False
|
||
|
||
if runtime_filters.models and not self._matches_any_token(model_lower, runtime_filters.models):
|
||
return False
|
||
if runtime_filters.exclude_models and self._matches_any_token(model_lower, runtime_filters.exclude_models):
|
||
return False
|
||
|
||
if runtime_filters.years and (record.year is None or record.year not in runtime_filters.years):
|
||
return False
|
||
if runtime_filters.exclude_years and record.year in runtime_filters.exclude_years:
|
||
return False
|
||
|
||
if runtime_filters.body_types and body_type_lower not in runtime_filters.body_types:
|
||
return False
|
||
if runtime_filters.exclude_body_types and body_type_lower in runtime_filters.exclude_body_types:
|
||
return False
|
||
|
||
if runtime_filters.colors and color_lower not in runtime_filters.colors:
|
||
return False
|
||
if runtime_filters.drives and drive_lower not in runtime_filters.drives:
|
||
return False
|
||
if runtime_filters.gearboxes and gearbox_lower not in runtime_filters.gearboxes:
|
||
return False
|
||
|
||
return True
|
||
|
||
def _canonical_vehicle_id_from_images(self, record: CarRecord) -> str | None:
|
||
current_vehicle_id = self._extract_vehicle_id(record.origin_url or "")
|
||
if not current_vehicle_id:
|
||
current_vehicle_id = record.origin_id.split(":", 1)[-1] if ":" in record.origin_id else record.origin_id
|
||
|
||
image_vehicle_ids: list[str] = []
|
||
for image in record.images:
|
||
url = image.fullres_image or ""
|
||
match = PHOTO_VEHICLE_ID_RE.search(url)
|
||
if match:
|
||
image_vehicle_ids.append(match.group(1))
|
||
|
||
if not image_vehicle_ids:
|
||
return None
|
||
|
||
canonical_vehicle_id, occurrences = Counter(image_vehicle_ids).most_common(1)[0]
|
||
if canonical_vehicle_id == current_vehicle_id:
|
||
return None
|
||
|
||
# Канонизируем только если mismatch подтверждается минимум двумя фото.
|
||
if occurrences < 2:
|
||
return None
|
||
|
||
return canonical_vehicle_id
|
||
|
||
def _apply_canonical_vehicle_id(self, record: CarRecord) -> bool:
|
||
canonical_vehicle_id = self._canonical_vehicle_id_from_images(record)
|
||
if not canonical_vehicle_id:
|
||
return False
|
||
|
||
record.parser_id = f"encar:{canonical_vehicle_id}"
|
||
record.origin_id = f"encar:{canonical_vehicle_id}"
|
||
record.origin_url = ENCAR_DETAIL_URL_TEMPLATE.format(vehicle_id=canonical_vehicle_id)
|
||
record.slug = self.car_mapper._build_slug(record.brand, record.model, canonical_vehicle_id)
|
||
return True
|
||
|
||
def collect_listing(
|
||
self,
|
||
limit: int | None = None,
|
||
page_size: int = DEFAULT_PAGE_SIZE,
|
||
filters: EncarFilters | None = None,
|
||
) -> dict[str, Any]:
|
||
"""Собирает список авто из API."""
|
||
filters = filters or EncarFilters()
|
||
vehicle_items: list[dict] = []
|
||
page = 0
|
||
total = None
|
||
|
||
logger.info("Collecting Encar listing with query: %s", filters.build_query())
|
||
|
||
while True:
|
||
offset = page * page_size
|
||
response = self._fetch_listing_page(offset=offset, limit=page_size, filters=filters)
|
||
results = response.get("SearchResults") or []
|
||
|
||
if total is None:
|
||
total = int(response.get("Count") or len(results))
|
||
logger.info("Total available: %d", total)
|
||
|
||
if not results:
|
||
break
|
||
|
||
for item in results:
|
||
vehicle_id = str(item.get("Id") or "")
|
||
if not vehicle_id:
|
||
continue
|
||
vehicle_items.append(item)
|
||
if limit is not None and len(vehicle_items) >= limit:
|
||
break
|
||
|
||
if limit is not None and len(vehicle_items) >= limit:
|
||
break
|
||
if len(results) < page_size:
|
||
break
|
||
|
||
page += 1
|
||
if page % 10 == 0:
|
||
logger.info("Collected %d items from %d pages...", len(vehicle_items), page)
|
||
|
||
logger.info("Finished collecting: %d items from %d pages", len(vehicle_items), page + 1)
|
||
|
||
return {
|
||
"items": vehicle_items,
|
||
"vehicle_urls": [ENCAR_DETAIL_URL_TEMPLATE.format(vehicle_id=item["Id"]) for item in vehicle_items],
|
||
"pages_collected": page + 1,
|
||
"items_collected": len(vehicle_items),
|
||
"total_available": total,
|
||
}
|
||
|
||
def scrape_vehicle_detail(self, vehicle_url: str) -> dict[str, Any]:
|
||
"""Получает детальную информацию по одному авто."""
|
||
vehicle_id = self._extract_vehicle_id(vehicle_url)
|
||
if not vehicle_id:
|
||
raise ValueError(f"Cannot extract Encar vehicle id from {vehicle_url}")
|
||
|
||
# Сначала пробуем получить из листинга (там больше данных)
|
||
payload = self._fetch_listing_item(vehicle_id)
|
||
|
||
# Если нет в листинге — берём из детального API
|
||
if payload is None:
|
||
payload = self._fetch_vehicle_detail_api(vehicle_id)
|
||
|
||
if payload is None:
|
||
raise ValueError(f"Encar vehicle {vehicle_id} not found")
|
||
|
||
record = self.car_mapper.map_to_car_record(vehicle_url, payload)
|
||
return {
|
||
"vehicle_id": vehicle_id,
|
||
"payload": payload,
|
||
"db_record": record.model_dump(mode="json"),
|
||
}
|
||
|
||
def init_db(self) -> dict[str, Any]:
|
||
self.persistence.create_tables()
|
||
return {"status": "ok"}
|
||
|
||
def sync_vehicle(self, vehicle_url: str, lane: str = "encar") -> dict[str, Any]:
|
||
"""Синхронизирует одно авто в БД."""
|
||
self.persistence.create_tables()
|
||
payload = self.scrape_vehicle_detail(vehicle_url)
|
||
record = CarRecord.model_validate(payload["db_record"])
|
||
result = self.persistence.upsert_car(record)
|
||
return {
|
||
"status": "success",
|
||
"lane": lane,
|
||
"vehicle_url": vehicle_url,
|
||
"origin_id": record.origin_id,
|
||
**result,
|
||
}
|
||
|
||
CHECKPOINT_KEY = "encar:sync:checkpoint"
|
||
CHECKPOINT_SAVE_INTERVAL = 10 # сохранять чекпоинт каждые N страниц
|
||
SHARD_MAX_ITEMS = 9500 # макс. элементов в одном шарде (API лимит ~10k)
|
||
|
||
# ------------------------------------------------------------------ sharding
|
||
|
||
def _get_listing_count(self, filters: EncarFilters) -> int:
|
||
"""Получить Count для фильтра без выгрузки данных."""
|
||
resp = self._fetch_listing_page(offset=0, limit=1, filters=filters)
|
||
return int(resp.get("Count") or 0)
|
||
|
||
def _build_sync_shards(self, base_filters: EncarFilters) -> list[EncarFilters]:
|
||
"""Авто-сегментация фильтров для обхода лимита пагинации API (~10k).
|
||
|
||
Быстрый алгоритм (~15 API calls вместо ~90):
|
||
1. total <= 9500 → 1 шард.
|
||
2. Разбиваем по car_type (Y/N).
|
||
3. Группируем старые годы декадами (2000-2009 → один шард).
|
||
4. Только годы > 9500 разбиваем на полугодия.
|
||
|
||
seen_vehicle_ids в sync_listing убирает дубли между шардами.
|
||
"""
|
||
MAX = self.SHARD_MAX_ITEMS
|
||
|
||
total = self._get_listing_count(base_filters)
|
||
if total <= MAX:
|
||
logger.info("Total %d ≤ %d — single shard", total, MAX)
|
||
return [base_filters]
|
||
|
||
logger.info("Total %d > %d — building shards", total, MAX)
|
||
|
||
if base_filters.car_type in (CAR_TYPE_ALL, None, ""):
|
||
car_types = ["Y", "N"]
|
||
else:
|
||
car_types = [base_filters.car_type]
|
||
|
||
shards: list[EncarFilters] = []
|
||
shard_items = 0
|
||
|
||
for ct in car_types:
|
||
ct_f = replace(base_filters, car_type=ct)
|
||
ct_count = self._get_listing_count(ct_f)
|
||
if ct_count == 0:
|
||
continue
|
||
if ct_count <= MAX:
|
||
shards.append(ct_f)
|
||
shard_items += ct_count
|
||
continue
|
||
|
||
# Сначала пробуем группировку: все годы до 2014 одним шардом
|
||
old_f = replace(ct_f, year_from=None, year_to=201412)
|
||
old_count = self._get_listing_count(old_f)
|
||
if old_count > 0:
|
||
if old_count <= MAX:
|
||
shards.append(old_f)
|
||
shard_items += old_count
|
||
else:
|
||
# Разбиваем на 2 части: до 2009 и 2010-2014
|
||
for ys, ye in [(None, 200912), (201001, 201412)]:
|
||
part_f = replace(ct_f, year_from=ys, year_to=ye)
|
||
part_c = self._get_listing_count(part_f)
|
||
if part_c > 0:
|
||
shards.append(part_f)
|
||
shard_items += part_c
|
||
|
||
# Годы 2015-2027: проверяем каждый
|
||
current_year = time.localtime().tm_year
|
||
for year in range(2015, current_year + 2):
|
||
yf = replace(ct_f, year_from=year, year_to=year)
|
||
yc = self._get_listing_count(yf)
|
||
if yc == 0:
|
||
continue
|
||
if yc <= MAX:
|
||
shards.append(yf)
|
||
shard_items += yc
|
||
else:
|
||
# Разбиваем на полугодия
|
||
for ms, me in [(year * 100 + 1, year * 100 + 6),
|
||
(year * 100 + 7, year * 100 + 12)]:
|
||
hf = replace(ct_f, year_from=ms, year_to=me)
|
||
hc = self._get_listing_count(hf)
|
||
if hc == 0:
|
||
continue
|
||
if hc <= MAX:
|
||
shards.append(hf)
|
||
shard_items += hc
|
||
else:
|
||
# Полугодие > MAX — разбиваем на кварталы
|
||
q_ranges = [(ms, ms + 2), (ms + 3, me)] if me - ms == 5 else [(ms, ms + 2), (ms + 3, me)]
|
||
for qs, qe in q_ranges:
|
||
qf = replace(ct_f, year_from=qs, year_to=qe)
|
||
qc = self._get_listing_count(qf)
|
||
if qc > 0:
|
||
if qc > MAX:
|
||
logger.warning("Quarter shard %s: %d items (> %d)", qf.build_query(), qc, MAX)
|
||
shards.append(qf)
|
||
shard_items += qc
|
||
|
||
logger.info("Generated %d shards (≈ %d total items)", len(shards), shard_items)
|
||
return shards
|
||
|
||
# ------------------------------------------------------------------ checkpoint
|
||
|
||
def _load_checkpoint(self, redis_client: Any | None) -> dict[str, Any] | None:
|
||
"""Загрузить чекпоинт из Redis."""
|
||
if redis_client is None:
|
||
return None
|
||
try:
|
||
data = redis_client.get(self.CHECKPOINT_KEY)
|
||
if data:
|
||
cp = json.loads(data)
|
||
logger.info(
|
||
"Resuming from checkpoint: shard=%d, page=%d, synced=%d, collected=%d",
|
||
cp.get("shard_idx", 0), cp.get("page", 0),
|
||
cp.get("synced", 0), cp.get("items_collected", 0),
|
||
)
|
||
return cp
|
||
except Exception as exc:
|
||
logger.warning("Failed to load checkpoint: %s", exc)
|
||
return None
|
||
|
||
def _save_checkpoint(self, redis_client: Any | None, state: dict[str, Any]) -> None:
|
||
"""Сохранить чекпоинт в Redis (TTL = 6 часов)."""
|
||
if redis_client is None:
|
||
return
|
||
try:
|
||
redis_client.set(self.CHECKPOINT_KEY, json.dumps(state), ex=21600)
|
||
except Exception as exc:
|
||
logger.warning("Failed to save checkpoint: %s", exc)
|
||
|
||
def _clear_checkpoint(self, redis_client: Any | None) -> None:
|
||
"""Удалить чекпоинт после успешного завершения."""
|
||
if redis_client is None:
|
||
return
|
||
try:
|
||
redis_client.delete(self.CHECKPOINT_KEY)
|
||
except Exception:
|
||
pass
|
||
|
||
def sync_listing(
|
||
self,
|
||
limit: int | None = None,
|
||
filters: EncarFilters | None = None,
|
||
lane: str = "encar",
|
||
only_new: bool = False,
|
||
batch_size: int = 1000,
|
||
redis_client: Any | None = None,
|
||
allowed_brands: set[str] | None = None,
|
||
excluded_brands: set[str] | None = None,
|
||
runtime_filters: FiltersConfig | None = None,
|
||
probe_all_photos: bool = False,
|
||
) -> dict[str, Any]:
|
||
"""Полная синхронизация листинга Encar.
|
||
|
||
Encar API ограничивает пагинацию ~10k записей на один query.
|
||
Для получения всех 160k+ авто запрос автоматически разбивается
|
||
на шарды по car_type → year → quarter (через _build_sync_shards).
|
||
|
||
1. Генерирует шарды (фильтры), каждый ≤ 9500 записей.
|
||
2. Для каждого шарда собирает все авто, upsert в БД.
|
||
3. Помечает is_sold=True авто, которых нет в листинге.
|
||
|
||
Чекпоинт через Redis: сохраняет shard_idx + page для продолжения
|
||
после рестарта.
|
||
"""
|
||
self.persistence.create_tables()
|
||
filters = filters or EncarFilters()
|
||
page_size = self.DEFAULT_PAGE_SIZE
|
||
compiled_runtime_filters = self._compile_runtime_filters(runtime_filters)
|
||
|
||
# --- Генерация шардов ---
|
||
if limit:
|
||
# С limit не разбиваем на шарды
|
||
shards = [filters]
|
||
else:
|
||
shards = self._build_sync_shards(filters)
|
||
|
||
# --- Загрузка чекпоинта ---
|
||
checkpoint = self._load_checkpoint(redis_client)
|
||
start_shard_idx = 0
|
||
start_page = 0
|
||
total_available = 0
|
||
items_collected = 0
|
||
synced = 0
|
||
failed = 0
|
||
failed_ids: list[str] = []
|
||
all_origin_ids: set[str] = set()
|
||
|
||
if checkpoint and not limit:
|
||
start_shard_idx = checkpoint.get("shard_idx", 0)
|
||
start_page = checkpoint.get("page", 0)
|
||
items_collected = checkpoint.get("items_collected", 0)
|
||
synced = checkpoint.get("synced", 0)
|
||
failed = checkpoint.get("failed", 0)
|
||
total_available = checkpoint.get("total_available", 0)
|
||
logger.info(
|
||
"Checkpoint loaded: shard=%d/%d, page=%d, synced=%d",
|
||
start_shard_idx, len(shards), start_page, synced,
|
||
)
|
||
|
||
logger.info(
|
||
"Full sync Encar: %d shards, base query: %s",
|
||
len(shards), filters.build_query(),
|
||
)
|
||
|
||
# --- Фаза 1: Сбор и upsert по шардам ---
|
||
seen_vehicle_ids: set[str] = set() # Глобальный дедуп между шардами
|
||
|
||
for shard_idx, shard_filters in enumerate(shards):
|
||
if shard_idx < start_shard_idx:
|
||
continue
|
||
|
||
page = start_page if shard_idx == start_shard_idx else 0
|
||
start_page = 0 # сбрасываем после первого шарда
|
||
|
||
batch: list[Any] = []
|
||
consecutive_errors = 0
|
||
max_consecutive_errors = 5
|
||
consecutive_zero_new = 0
|
||
max_consecutive_zero_new = 5
|
||
|
||
shard_query = shard_filters.build_query()
|
||
shard_collected = 0
|
||
|
||
while True:
|
||
offset = page * page_size
|
||
try:
|
||
response = self._fetch_listing_page(
|
||
offset=offset, limit=page_size, filters=shard_filters,
|
||
)
|
||
except Exception as exc:
|
||
consecutive_errors += 1
|
||
logger.warning(
|
||
"Shard %d page %d fetch error (%d/%d): %s",
|
||
shard_idx, page, consecutive_errors, max_consecutive_errors, exc,
|
||
)
|
||
if consecutive_errors >= max_consecutive_errors:
|
||
logger.error("Too many errors in shard %d, moving to next", shard_idx)
|
||
break
|
||
time.sleep(5 * consecutive_errors)
|
||
continue
|
||
|
||
consecutive_errors = 0
|
||
results = response.get("SearchResults") or []
|
||
|
||
if page == 0:
|
||
shard_total = int(response.get("Count") or len(results))
|
||
if shard_idx == 0 or total_available == 0:
|
||
total_available = self._get_listing_count(filters)
|
||
if len(shards) > 1:
|
||
logger.info(
|
||
"Shard %d/%d (%s): %d items",
|
||
shard_idx + 1, len(shards), shard_query, shard_total,
|
||
)
|
||
|
||
if not results:
|
||
break
|
||
|
||
page += 1
|
||
|
||
new_in_page = 0
|
||
for item in results:
|
||
vehicle_id = str(item.get("Id") or "")
|
||
if not vehicle_id:
|
||
continue
|
||
if vehicle_id in seen_vehicle_ids:
|
||
continue
|
||
seen_vehicle_ids.add(vehicle_id)
|
||
batch.append(item)
|
||
items_collected += 1
|
||
shard_collected += 1
|
||
new_in_page += 1
|
||
|
||
if limit is not None and items_collected >= limit:
|
||
break
|
||
|
||
# API зацикливает после ~10k → выход из шарда
|
||
if new_in_page == 0:
|
||
consecutive_zero_new += 1
|
||
if consecutive_zero_new >= max_consecutive_zero_new:
|
||
logger.info(
|
||
"Shard %d pagination exhausted at page=%d, collected=%d",
|
||
shard_idx, page, shard_collected,
|
||
)
|
||
break
|
||
else:
|
||
consecutive_zero_new = 0
|
||
|
||
# Flush batch
|
||
if len(batch) >= batch_size:
|
||
s, f, fids = self._flush_batch(
|
||
batch, all_origin_ids, allowed_brands,
|
||
excluded_brands,
|
||
runtime_filters=compiled_runtime_filters,
|
||
probe_all_photos=probe_all_photos,
|
||
)
|
||
synced += s
|
||
failed += f
|
||
failed_ids.extend(fids)
|
||
batch = []
|
||
|
||
logger.info(
|
||
"Progress: shard=%d/%d, page=%d, collected=%d, synced=%d, failed=%d",
|
||
shard_idx + 1, len(shards), page, items_collected, synced, failed,
|
||
)
|
||
|
||
# Чекпоинт
|
||
if page % self.CHECKPOINT_SAVE_INTERVAL == 0:
|
||
self._save_checkpoint(redis_client, {
|
||
"shard_idx": shard_idx,
|
||
"page": page,
|
||
"items_collected": items_collected,
|
||
"synced": synced,
|
||
"failed": failed,
|
||
"total_available": total_available,
|
||
})
|
||
|
||
if limit is not None and items_collected >= limit:
|
||
break
|
||
|
||
if len(results) < page_size:
|
||
break
|
||
|
||
if page % 10 == 0:
|
||
time.sleep(0.5)
|
||
|
||
# Flush remaining batch после каждого шарда
|
||
if batch:
|
||
s, f, fids = self._flush_batch(
|
||
batch, all_origin_ids, allowed_brands,
|
||
excluded_brands,
|
||
runtime_filters=compiled_runtime_filters,
|
||
probe_all_photos=probe_all_photos,
|
||
)
|
||
synced += s
|
||
failed += f
|
||
failed_ids.extend(fids)
|
||
batch = []
|
||
|
||
# Чекпоинт после завершения шарда
|
||
self._save_checkpoint(redis_client, {
|
||
"shard_idx": shard_idx + 1,
|
||
"page": 0,
|
||
"items_collected": items_collected,
|
||
"synced": synced,
|
||
"failed": failed,
|
||
"total_available": total_available,
|
||
})
|
||
|
||
if limit is not None and items_collected >= limit:
|
||
break
|
||
|
||
# Пауза между шардами
|
||
if shard_idx < len(shards) - 1:
|
||
time.sleep(1)
|
||
|
||
if items_collected == 0:
|
||
logger.warning("No items collected from listing")
|
||
return {
|
||
"status": "empty",
|
||
"lane": lane,
|
||
"total_available": total_available,
|
||
"items_collected": 0,
|
||
"cars_synced": 0,
|
||
"cars_failed": 0,
|
||
"cars_marked_sold": 0,
|
||
}
|
||
|
||
# --- Фаза 2: Пометить sold ---
|
||
# Помечаем sold только если прошли ВСЕ шарды полностью (без resume из середины).
|
||
# Иначе all_origin_ids неполный и мы ошибочно пометим активные авто.
|
||
marked_sold = 0
|
||
if not limit and checkpoint is None:
|
||
marked_sold = self.persistence.mark_sold_not_in_listing(
|
||
active_origin_ids=all_origin_ids, lane=lane,
|
||
)
|
||
elif not limit:
|
||
logger.info(
|
||
"Skipping mark_sold: sync resumed from shard %d, origin_ids may be incomplete",
|
||
start_shard_idx,
|
||
)
|
||
|
||
self._clear_checkpoint(redis_client)
|
||
|
||
# Закрываем pool batch API
|
||
if self._batch_pool is not None:
|
||
self._batch_pool.close()
|
||
self._batch_pool = None
|
||
|
||
logger.info(
|
||
"Full sync complete: %d shards, %d collected, %d synced, %d failed, %d marked sold",
|
||
len(shards), items_collected, synced, failed, marked_sold,
|
||
)
|
||
|
||
return {
|
||
"status": "success",
|
||
"lane": lane,
|
||
"total_available": total_available,
|
||
"items_collected": items_collected,
|
||
"cars_synced": synced,
|
||
"cars_failed": failed,
|
||
"cars_marked_sold": marked_sold,
|
||
"shards_count": len(shards),
|
||
"failed_ids": failed_ids[:20],
|
||
}
|
||
|
||
def _flush_batch(
|
||
self,
|
||
items: list[Any],
|
||
all_origin_ids: set[str],
|
||
allowed_brands: set[str] | None = None,
|
||
excluded_brands: set[str] | None = None,
|
||
runtime_filters: RuntimeFilterSpec | None = None,
|
||
probe_all_photos: bool = False,
|
||
) -> tuple[int, int, list[str]]:
|
||
"""Маппит и upsert'ит батч items в БД. Возвращает (synced, failed, failed_ids)."""
|
||
records: list[CarRecord] = []
|
||
# Для пакетного пробинга: (record_index, item) пар, прошедших фильтр
|
||
probe_tasks: list[tuple[int, Any]] = []
|
||
failed = 0
|
||
failed_ids: list[str] = []
|
||
skipped_brands = 0
|
||
skipped_runtime = 0
|
||
canonicalized_records = 0
|
||
|
||
# Фаза 1: маппинг без пробинга + фильтр брендов
|
||
for item in items:
|
||
vehicle_id = str(item.get("Id") or "")
|
||
try:
|
||
vehicle_url = ENCAR_DETAIL_URL_TEMPLATE.format(vehicle_id=vehicle_id)
|
||
record = self.car_mapper.map_to_car_record(vehicle_url, item, probe_all_photos=False)
|
||
|
||
brand_lower = record.brand.lower() if record.brand else ""
|
||
if allowed_brands and brand_lower not in allowed_brands:
|
||
skipped_brands += 1
|
||
continue
|
||
if excluded_brands and brand_lower in excluded_brands:
|
||
skipped_brands += 1
|
||
continue
|
||
if not self._record_passes_runtime_filters(record, runtime_filters):
|
||
skipped_runtime += 1
|
||
continue
|
||
|
||
records.append(record)
|
||
probe_tasks.append((len(records) - 1, item))
|
||
except Exception as exc:
|
||
logger.warning("Failed to map vehicle %s: %s", vehicle_id, exc)
|
||
failed += 1
|
||
failed_ids.append(vehicle_id)
|
||
|
||
# Фаза 2: фото — берём из photo API (все фото) или только из листинга (4 фото)
|
||
if records and probe_all_photos:
|
||
t0 = time.monotonic()
|
||
self._batch_fetch_photo_lists(records, probe_tasks)
|
||
probe_time = time.monotonic() - t0
|
||
total_images_fetched = sum(len(r.images) for r in records)
|
||
logger.info(
|
||
"Photo API fetch: %d cars, %d total images in %.1fs",
|
||
len(records), total_images_fetched, probe_time,
|
||
)
|
||
elif records:
|
||
# Без пробинга: берём только фото из API (как reference-проект)
|
||
self._apply_api_photos(records, probe_tasks)
|
||
|
||
if records:
|
||
for record in records:
|
||
if self._apply_canonical_vehicle_id(record):
|
||
canonicalized_records += 1
|
||
|
||
deduped_records: list[CarRecord] = []
|
||
seen_origin_ids: set[str] = set()
|
||
collapsed_duplicates = 0
|
||
for record in records:
|
||
if record.origin_id in seen_origin_ids:
|
||
collapsed_duplicates += 1
|
||
continue
|
||
seen_origin_ids.add(record.origin_id)
|
||
deduped_records.append(record)
|
||
all_origin_ids.add(record.origin_id)
|
||
records = deduped_records
|
||
|
||
if canonicalized_records or collapsed_duplicates:
|
||
logger.info(
|
||
"Canonicalized %d Encar records by photo vehicle id and collapsed %d duplicates in batch",
|
||
canonicalized_records,
|
||
collapsed_duplicates,
|
||
)
|
||
|
||
# Фаза 3: upsert в БД
|
||
if records:
|
||
try:
|
||
result = self.persistence.upsert_cars_batch(records)
|
||
synced = result.get("inserted", 0) + result.get("updated", 0)
|
||
except Exception as exc:
|
||
logger.error("Batch upsert failed: %s", exc)
|
||
synced = 0
|
||
failed += len(records)
|
||
failed_ids.extend(r.origin_id for r in records)
|
||
else:
|
||
synced = 0
|
||
|
||
total_images = sum(len(r.images) for r in records)
|
||
logger.info(
|
||
"Batch stats: %d items → %d after runtime filter (-%d brands, -%d runtime), %d synced, %d total images",
|
||
len(items), len(records), skipped_brands, skipped_runtime, synced, total_images,
|
||
)
|
||
|
||
return synced, failed, failed_ids
|
||
|
||
def _apply_api_photos(
|
||
self,
|
||
records: list[CarRecord],
|
||
probe_tasks: list[tuple[int, Any]],
|
||
) -> None:
|
||
"""Применяет фото из API (Photo + Photos) без HEAD-пробинга."""
|
||
for rec_idx, item in probe_tasks:
|
||
photos = item.get("Photos") or []
|
||
photo_prefix = self.car_mapper._as_str(item.get("Photo") or "")
|
||
|
||
images: list[ImageRecord] = []
|
||
for raw in (photos if isinstance(photos, list) else []):
|
||
if not isinstance(raw, dict):
|
||
continue
|
||
location = raw.get("location") or raw.get("Location") or ""
|
||
if not location:
|
||
continue
|
||
url = self.car_mapper._build_image_url(str(location))
|
||
order = int(raw.get("ordering") or 0)
|
||
images.append(ImageRecord(fullres_image=url, preview_image=url, order_index=order))
|
||
|
||
# Если Photos пусто, но Photo есть — создаём хотя бы одну запись
|
||
if not images and photo_prefix:
|
||
url = self.car_mapper._build_image_url(photo_prefix.rstrip("_") + "_001.jpg")
|
||
images.append(ImageRecord(fullres_image=url, preview_image=url, order_index=1))
|
||
|
||
images.sort(key=lambda x: x.order_index)
|
||
records[rec_idx].images = images
|
||
|
||
def _batch_fetch_photo_lists(
|
||
self,
|
||
records: list[CarRecord],
|
||
probe_tasks: list[tuple[int, Any]],
|
||
) -> None:
|
||
"""Получает полные списки фото через batch API (/v1/readside/vehicles).
|
||
|
||
Batch endpoint найден в JS-бандле мобильной SPA (fem.encar.com).
|
||
Принимает до 20 vehicleIds за запрос, возвращает все фото каждой машины.
|
||
В ~4.4x быстрее legacy photo API (370 vs 61 машин/с) и без connection errors.
|
||
"""
|
||
tasks: list[tuple[int, str]] = []
|
||
for rec_idx, item in probe_tasks:
|
||
vehicle_id = str(item.get("Id") or "")
|
||
if vehicle_id:
|
||
tasks.append((rec_idx, vehicle_id))
|
||
|
||
if not tasks:
|
||
return
|
||
|
||
pool = self._get_batch_pool()
|
||
|
||
# Разбиваем на чанки по BATCH_VEHICLES_CHUNK_SIZE (20) ID
|
||
chunks: list[list[tuple[int, str]]] = []
|
||
for i in range(0, len(tasks), BATCH_VEHICLES_CHUNK_SIZE):
|
||
chunks.append(tasks[i:i + BATCH_VEHICLES_CHUNK_SIZE])
|
||
|
||
def _fetch_chunk(chunk: list[tuple[int, str]]) -> list[tuple[int, list[ImageRecord]]]:
|
||
"""Запрашивает фото для чанка машин одним batch-запросом."""
|
||
ids_str = ",".join(vid for _, vid in chunk)
|
||
path = f"/v1/readside/vehicles?vehicleIds={ids_str}&include=PHOTOS"
|
||
for attempt in range(2):
|
||
try:
|
||
resp = pool.request("GET", path, timeout=15, retries=False)
|
||
if resp.status == 200:
|
||
if not resp.data:
|
||
return [(ri, []) for ri, _ in chunk]
|
||
vehicles = json.loads(resp.data.decode("utf-8", errors="ignore"))
|
||
# Маппим по vehicleId, а не по индексу — API может пропускать удалённые ID
|
||
id_to_rec: dict[str, int] = {vid: ri for ri, vid in chunk}
|
||
result: list[tuple[int, list[ImageRecord]]] = []
|
||
for vehicle_data in vehicles:
|
||
vid = str(vehicle_data.get("vehicleId") or vehicle_data.get("id") or "")
|
||
rec_idx = id_to_rec.get(vid)
|
||
if rec_idx is None:
|
||
continue
|
||
photos_raw = vehicle_data.get("photos") or []
|
||
images: list[ImageRecord] = []
|
||
seen_paths: set[str] = set()
|
||
for photo in photos_raw:
|
||
p = photo.get("path") or ""
|
||
if not p or p in seen_paths:
|
||
continue
|
||
seen_paths.add(p)
|
||
full_url = self.car_mapper._build_image_url(p)
|
||
code = photo.get("code") or "0"
|
||
order = int(code) if code.isdigit() else 0
|
||
images.append(ImageRecord(
|
||
fullres_image=full_url,
|
||
preview_image=full_url,
|
||
order_index=order,
|
||
))
|
||
images.sort(key=lambda x: x.order_index)
|
||
result.append((rec_idx, images))
|
||
return result
|
||
if resp.status in (429, 503):
|
||
time.sleep(1 + attempt * 2)
|
||
continue
|
||
return [(ri, []) for ri, _ in chunk]
|
||
except json.JSONDecodeError:
|
||
return [(ri, []) for ri, _ in chunk]
|
||
except Exception as exc:
|
||
logger.debug("Batch photo chunk error (attempt %d): %s", attempt + 1, exc)
|
||
if attempt < 1:
|
||
time.sleep(1)
|
||
return [(ri, []) for ri, _ in chunk]
|
||
|
||
results: dict[int, list[ImageRecord]] = {}
|
||
with ThreadPoolExecutor(max_workers=30) as executor:
|
||
futures = {executor.submit(_fetch_chunk, ch): ch for ch in chunks}
|
||
for future in as_completed(futures):
|
||
try:
|
||
for rec_idx, images in future.result():
|
||
if images:
|
||
results[rec_idx] = images
|
||
except Exception as exc:
|
||
logger.warning("Batch photo future failed: %s", exc)
|
||
|
||
# Применяем: batch API → fallback на Photos из листинга
|
||
api_ok = 0
|
||
fallback_count = 0
|
||
for rec_idx, item in probe_tasks:
|
||
if rec_idx in results:
|
||
records[rec_idx].images = results[rec_idx]
|
||
api_ok += 1
|
||
else:
|
||
fallback_count += 1
|
||
photos = item.get("Photos") or []
|
||
photo_prefix = self.car_mapper._as_str(item.get("Photo") or "")
|
||
images: list[ImageRecord] = []
|
||
for raw in (photos if isinstance(photos, list) else []):
|
||
if not isinstance(raw, dict):
|
||
continue
|
||
location = raw.get("location") or raw.get("Location") or ""
|
||
if not location:
|
||
continue
|
||
url = self.car_mapper._build_image_url(str(location))
|
||
order = int(raw.get("ordering") or 0)
|
||
images.append(ImageRecord(fullres_image=url, preview_image=url, order_index=order))
|
||
if not images and photo_prefix:
|
||
url = self.car_mapper._build_image_url(photo_prefix.rstrip("_") + "_001.jpg")
|
||
images.append(ImageRecord(fullres_image=url, preview_image=url, order_index=1))
|
||
images.sort(key=lambda x: x.order_index)
|
||
records[rec_idx].images = images
|
||
if fallback_count:
|
||
logger.info("Batch photo API: %d OK, %d fallback to listing photos", api_ok, fallback_count)
|
||
|
||
def _fetch_listing_page(self, offset: int, limit: int, filters: EncarFilters | None = None) -> dict[str, Any]:
|
||
filters = filters or EncarFilters()
|
||
query = filters.build_query()
|
||
params = {
|
||
"count": "true",
|
||
"q": query,
|
||
# CreatedDate — стабильная сортировка (дата создания объявления не меняется).
|
||
# ModifiedDate нестабильна: машины "всплывают" при обновлении цены/фото,
|
||
# вызывая до 50% дублей на поздних страницах длинного синка.
|
||
"sr": f"|CreatedDate|{offset}|{limit}",
|
||
}
|
||
url = f"{ENCAR_LISTING_API}?{urlencode(params)}"
|
||
return self._fetch_json(url)
|
||
|
||
def _fetch_listing_item(self, vehicle_id: str) -> dict[str, Any] | None:
|
||
"""Ищет авто в листинге по ID (макс 5 страниц)."""
|
||
page = 0
|
||
page_size = 100
|
||
max_pages = 5 # Ограничиваем поиск
|
||
|
||
while page < max_pages:
|
||
response = self._fetch_listing_page(offset=page * page_size, limit=page_size)
|
||
results = response.get("SearchResults") or []
|
||
if not results:
|
||
break
|
||
for item in results:
|
||
if str(item.get("Id")) == vehicle_id:
|
||
return item
|
||
if len(results) < page_size:
|
||
break
|
||
page += 1
|
||
return None
|
||
|
||
def _fetch_vehicle_detail_api(self, vehicle_id: str) -> dict[str, Any] | None:
|
||
url = ENCAR_DETAIL_API_TEMPLATE.format(vehicle_id=vehicle_id)
|
||
params = {
|
||
"include": ",".join(["ADVERTISEMENT", "VEHICLE", "SPECS", "IMAGE", "CONDITION"])
|
||
}
|
||
full_url = f"{url}?{urlencode(params)}"
|
||
payload = self._fetch_json(full_url, headers={"Referer": ENCAR_DETAIL_URL_TEMPLATE.format(vehicle_id=vehicle_id)})
|
||
return payload if payload and payload.get("vehicleId") else None
|
||
|
||
def _fetch_json(self, url: str, headers: dict[str, str] | None = None) -> dict[str, Any]:
|
||
request_headers = {
|
||
"User-Agent": self.settings.fingerprint.user_agent,
|
||
"Accept": "application/json, text/plain, */*",
|
||
"Accept-Language": "ko-KR,ko;q=0.9,en;q=0.8",
|
||
"Referer": "https://www.encar.com/",
|
||
}
|
||
if headers:
|
||
request_headers.update(headers)
|
||
req = Request(url, headers=request_headers)
|
||
for attempt in range(3):
|
||
try:
|
||
with urlopen(req, timeout=30) as resp:
|
||
content = resp.read().decode("utf-8", errors="ignore")
|
||
return json.loads(content)
|
||
except (HTTPError, URLError, ConnectionError, TimeoutError, ValueError, OSError) as exc:
|
||
logger.warning("Encar JSON fetch failed (%s) for %s; attempt %d", exc, url, attempt + 1)
|
||
if attempt == 2:
|
||
raise
|
||
time.sleep(2 + attempt * 2)
|
||
return {}
|
||
|
||
def _extract_vehicle_id(self, vehicle_url: str) -> str | None:
|
||
match = VEHICLE_ID_RE.search(vehicle_url)
|
||
if match:
|
||
return match.group(1)
|
||
return None
|