Files
dubizzle/scripts/export_for_upload.py
2026-07-01 14:55:15 +03:00

353 lines
13 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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:
"""Исходный формат имени папки.
Формат:
<source>_<source_id>__<brand>_<model>_<year>
"""
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()