From a4d0f35279ef8a873c3a6467f4fbc0eb76827ee5 Mon Sep 17 00:00:00 2001 From: qananasikq Date: Wed, 1 Jul 2026 13:49:30 +0300 Subject: [PATCH] fix car mapping --- .gitignore | 7 +- dubizzle_scraper/core/config.py | 2 +- dubizzle_scraper/discovery/algolia.py | 29 +- scripts/export_for_upload.py | 368 ++++++++++++++++++++++++++ 4 files changed, 399 insertions(+), 7 deletions(-) create mode 100644 scripts/export_for_upload.py diff --git a/.gitignore b/.gitignore index 1b10c3c..e0b0433 100644 --- a/.gitignore +++ b/.gitignore @@ -20,4 +20,9 @@ build/ artifacts/ celerybeat-schedule* tokens_data/ -tokens.json \ No newline at end of file +tokens.json + + +tesst/ +export_output/ +upload_export/ \ No newline at end of file diff --git a/dubizzle_scraper/core/config.py b/dubizzle_scraper/core/config.py index dac1c6b..d3c4c9d 100644 --- a/dubizzle_scraper/core/config.py +++ b/dubizzle_scraper/core/config.py @@ -222,7 +222,7 @@ class AlgoliaConfig: index_name: str = _env_str("DUBIZZLE_ALGOLIA_INDEX_NAME", "motors.com") category_slug: str = _env_str("DUBIZZLE_ALGOLIA_CATEGORY_SLUG", "motors/used-cars") base_url: str = _env_str("DUBIZZLE_ALGOLIA_BASE_URL", "https://WD0PTZ13ZS-dsn.algolia.net") - hits_per_page: int = _env_int("DUBIZZLE_ALGOLIA_HITS_PER_PAGE", 100) + hits_per_page: int = _env_int("DUBIZZLE_ALGOLIA_HITS_PER_PAGE", 20) # --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) --- diff --git a/dubizzle_scraper/discovery/algolia.py b/dubizzle_scraper/discovery/algolia.py index 1ed80f2..2a3384d 100644 --- a/dubizzle_scraper/discovery/algolia.py +++ b/dubizzle_scraper/discovery/algolia.py @@ -99,11 +99,22 @@ class _AlgoliaHttpClient: }, method="POST", ) - with urlopen(request, timeout=30) as response: - body = response.read().decode("utf-8", errors="ignore") - parsed = json.loads(body) - results = parsed.get("results") or [] - return results[0] if results and isinstance(results[0], dict) else {} + + last_exc: Exception | None = None + for attempt in range(1, 4): + try: + with urlopen(request, timeout=15) as response: + body = response.read().decode("utf-8", errors="ignore") + parsed = json.loads(body) + results = parsed.get("results") or [] + return results[0] if results and isinstance(results[0], dict) else {} + except Exception as exc: + last_exc = exc + if attempt >= 3: + break + time.sleep(0.8 * attempt) + + raise last_exc or RuntimeError("Algolia request failed") def _build_origin_id_from_hit(hit: dict[str, Any]) -> str | None: @@ -282,6 +293,7 @@ def _fetch_shard_hits( year_min: int | None, year_max: int | None, known_origin_ids: set[str] | None, + limit: int | None, threshold: float, max_duration_seconds: float | None, started_at: float, @@ -347,6 +359,12 @@ def _fetch_shard_hits( hit_records[url] = raw_hit if origin_id: origin_ids_by_url[url] = origin_id + if limit is not None and limit > 0 and len(urls) >= limit: + early_stopped = True + break + + if limit is not None and limit > 0 and len(urls) >= limit: + break if known_origin_ids is not None and threshold > 0 and (page_known + page_new) > 0: ratio = page_known / (page_known + page_new) @@ -438,6 +456,7 @@ def discover_vehicle_hits_from_algolia( year_min=year_min, year_max=year_max, known_origin_ids=known_origin_ids, + limit=(limit - len(collected_urls)) if limit is not None and limit > 0 else None, threshold=threshold, max_duration_seconds=max_duration_seconds, started_at=started_at, diff --git a/scripts/export_for_upload.py b/scripts/export_for_upload.py new file mode 100644 index 0000000..fb6b176 --- /dev/null +++ b/scripts/export_for_upload.py @@ -0,0 +1,368 @@ +"""Экспорт авто из Dubizzle в папки для загрузки на сайт. + +Делает полный путь данных без БД/Docker: + Algolia discovery -> CarMapper (маппинг) -> папка на каждое авто (car.json + фото). + +Каждое авто кладётся в отдельную папку: + /______/ + ├── car.json # замапленная запись (CarRecord) + сырой Algolia hit + └── images/ # скачанные полноразмерные фото (01.jpg, 02.jpg, ...) + +Скрипт идемпотентный и возобновляемый: повторный запуск пропускает уже +выгруженные авто и уже скачанные фото. + +Запуск: + python scripts/export_for_upload.py --target 11000 --output tesst +""" + +from __future__ import annotations + +import argparse +import json +import re +import sys +import time +from concurrent.futures import ThreadPoolExecutor, as_completed +from pathlib import Path +from urllib.request import Request, urlopen + +# Делаем пакет импортируемым при запуске из корня репозитория. +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) + +from dubizzle_scraper.core.config import settings # noqa: E402 +from dubizzle_scraper.core.config import _AUTO_YEAR_SPLITS # noqa: E402 +from dubizzle_scraper.discovery.algolia import ( # noqa: E402 + AlgoliaDiscoveryError, + discover_vehicle_hits_from_algolia, +) +from dubizzle_scraper.parsing.mapper import CarMapper # noqa: E402 + + +_SAFE_RE = re.compile(r"[^A-Za-z0-9._-]+") + + +def _safe(text: str, max_len: int = 40) -> str: + text = (text or "").strip().replace(" ", "-") + text = _SAFE_RE.sub("", text) + return text[:max_len] or "NA" + + +def _folder_name(record) -> str: + """Исходный формат имени папки. + + Формат: + _____ + """ + origin = str(getattr(record, "origin_id", "") or "") + source = str(getattr(record, "origin", "") or "dubizzle").split(":", 1)[0].strip().lower() or "dubizzle" + source_id = origin.split(":", 1)[-1].strip() if ":" in origin else origin.strip() + source_id = source_id or origin or "NA" + return ( + f"{_safe(source, 20)}_{_safe(source_id, 40)}__" + f"{_safe(record.brand)}_{_safe(record.model)}_{record.year or 'NA'}" + ) + + +def _scan_existing_origin_ids(output_dir: Path) -> set[str]: + """Собирает raw origin_id уже выгруженных авто из car.json.""" + existing: set[str] = set() + if not output_dir.exists(): + return existing + for child in output_dir.iterdir(): + if not child.is_dir(): + continue + car_json = child / "car.json" + if not car_json.exists(): + continue + try: + payload = json.loads(car_json.read_text(encoding="utf-8")) + except Exception: + continue + origin_id = (((payload or {}).get("car") or {}).get("origin_id") or "").strip() + if origin_id: + existing.add(origin_id) + return existing + + +def _build_payload_insights(hit: dict) -> dict: + return { + "vehicle_core": { + "lot_number": hit.get("id") or hit.get("objectID"), + "year": hit.get("year"), + "make": hit.get("make"), + "model": hit.get("model"), + "trim": hit.get("trim") or hit.get("motors_trim"), + "odometer": hit.get("kilometers") or hit.get("odometer"), + "body_type": hit.get("body_type"), + "gearbox": hit.get("transmission_type") or hit.get("transmission"), + "drive": hit.get("drive_type"), + "fuel_type": hit.get("fuel_type"), + "color": hit.get("exterior_color") or hit.get("color"), + "seller": hit.get("seller_type"), + "location": hit.get("location_name") or hit.get("location"), + "title": hit.get("title") or hit.get("name"), + "country": "AE", + "selling_type": "STOCK", + }, + "pricing": { + "buy_now": hit.get("price"), + "current_bid": None, + "actual_cash_value": None, + "currency": hit.get("price_currency") or "AED", + }, + "damage": {}, + "auction": {"sale_status": hit.get("status")}, + "images": { + "urls": hit.get("photo_mains") or hit.get("photo_thumbnails") or [], + }, + } + + +def discover_cars( + target: int, + max_seconds: float | None, + skip_ids: set[str] | None = None, +) -> dict[str, dict]: + """Собирает hit-записи по годовым сегментам, дедупит по origin_id. + + ``skip_ids`` — raw origin_id уже выгруженных авто, которые не нужно + считать как новые. Цель ``target`` считается именно по новым записям. + """ + from dubizzle_scraper.discovery.algolia import _build_origin_id_from_hit + + skip_ids = skip_ids or set() + collected: dict[str, dict] = {} # origin_id -> hit + started = time.perf_counter() + + for yr_min, yr_max in _AUTO_YEAR_SPLITS: + if len(collected) >= target: + break + + print(f" [start] segment {yr_min}-{yr_max}", flush=True) + attempts = 0 + segment_hits: dict[str, dict] = {} + while attempts < 5 and not segment_hits: + attempts += 1 + remaining_limit = max(1, target - len(collected)) + remaining_time = None + if max_seconds is not None: + remaining_time = max_seconds - (time.perf_counter() - started) + if remaining_time <= 0: + break + try: + result = discover_vehicle_hits_from_algolia( + settings=settings, + year_min=yr_min, + year_max=yr_max, + limit=remaining_limit, + known_origin_ids=skip_ids, + max_duration_seconds=remaining_time, + ) + segment_hits = result.hit_records + except Exception as exc: + print(f" [warn] segment {yr_min}-{yr_max} attempt {attempts} failed: {exc}", flush=True) + if attempts < 5: + time.sleep(min(15, 2 * attempts)) + + # Fallback: при сетевых обрывах Algolia уменьшаем page size, + # это заметно стабильнее на больших сегментах. + try: + original_hpp = settings.algolia.hits_per_page + settings.algolia.hits_per_page = min(20, original_hpp) + result = discover_vehicle_hits_from_algolia( + settings=settings, + year_min=yr_min, + year_max=yr_max, + limit=remaining_limit, + known_origin_ids=skip_ids, + max_duration_seconds=remaining_time, + ) + segment_hits = result.hit_records + except Exception as exc2: + print(f" [warn] segment {yr_min}-{yr_max} fallback failed: {exc2}", flush=True) + finally: + settings.algolia.hits_per_page = original_hpp + + if not segment_hits: + continue + + added = 0 + for url, hit in segment_hits.items(): + origin_id = _build_origin_id_from_hit(hit) or url + if origin_id in collected or origin_id in skip_ids: + continue + hit.setdefault("_vehicle_url", url) + collected[origin_id] = hit + added += 1 + if len(collected) >= target: + break + + print( + f" segment {yr_min}-{yr_max}: +{added} new (total {len(collected)}/{target})", + flush=True, + ) + + return collected + + +def _download_image(url: str, dest: Path, user_agent: str) -> bool: + if dest.exists() and dest.stat().st_size > 0: + return True + try: + req = Request(url, headers={"user-agent": user_agent, "accept": "image/*"}) + with urlopen(req, timeout=30) as resp: + data = resp.read() + if not data: + return False + dest.write_bytes(data) + return True + except Exception: + return False + + +def export_car( + index: int, + hit: dict, + output_dir: Path, + mapper: CarMapper, + user_agent: str, + image_workers: int, +) -> tuple[int, int]: + """Возвращает (скачано_фото, всего_фото).""" + url = hit.get("_vehicle_url") or "" + record = mapper.map_to_car_record( + vehicle_url=url, + vehicle_summary=hit, + payload_insights=_build_payload_insights(hit), + ) + data = record.model_dump(mode="json") + + car_dir = output_dir / _folder_name(record) + images_dir = car_dir / "images" + images_dir.mkdir(parents=True, exist_ok=True) + + # car.json: только нужные замапленные данные для загрузки. + # Сырой Algolia payload не сохраняем — он сильно раздувает JSON и в БД не нужен. + payload = { + "car": data, + "source_url": url, + } + (car_dir / "car.json").write_text( + json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8" + ) + + images = data.get("images") or [] + total = len(images) + if total == 0: + return 0, 0 + + tasks = [] + for img in images: + order = int(img.get("order_index", 0)) + src = img.get("fullres_image") + if not src: + continue + dest = images_dir / f"{order + 1:02d}.jpg" + tasks.append((src, dest)) + + downloaded = 0 + with ThreadPoolExecutor(max_workers=max(1, image_workers)) as pool: + futures = {pool.submit(_download_image, src, dest, user_agent): dest for src, dest in tasks} + for fut in as_completed(futures): + if fut.result(): + downloaded += 1 + + return downloaded, total + + +def main() -> None: + parser = argparse.ArgumentParser(description="Экспорт авто Dubizzle в папки для загрузки") + parser.add_argument("--target", type=int, default=11000, help="Сколько авто выгрузить") + parser.add_argument("--output", default="tesst", help="Папка вывода") + parser.add_argument("--car-workers", type=int, default=6, help="Параллельных авто") + parser.add_argument("--image-workers", type=int, default=8, help="Параллельных фото на авто") + parser.add_argument("--discovery-timeout", type=float, default=None, help="Лимит времени discovery, сек") + parser.add_argument("--no-images", action="store_true", help="Только car.json без скачивания фото") + args = parser.parse_args() + + output_dir = Path(args.output) + output_dir.mkdir(parents=True, exist_ok=True) + user_agent = settings.fingerprint.user_agent + mapper = CarMapper() + + existing_ids = _scan_existing_origin_ids(output_dir) + if existing_ids: + print( + f" Найдено {len(existing_ids)} уже выгруженных авто — будут пропущены.", + flush=True, + ) + + print(f"[1/2] Discovery: собираю до {args.target} новых авто через Algolia ...", flush=True) + cars = discover_cars(args.target, args.discovery_timeout, skip_ids=existing_ids) + hits = list(cars.values()) + print(f" Найдено {len(hits)} новых уникальных авто.", flush=True) + + if not hits: + print("Нет данных для экспорта. Проверь сеть/Algolia.", flush=True) + return + + print(f"[2/2] Экспорт в '{output_dir}/' ...", flush=True) + total_cars = len(hits) + total_imgs = 0 + done_cars = 0 + image_workers = 0 if args.no_images else args.image_workers + + def _work(idx_hit): + idx, hit = idx_hit + if args.no_images: + return _export_no_images(idx, hit, output_dir, mapper) + return export_car(idx, hit, output_dir, mapper, user_agent, image_workers) + + with ThreadPoolExecutor(max_workers=max(1, args.car_workers)) as pool: + futures = { + pool.submit(_work, (idx, hit)): idx + for idx, hit in enumerate(hits, start=1) + } + for fut in as_completed(futures): + idx = futures[fut] + try: + dl, tot = fut.result() + total_imgs += dl + except Exception as exc: + print(f" [err] car #{idx}: {exc}", flush=True) + continue + done_cars += 1 + if done_cars % 100 == 0 or done_cars == total_cars: + print( + f" {done_cars}/{total_cars} авто, фото скачано: {total_imgs}", + flush=True, + ) + + print( + f"Готово: {done_cars} авто в '{output_dir}/', всего фото скачано: {total_imgs}.", + flush=True, + ) + + +def _export_no_images(index: int, hit: dict, output_dir: Path, mapper: CarMapper) -> tuple[int, int]: + url = hit.get("_vehicle_url") or "" + record = mapper.map_to_car_record( + vehicle_url=url, + vehicle_summary=hit, + payload_insights=_build_payload_insights(hit), + ) + data = record.model_dump(mode="json") + car_dir = output_dir / _folder_name(record) + car_dir.mkdir(parents=True, exist_ok=True) + payload = { + "car": data, + "source_url": url, + } + (car_dir / "car.json").write_text( + json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8" + ) + return 0, len(data.get("images") or []) + + +if __name__ == "__main__": + main()