fix car mapping
This commit is contained in:
7
.gitignore
vendored
7
.gitignore
vendored
@@ -20,4 +20,9 @@ build/
|
|||||||
artifacts/
|
artifacts/
|
||||||
celerybeat-schedule*
|
celerybeat-schedule*
|
||||||
tokens_data/
|
tokens_data/
|
||||||
tokens.json
|
tokens.json
|
||||||
|
|
||||||
|
|
||||||
|
tesst/
|
||||||
|
export_output/
|
||||||
|
upload_export/
|
||||||
@@ -222,7 +222,7 @@ class AlgoliaConfig:
|
|||||||
index_name: str = _env_str("DUBIZZLE_ALGOLIA_INDEX_NAME", "motors.com")
|
index_name: str = _env_str("DUBIZZLE_ALGOLIA_INDEX_NAME", "motors.com")
|
||||||
category_slug: str = _env_str("DUBIZZLE_ALGOLIA_CATEGORY_SLUG", "motors/used-cars")
|
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")
|
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-расписание) ---
|
# --- Конфиг Celery (лимиты задач, concurrency, beat-расписание) ---
|
||||||
|
|||||||
@@ -99,11 +99,22 @@ class _AlgoliaHttpClient:
|
|||||||
},
|
},
|
||||||
method="POST",
|
method="POST",
|
||||||
)
|
)
|
||||||
with urlopen(request, timeout=30) as response:
|
|
||||||
body = response.read().decode("utf-8", errors="ignore")
|
last_exc: Exception | None = None
|
||||||
parsed = json.loads(body)
|
for attempt in range(1, 4):
|
||||||
results = parsed.get("results") or []
|
try:
|
||||||
return results[0] if results and isinstance(results[0], dict) else {}
|
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:
|
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_min: int | None,
|
||||||
year_max: int | None,
|
year_max: int | None,
|
||||||
known_origin_ids: set[str] | None,
|
known_origin_ids: set[str] | None,
|
||||||
|
limit: int | None,
|
||||||
threshold: float,
|
threshold: float,
|
||||||
max_duration_seconds: float | None,
|
max_duration_seconds: float | None,
|
||||||
started_at: float,
|
started_at: float,
|
||||||
@@ -347,6 +359,12 @@ def _fetch_shard_hits(
|
|||||||
hit_records[url] = raw_hit
|
hit_records[url] = raw_hit
|
||||||
if origin_id:
|
if origin_id:
|
||||||
origin_ids_by_url[url] = 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:
|
if known_origin_ids is not None and threshold > 0 and (page_known + page_new) > 0:
|
||||||
ratio = page_known / (page_known + page_new)
|
ratio = page_known / (page_known + page_new)
|
||||||
@@ -438,6 +456,7 @@ def discover_vehicle_hits_from_algolia(
|
|||||||
year_min=year_min,
|
year_min=year_min,
|
||||||
year_max=year_max,
|
year_max=year_max,
|
||||||
known_origin_ids=known_origin_ids,
|
known_origin_ids=known_origin_ids,
|
||||||
|
limit=(limit - len(collected_urls)) if limit is not None and limit > 0 else None,
|
||||||
threshold=threshold,
|
threshold=threshold,
|
||||||
max_duration_seconds=max_duration_seconds,
|
max_duration_seconds=max_duration_seconds,
|
||||||
started_at=started_at,
|
started_at=started_at,
|
||||||
|
|||||||
368
scripts/export_for_upload.py
Normal file
368
scripts/export_for_upload.py
Normal file
@@ -0,0 +1,368 @@
|
|||||||
|
"""Экспорт авто из Dubizzle в папки для загрузки на сайт.
|
||||||
|
|
||||||
|
Делает полный путь данных без БД/Docker:
|
||||||
|
Algolia discovery -> CarMapper (маппинг) -> папка на каждое авто (car.json + фото).
|
||||||
|
|
||||||
|
Каждое авто кладётся в отдельную папку:
|
||||||
|
<output>/<NNNNN>__<brand>_<model>_<year>__<origin_id>/
|
||||||
|
├── 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:
|
||||||
|
"""Исходный формат имени папки.
|
||||||
|
|
||||||
|
Формат:
|
||||||
|
<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()
|
||||||
Reference in New Issue
Block a user