Compare commits
5 Commits
42bd4c64d6
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1b5659f663 | ||
|
|
a4d0f35279 | ||
|
|
6a46b2f022 | ||
|
|
66bb930f34 | ||
|
|
046e9f125f |
5
.gitignore
vendored
5
.gitignore
vendored
@@ -21,3 +21,8 @@ artifacts/
|
||||
celerybeat-schedule*
|
||||
tokens_data/
|
||||
tokens.json
|
||||
|
||||
|
||||
tesst/
|
||||
export_output/
|
||||
upload_export/
|
||||
@@ -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-расписание) ---
|
||||
|
||||
@@ -99,11 +99,22 @@ class _AlgoliaHttpClient:
|
||||
},
|
||||
method="POST",
|
||||
)
|
||||
with urlopen(request, timeout=30) as response:
|
||||
|
||||
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,
|
||||
|
||||
@@ -8,6 +8,11 @@ from ..core.utils import LOT_RE, PRICE_RE, deep_find_all_keys, deep_find_key, fi
|
||||
|
||||
logger = logging.getLogger("dubizzle_scraper.parsers")
|
||||
|
||||
NEXT_DATA_RE = re.compile(
|
||||
r'<script[^>]+id=["\']__NEXT_DATA__["\'][^>]*type=["\']application/json["\'][^>]*>(.*?)</script>',
|
||||
re.IGNORECASE | re.DOTALL,
|
||||
)
|
||||
|
||||
|
||||
class VehicleParser:
|
||||
# Парсер страницы авто.
|
||||
@@ -157,6 +162,11 @@ class VehicleParser:
|
||||
network_dump = network_dump or {}
|
||||
responses = network_dump.get("json_responses", [])
|
||||
payloads = [item.get("payload") for item in responses if isinstance(item.get("payload"), (dict, list))]
|
||||
embedded = self._extract_embedded_json(page_html)
|
||||
for item in embedded:
|
||||
payload = item.get("payload")
|
||||
if isinstance(payload, (dict, list)):
|
||||
payloads.append(payload)
|
||||
dom_kv = self._parse_dom_key_value_pairs(dom_text)
|
||||
page_title = ""
|
||||
title_match = re.search(r"<title[^>]*>(.*?)</title>", page_html or "", re.IGNORECASE | re.DOTALL)
|
||||
@@ -194,11 +204,9 @@ class VehicleParser:
|
||||
if not summary.get("current_bid") and len(prices) > 1:
|
||||
summary["current_bid"] = prices[1]
|
||||
|
||||
embedded = self._extract_embedded_json(page_html)
|
||||
for item in embedded:
|
||||
p = item.get("payload")
|
||||
if isinstance(p, (dict, list)):
|
||||
payloads.append(p)
|
||||
# Доп. проход по JSON.
|
||||
extra = deep_find_all_keys([p], self.SUMMARY_KEY_MAP)
|
||||
for field, vals in extra.items():
|
||||
@@ -299,6 +307,11 @@ class VehicleParser:
|
||||
|
||||
@staticmethod
|
||||
def _guess_currency(summary: dict[str, Any]) -> str:
|
||||
for key in ["currency", "price", "buy_now", "current_bid", "actual_cash_value", "estimated_repair_cost"]:
|
||||
value = str(summary.get(key) or "")
|
||||
upper = value.upper()
|
||||
if "AED" in upper or "د.إ" in value:
|
||||
return "AED"
|
||||
for key in ["buy_now", "current_bid", "actual_cash_value", "estimated_repair_cost"]:
|
||||
value = str(summary.get(key) or "")
|
||||
if "$" in value:
|
||||
@@ -322,9 +335,18 @@ class VehicleParser:
|
||||
def _extract_embedded_json(html: str) -> list[dict[str, Any]]:
|
||||
scripts = re.findall(r"<script[^>]*>(.*?)</script>", html or "", flags=re.DOTALL | re.IGNORECASE)
|
||||
extracted: list[dict[str, Any]] = []
|
||||
next_match = NEXT_DATA_RE.search(html or "")
|
||||
if next_match:
|
||||
try:
|
||||
next_data = json.loads(html_module.unescape(next_match.group(1).strip()))
|
||||
extracted.append({"type": "next_data", "payload": next_data})
|
||||
except Exception:
|
||||
pass
|
||||
for script_text in scripts:
|
||||
if "{" not in script_text and "[" not in script_text:
|
||||
continue
|
||||
if "__NEXT_DATA__" in script_text:
|
||||
continue
|
||||
# Пропускаем большие блоки.
|
||||
if len(script_text) > 51_200:
|
||||
continue
|
||||
|
||||
@@ -7,6 +7,7 @@ import signal
|
||||
import time
|
||||
import uuid
|
||||
import html as html_module
|
||||
from collections.abc import Mapping
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from concurrent.futures import TimeoutError as FuturesTimeoutError
|
||||
from datetime import datetime, timezone
|
||||
@@ -44,6 +45,10 @@ VEHICLE_ID_RE = re.compile(
|
||||
)
|
||||
HTML_TAG_RE = re.compile(r"<[^>]+>")
|
||||
SCRIPT_STYLE_RE = re.compile(r"<(script|style)[^>]*>.*?</\1>", re.IGNORECASE | re.DOTALL)
|
||||
NEXT_DATA_RE = re.compile(
|
||||
r'<script[^>]+id=["\']__NEXT_DATA__["\'][^>]*type=["\']application/json["\'][^>]*>(.*?)</script>',
|
||||
re.IGNORECASE | re.DOTALL,
|
||||
)
|
||||
|
||||
# Таймауты операций страницы.
|
||||
_PAGE_COLLECT_TIMEOUT_S = 90
|
||||
@@ -797,6 +802,11 @@ class DUBIZZLEScraper:
|
||||
|
||||
total = len(records) if effective_only_new else len(all_raw_urls)
|
||||
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_raw_urls if url}
|
||||
all_listing_origin_ids = {
|
||||
origin_id
|
||||
for origin_id in (self._extract_db_origin_id_from_url(url) for url in all_raw_urls)
|
||||
if origin_id
|
||||
}
|
||||
return {
|
||||
"listing": {
|
||||
"status": "ok",
|
||||
@@ -818,6 +828,7 @@ class DUBIZZLEScraper:
|
||||
"protection_events": 0,
|
||||
"failures": failures,
|
||||
"all_listing_origin_urls": all_listing_origin_urls,
|
||||
"all_listing_origin_ids": all_listing_origin_ids,
|
||||
}
|
||||
|
||||
def _sync_listing_streaming(
|
||||
@@ -952,6 +963,11 @@ class DUBIZZLEScraper:
|
||||
|
||||
all_raw_urls = list(algolia_result.vehicle_urls)
|
||||
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_raw_urls if url}
|
||||
all_listing_origin_ids = {
|
||||
str(origin_id).strip()
|
||||
for origin_id in algolia_result.origin_ids_by_url.values()
|
||||
if str(origin_id).strip()
|
||||
}
|
||||
pages_info = list(algolia_result.pages)
|
||||
effective_urls = list(all_raw_urls)
|
||||
|
||||
@@ -1065,6 +1081,7 @@ class DUBIZZLEScraper:
|
||||
"protection_events": protection_events,
|
||||
"failures": failures,
|
||||
"all_listing_origin_urls": all_listing_origin_urls,
|
||||
"all_listing_origin_ids": all_listing_origin_ids,
|
||||
}
|
||||
|
||||
def _sync_listing_streaming_sitemap_fallback(
|
||||
@@ -1107,6 +1124,11 @@ class DUBIZZLEScraper:
|
||||
failures.append({"vehicle_url": f"sitemap_batch_{i}", "error": str(exc)})
|
||||
|
||||
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_urls if url}
|
||||
all_listing_origin_ids = {
|
||||
origin_id
|
||||
for origin_id in (self._extract_db_origin_id_from_url(url) for url in all_urls)
|
||||
if origin_id
|
||||
}
|
||||
return {
|
||||
"listing": {
|
||||
"status": "ok",
|
||||
@@ -1128,6 +1150,7 @@ class DUBIZZLEScraper:
|
||||
"protection_events": 0,
|
||||
"failures": failures,
|
||||
"all_listing_origin_urls": all_listing_origin_urls,
|
||||
"all_listing_origin_ids": all_listing_origin_ids,
|
||||
}
|
||||
|
||||
def _sync_listing_streaming_browser_fallback(
|
||||
@@ -1275,6 +1298,11 @@ class DUBIZZLEScraper:
|
||||
page.close()
|
||||
|
||||
all_listing_origin_urls = {self._normalize_vehicle_url(url) for url in all_raw_urls if url}
|
||||
all_listing_origin_ids = {
|
||||
origin_id
|
||||
for origin_id in (self._extract_db_origin_id_from_url(url) for url in all_raw_urls)
|
||||
if origin_id
|
||||
}
|
||||
total = len(all_raw_urls) if not effective_only_new else max(0, len(all_raw_urls) - skipped_existing)
|
||||
return {
|
||||
"listing": {
|
||||
@@ -1296,6 +1324,7 @@ class DUBIZZLEScraper:
|
||||
"protection_events": 0,
|
||||
"failures": failures,
|
||||
"all_listing_origin_urls": all_listing_origin_urls,
|
||||
"all_listing_origin_ids": all_listing_origin_ids,
|
||||
}
|
||||
|
||||
def _get_page(self) -> Page:
|
||||
@@ -1314,6 +1343,48 @@ class DUBIZZLEScraper:
|
||||
_JS_EXTRACT = """
|
||||
() => {
|
||||
try {
|
||||
const nextData = window.__NEXT_DATA__;
|
||||
const listing = nextData?.props?.pageProps?.reduxWrapperActionsGIPP
|
||||
?.find((entry) => entry && entry.payload && entry.payload.listing)
|
||||
?.payload?.listing;
|
||||
if (listing && listing.details) {
|
||||
const detailSections = listing.details;
|
||||
const normalizeSection = (items) => Array.isArray(items)
|
||||
? items.map((item) => ({
|
||||
label: item?.label || '',
|
||||
value: item?.value ?? '',
|
||||
slug: item?.slug || '',
|
||||
}))
|
||||
: [];
|
||||
return {
|
||||
ok: true,
|
||||
source: 'next_data',
|
||||
listing: {
|
||||
name: listing.name || '',
|
||||
description: listing.description || listing.long_description || '',
|
||||
absolute_url: listing.absolute_url || {},
|
||||
short_url: listing.short_url || '',
|
||||
price: listing.price || {},
|
||||
location: listing.location || {},
|
||||
posted_timestamp: listing.posted_timestamp || null,
|
||||
tracking: listing.tracking || {},
|
||||
categories: Array.isArray(listing.categories) ? listing.categories : [],
|
||||
details: {
|
||||
make_model_trim: normalizeSection(detailSections.make_model_trim),
|
||||
primary: normalizeSection(detailSections.primary),
|
||||
secondary: normalizeSection(detailSections.secondary),
|
||||
rental_details: normalizeSection(detailSections.rental_details),
|
||||
requirements: normalizeSection(detailSections.requirements),
|
||||
},
|
||||
photos: Array.isArray(listing.photos_combined)
|
||||
? listing.photos_combined.map((photo) => photo?.url || photo?.large || photo?.medium || photo?.small || '').filter(Boolean)
|
||||
: Array.isArray(listing.photos)
|
||||
? listing.photos.map((photo) => photo?.url || photo?.large || photo?.medium || photo?.small || '').filter(Boolean)
|
||||
: [],
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
const scripts = document.querySelectorAll('script:not([src])');
|
||||
for (const s of scripts) {
|
||||
const t = s.textContent || '';
|
||||
@@ -1383,6 +1454,142 @@ class DUBIZZLEScraper:
|
||||
if not js_data or not js_data.get("ok"):
|
||||
return {}
|
||||
|
||||
if js_data.get("source") == "next_data":
|
||||
listing = js_data.get("listing") or {}
|
||||
details = listing.get("details") or {}
|
||||
categories = listing.get("categories") or []
|
||||
absolute_url = listing.get("absolute_url") or {}
|
||||
location = listing.get("location") or {}
|
||||
tracking = listing.get("tracking") or {}
|
||||
|
||||
def detail_value_by_slug(target_slug: str) -> Any:
|
||||
target = target_slug.strip().lower()
|
||||
for section_items in details.values():
|
||||
if not isinstance(section_items, list):
|
||||
continue
|
||||
for item in section_items:
|
||||
if not isinstance(item, Mapping):
|
||||
continue
|
||||
if str(item.get("slug") or "").strip().lower() == target:
|
||||
return item.get("value")
|
||||
return None
|
||||
|
||||
def pick_photo_url(photo: Any) -> str:
|
||||
if isinstance(photo, str):
|
||||
return photo.strip()
|
||||
if not isinstance(photo, Mapping):
|
||||
return ""
|
||||
for key in ("url", "main", "large", "medium", "small", "micro"):
|
||||
value = photo.get(key)
|
||||
if value:
|
||||
return str(value).strip()
|
||||
return ""
|
||||
|
||||
detail_sections: dict[str, list[dict[str, Any]]] = {}
|
||||
flat_details: dict[str, dict[str, dict[str, Any]]] = {}
|
||||
for section_name, items in details.items():
|
||||
if not isinstance(items, list):
|
||||
continue
|
||||
normalized_items: list[dict[str, Any]] = []
|
||||
for item in items:
|
||||
if not isinstance(item, Mapping):
|
||||
continue
|
||||
label = str(item.get("label") or "").strip()
|
||||
value = item.get("value")
|
||||
slug = str(item.get("slug") or "").strip()
|
||||
normalized_items.append({"label": {"en": label}, "value": {"en": value}, "slug": slug})
|
||||
if label:
|
||||
flat_details[label] = {"en": {"label": label, "value": value}}
|
||||
if normalized_items:
|
||||
detail_sections[section_name] = normalized_items
|
||||
|
||||
category_v2 = None
|
||||
category_make = None
|
||||
category_model = None
|
||||
if categories:
|
||||
names_en = [str(item.get("name") or "").strip() for item in categories if isinstance(item, Mapping)]
|
||||
slug_paths = [str(item.get("full_slug") or item.get("slug") or "").strip() for item in categories if isinstance(item, Mapping)]
|
||||
ids = [item.get("legacy_id") for item in categories if isinstance(item, Mapping)]
|
||||
category_v2 = {
|
||||
"names_en": names_en,
|
||||
"slug_paths": slug_paths,
|
||||
"ids": ids,
|
||||
}
|
||||
if len(names_en) >= 4:
|
||||
category_make = names_en[2] or None
|
||||
category_model = names_en[3] or None
|
||||
|
||||
photos: list[str] = []
|
||||
for photo in listing.get("photos_combined") or []:
|
||||
picked = pick_photo_url(photo)
|
||||
if picked and picked not in photos:
|
||||
photos.append(picked)
|
||||
for photo in listing.get("photos") or []:
|
||||
picked = pick_photo_url(photo)
|
||||
if picked and picked not in photos:
|
||||
photos.append(picked)
|
||||
|
||||
posted_timestamp = listing.get("posted_timestamp")
|
||||
posted_at_iso = None
|
||||
if isinstance(posted_timestamp, (int, float)) and posted_timestamp > 0:
|
||||
posted_at_iso = datetime.fromtimestamp(posted_timestamp, tz=timezone.utc).isoformat()
|
||||
|
||||
price = listing.get("price") if isinstance(listing.get("price"), Mapping) else {}
|
||||
price_raw = price.get("raw") or price.get("formatted")
|
||||
neighborhood_name = tracking.get("neighbourhood", {}).get("name") if isinstance(tracking.get("neighbourhood"), Mapping) else None
|
||||
year = detail_value_by_slug("year")
|
||||
body_type = detail_value_by_slug("body_type")
|
||||
kilometers = detail_value_by_slug("kilometers")
|
||||
engine_capacity = detail_value_by_slug("engine_capacity_cc")
|
||||
transmission_type = detail_value_by_slug("transmission_type")
|
||||
steering_side = detail_value_by_slug("steering_side")
|
||||
exterior_color = detail_value_by_slug("exterior_color")
|
||||
seller_type = detail_value_by_slug("seller_type")
|
||||
trim = detail_value_by_slug("motors_trim")
|
||||
|
||||
summary: dict[str, Any] = {
|
||||
"source_url": vehicle_url,
|
||||
"name": {"en": listing.get("name")},
|
||||
"title": listing.get("name"),
|
||||
"description": listing.get("description") or listing.get("long_description"),
|
||||
"make": category_make,
|
||||
"model": category_model,
|
||||
"trim": trim,
|
||||
"year": year,
|
||||
"body_type": body_type,
|
||||
"odometer": kilometers,
|
||||
"kilometers": kilometers,
|
||||
"engine": engine_capacity,
|
||||
"engine_volume": engine_capacity,
|
||||
"transmission_type": transmission_type,
|
||||
"gearbox": transmission_type,
|
||||
"steering_side": steering_side,
|
||||
"steering_wheel": steering_side,
|
||||
"exterior_color": exterior_color,
|
||||
"color": exterior_color,
|
||||
"seller": seller_type,
|
||||
"price": price_raw,
|
||||
"buy_now": price_raw,
|
||||
"currency": price.get("currency") or "AED",
|
||||
"details_v2": detail_sections,
|
||||
"details": flat_details,
|
||||
"image_urls": photos,
|
||||
"photo_mains": photos,
|
||||
"category_v2": category_v2,
|
||||
"location": location.get("name"),
|
||||
"location_name": location.get("name"),
|
||||
"site": {"en": "UAE"},
|
||||
"posted_at": posted_at_iso,
|
||||
"permalink": listing.get("short_url") or vehicle_url,
|
||||
"absolute_url": absolute_url,
|
||||
"id": tracking.get("legacy_id") or listing.get("listing_id") or listing.get("object_id"),
|
||||
"objectID": listing.get("encoded_object_id") or listing.get("object_id"),
|
||||
"uuid": listing.get("listing_uuid") or listing.get("uuid"),
|
||||
}
|
||||
if neighborhood_name:
|
||||
summary["neighbourhood"] = {"en": neighborhood_name}
|
||||
return summary
|
||||
|
||||
brnch = str(js_data.get("BranchNumber", "") or "").strip()
|
||||
img_keys = js_data.get("imageKeys") or []
|
||||
image_urls = [
|
||||
@@ -1421,16 +1628,18 @@ class DUBIZZLEScraper:
|
||||
|
||||
@staticmethod
|
||||
def _build_payload_insights(vehicle_summary: dict) -> dict:
|
||||
currency = vehicle_summary.get("currency") or "USD"
|
||||
effective_price = vehicle_summary.get("buy_now") or vehicle_summary.get("price")
|
||||
return {
|
||||
"vehicle_core": vehicle_summary,
|
||||
"pricing": {
|
||||
"buy_now": vehicle_summary.get("buy_now"),
|
||||
"buy_now": effective_price,
|
||||
"current_bid": vehicle_summary.get("current_bid"),
|
||||
"actual_cash_value": vehicle_summary.get("actual_cash_value"),
|
||||
"estimated_repair_cost": vehicle_summary.get("estimated_repair_cost"),
|
||||
"currency": "USD",
|
||||
"currency": currency,
|
||||
},
|
||||
"bids": {"amount": vehicle_summary.get("current_bid"), "currency": "USD"},
|
||||
"bids": {"amount": vehicle_summary.get("current_bid"), "currency": currency},
|
||||
"damage": {"primary": vehicle_summary.get("primary_damage"), "secondary": vehicle_summary.get("secondary_damage")},
|
||||
"auction": {"auction_date": vehicle_summary.get("auction_date"), "branch": vehicle_summary.get("location")},
|
||||
"images": {"count": len(vehicle_summary.get("image_urls") or []), "urls": vehicle_summary.get("image_urls") or []},
|
||||
@@ -1530,6 +1739,27 @@ class DUBIZZLEScraper:
|
||||
Использует json.JSONDecoder.raw_decode для быстрого поиска JSON
|
||||
вместо посимвольного сканирования скобок.
|
||||
"""
|
||||
next_match = NEXT_DATA_RE.search(html)
|
||||
if next_match:
|
||||
try:
|
||||
next_data = json.loads(html_module.unescape(next_match.group(1)))
|
||||
actions = (
|
||||
next_data.get("props", {})
|
||||
.get("pageProps", {})
|
||||
.get("reduxWrapperActionsGIPP", [])
|
||||
)
|
||||
for entry in actions:
|
||||
payload = entry.get("payload") if isinstance(entry, dict) else None
|
||||
listing = payload.get("listing") if isinstance(payload, dict) else None
|
||||
if isinstance(listing, dict) and isinstance(listing.get("details"), dict):
|
||||
return {
|
||||
"ok": True,
|
||||
"source": "next_data",
|
||||
"listing": listing,
|
||||
}
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
search = "inventoryView"
|
||||
pos = html.find(search)
|
||||
if pos < 0:
|
||||
@@ -2259,6 +2489,7 @@ class DUBIZZLEScraper:
|
||||
protection_events = int(stream_result.get("protection_events", 0))
|
||||
failures.extend(stream_result["failures"])
|
||||
all_listing_origin_urls = stream_result["all_listing_origin_urls"]
|
||||
all_listing_origin_ids = stream_result.get("all_listing_origin_ids", set())
|
||||
|
||||
logger.info("Streaming sync processed %d vehicles", total)
|
||||
|
||||
@@ -2271,6 +2502,13 @@ class DUBIZZLEScraper:
|
||||
)
|
||||
if skip_mark_sold:
|
||||
logger.debug("Skipping mark_sold: caller requested")
|
||||
elif all_listing_origin_ids and not is_partial_scan:
|
||||
try:
|
||||
sold_count = self.persistence.mark_sold_not_in_listing(all_listing_origin_ids)
|
||||
if sold_count:
|
||||
logger.info("Marked %d cars as sold by origin_id", sold_count)
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to mark sold cars by origin_id: %s", exc)
|
||||
elif all_listing_origin_urls and not is_partial_scan:
|
||||
try:
|
||||
sold_count = self.persistence.mark_sold_not_in_listing_by_urls(all_listing_origin_urls)
|
||||
@@ -2363,6 +2601,8 @@ class DUBIZZLEScraper:
|
||||
segment_results: list[dict[str, Any]] = []
|
||||
completed_all = True
|
||||
skipped_segments = 0
|
||||
aggregated_active_urls: set[str] = set()
|
||||
aggregated_active_ids: set[str] = set()
|
||||
|
||||
# В segmented-режиме полный прогон должен проходить ВСЕ сегменты.
|
||||
# Runtime-фильтры применяются позже (на уровне конкретных карточек),
|
||||
@@ -2425,6 +2665,7 @@ class DUBIZZLEScraper:
|
||||
listing_url=seg_url,
|
||||
year_min=seg_year_min,
|
||||
year_max=seg_year_max,
|
||||
skip_mark_sold=True,
|
||||
)
|
||||
total_cars_upserted += result.get("cars_upserted", 0)
|
||||
total_cars_failed += result.get("cars_failed", 0)
|
||||
@@ -2432,6 +2673,8 @@ class DUBIZZLEScraper:
|
||||
total_skipped += result.get("skipped_existing", 0)
|
||||
total_discovered += result.get("listing", {}).get("vehicles_collected", 0)
|
||||
all_failures.extend(result.get("failures", []))
|
||||
aggregated_active_urls.update(result.get("all_listing_origin_urls", set()) or set())
|
||||
aggregated_active_ids.update(result.get("all_listing_origin_ids", set()) or set())
|
||||
segment_results.append({
|
||||
"segment": seg,
|
||||
"segment_index": seg_idx,
|
||||
@@ -2514,6 +2757,15 @@ class DUBIZZLEScraper:
|
||||
finally:
|
||||
self.set_progress_callback(outer_progress_callback)
|
||||
|
||||
if completed_all:
|
||||
try:
|
||||
if aggregated_active_ids:
|
||||
self.persistence.mark_sold_not_in_listing(aggregated_active_ids)
|
||||
elif aggregated_active_urls:
|
||||
self.persistence.mark_sold_not_in_listing_by_urls(aggregated_active_urls)
|
||||
except Exception:
|
||||
logger.warning("Segmented sync final sold-mark failed", exc_info=True)
|
||||
|
||||
elapsed = round(time.perf_counter() - started_at, 3)
|
||||
status = "success" if not all_failures else "partial_success" if total_cars_upserted else "failed"
|
||||
|
||||
@@ -2537,6 +2789,8 @@ class DUBIZZLEScraper:
|
||||
"images_upserted": total_images_upserted,
|
||||
"skipped_existing": total_skipped,
|
||||
"total_discovered": total_discovered,
|
||||
"all_listing_origin_urls": aggregated_active_urls,
|
||||
"all_listing_origin_ids": aggregated_active_ids,
|
||||
"elapsed_seconds": elapsed,
|
||||
"failures": all_failures,
|
||||
"segment_results": segment_results,
|
||||
|
||||
@@ -64,6 +64,7 @@ SYNC_SEGMENTS_PROGRESS_KEY = "dubizzle:state:sync_segments_progress"
|
||||
SYNC_SEGMENTS_TOTAL_KEY = "dubizzle:state:sync_segments_total"
|
||||
SYNC_SEGMENTED_SCAN_ACTIVE_KEY = "dubizzle:state:sync_segmented_scan_active"
|
||||
SYNC_SEGMENTED_ACTIVE_URLS_KEY = "dubizzle:state:sync_segmented_active_urls"
|
||||
SYNC_SEGMENTED_ACTIVE_IDS_KEY = "dubizzle:state:sync_segmented_active_ids"
|
||||
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
|
||||
TASK_PROGRESS_KEY_FMT = "dubizzle:state:task_progress:{task_id}"
|
||||
GLOBAL_PROGRESS_TS_KEY = "dubizzle:state:last_progress_ts"
|
||||
@@ -91,6 +92,7 @@ SYNC_SEGMENTS_PROGRESS_KEY = _runtime_key(SYNC_SEGMENTS_PROGRESS_KEY)
|
||||
SYNC_SEGMENTS_TOTAL_KEY = _runtime_key(SYNC_SEGMENTS_TOTAL_KEY)
|
||||
SYNC_SEGMENTED_SCAN_ACTIVE_KEY = _runtime_key(SYNC_SEGMENTED_SCAN_ACTIVE_KEY)
|
||||
SYNC_SEGMENTED_ACTIVE_URLS_KEY = _runtime_key(SYNC_SEGMENTED_ACTIVE_URLS_KEY)
|
||||
SYNC_SEGMENTED_ACTIVE_IDS_KEY = _runtime_key(SYNC_SEGMENTED_ACTIVE_IDS_KEY)
|
||||
TASK_PROGRESS_KEY_FMT = _runtime_key(TASK_PROGRESS_KEY_FMT)
|
||||
GLOBAL_PROGRESS_TS_KEY = _runtime_key(GLOBAL_PROGRESS_TS_KEY)
|
||||
SITEMAP_HOURLY_LAST_COUNT_KEY = _runtime_key(SITEMAP_HOURLY_LAST_COUNT_KEY)
|
||||
@@ -759,6 +761,27 @@ def _add_segmented_active_urls(redis_client: Redis, urls: list[str] | set[str])
|
||||
return 0
|
||||
|
||||
|
||||
def _add_segmented_active_ids(redis_client: Redis, origin_ids: list[str] | set[str]) -> int:
|
||||
normalized_ids = [str(origin_id).strip() for origin_id in origin_ids if str(origin_id).strip()]
|
||||
if not normalized_ids:
|
||||
return 0
|
||||
|
||||
try:
|
||||
added = 0
|
||||
pipe = redis_client.pipeline()
|
||||
chunk_size = 5000
|
||||
for start in range(0, len(normalized_ids), chunk_size):
|
||||
chunk = normalized_ids[start:start + chunk_size]
|
||||
pipe.sadd(SYNC_SEGMENTED_ACTIVE_IDS_KEY, *chunk)
|
||||
pipe.expire(SYNC_SEGMENTED_ACTIVE_IDS_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
|
||||
results = pipe.execute()
|
||||
added += int(results[0] or 0)
|
||||
return added
|
||||
except Exception:
|
||||
logger.warning("Failed to persist segmented active origin_ids", exc_info=True)
|
||||
return 0
|
||||
|
||||
|
||||
def _load_segmented_active_urls(redis_client: Redis) -> set[str]:
|
||||
try:
|
||||
raw = redis_client.smembers(SYNC_SEGMENTED_ACTIVE_URLS_KEY) or set()
|
||||
@@ -768,6 +791,15 @@ def _load_segmented_active_urls(redis_client: Redis) -> set[str]:
|
||||
return {str(url).strip() for url in raw if str(url).strip()}
|
||||
|
||||
|
||||
def _load_segmented_active_ids(redis_client: Redis) -> set[str]:
|
||||
try:
|
||||
raw = redis_client.smembers(SYNC_SEGMENTED_ACTIVE_IDS_KEY) or set()
|
||||
except Exception:
|
||||
logger.warning("Failed to read segmented active origin_ids", exc_info=True)
|
||||
return set()
|
||||
return {str(origin_id).strip() for origin_id in raw if str(origin_id).strip()}
|
||||
|
||||
|
||||
def _clear_segmented_active_urls(redis_client: Redis) -> None:
|
||||
try:
|
||||
redis_client.delete(SYNC_SEGMENTED_ACTIVE_URLS_KEY)
|
||||
@@ -775,6 +807,13 @@ def _clear_segmented_active_urls(redis_client: Redis) -> None:
|
||||
logger.warning("Failed to clear segmented active URLs", exc_info=True)
|
||||
|
||||
|
||||
def _clear_segmented_active_ids(redis_client: Redis) -> None:
|
||||
try:
|
||||
redis_client.delete(SYNC_SEGMENTED_ACTIVE_IDS_KEY)
|
||||
except Exception:
|
||||
logger.warning("Failed to clear segmented active origin_ids", exc_info=True)
|
||||
|
||||
|
||||
def _is_segmented_scan_in_progress(redis_client: Redis, celery_app) -> bool:
|
||||
try:
|
||||
marker = redis_client.get(SYNC_SEGMENTED_SCAN_ACTIVE_KEY)
|
||||
@@ -1036,6 +1075,7 @@ def _reset_segments_progress(redis_client: Redis, total: int) -> None:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.delete(SYNC_SEGMENTS_PROGRESS_KEY)
|
||||
pipe.delete(SYNC_SEGMENTED_ACTIVE_URLS_KEY)
|
||||
pipe.delete(SYNC_SEGMENTED_ACTIVE_IDS_KEY)
|
||||
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
|
||||
pipe.execute()
|
||||
except Exception:
|
||||
@@ -1152,9 +1192,12 @@ def sync_segment_task(
|
||||
segment_done = bool(result.get("full_scan_completed", False))
|
||||
listing_payload = result.get("listing") if isinstance(result.get("listing"), dict) else {}
|
||||
active_urls = listing_payload.get("vehicle_urls") or []
|
||||
active_ids = result.get("all_listing_origin_ids") or []
|
||||
active_urls_count = len(active_urls)
|
||||
if active_urls:
|
||||
_add_segmented_active_urls(redis_client, active_urls)
|
||||
if active_ids:
|
||||
_add_segmented_active_ids(redis_client, active_ids)
|
||||
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
@@ -1176,8 +1219,22 @@ def sync_segment_task(
|
||||
)
|
||||
if total > 0 and completed >= total:
|
||||
sold_count = 0
|
||||
active_ids_for_scan = _load_segmented_active_ids(redis_client)
|
||||
active_urls_for_scan = _load_segmented_active_urls(redis_client)
|
||||
if active_urls_for_scan:
|
||||
if active_ids_for_scan:
|
||||
try:
|
||||
sold_count = persistence.mark_sold_not_in_listing(
|
||||
active_ids_for_scan,
|
||||
lane=_origin_prefix().rstrip(":"),
|
||||
)
|
||||
logger.info(
|
||||
"Segmented full scan sold-mark completed by origin_id: sold=%d active_ids=%d",
|
||||
sold_count,
|
||||
len(active_ids_for_scan),
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Segmented full scan sold-mark by origin_id failed", exc_info=True)
|
||||
elif active_urls_for_scan:
|
||||
try:
|
||||
sold_count = persistence.mark_sold_not_in_listing_by_urls(
|
||||
active_urls_for_scan,
|
||||
@@ -1198,6 +1255,7 @@ def sync_segment_task(
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_segmented_scan_active(redis_client)
|
||||
_clear_segmented_active_urls(redis_client)
|
||||
_clear_segmented_active_ids(redis_client)
|
||||
_clear_bootstrap_failure_streak(redis_client)
|
||||
_clear_bootstrap_continuation_streak(redis_client)
|
||||
if always_full_scan:
|
||||
|
||||
352
scripts/export_for_upload.py
Normal file
352
scripts/export_for_upload.py
Normal file
@@ -0,0 +1,352 @@
|
||||
|
||||
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()
|
||||
@@ -122,6 +122,50 @@ class TestPersistenceServiceIntegration(unittest.TestCase):
|
||||
self.assertEqual(len(cars), 1)
|
||||
self.assertEqual(cars[0].origin_id, "NEW-ID")
|
||||
|
||||
def test_mark_sold_by_urls_marks_missing_records(self) -> None:
|
||||
first = self._record("111")
|
||||
second = self._record("222")
|
||||
first.origin_id = "dubizzle:111"
|
||||
second.origin_id = "dubizzle:222"
|
||||
self.persistence.upsert_car(first)
|
||||
self.persistence.upsert_car(second)
|
||||
|
||||
marked = self.persistence.mark_sold_not_in_listing_by_urls({first.origin_url})
|
||||
|
||||
self.assertEqual(marked, 1)
|
||||
with self.persistence.session_scope() as session:
|
||||
cars = {
|
||||
car.origin_id: car
|
||||
for car in session.execute(select(Car)).scalars().all()
|
||||
}
|
||||
|
||||
self.assertFalse(cars["dubizzle:111"].is_sold)
|
||||
self.assertTrue(cars["dubizzle:222"].is_sold)
|
||||
|
||||
def test_upsert_reactivates_previously_sold_car(self) -> None:
|
||||
record = self._record("333", price=1000)
|
||||
record.origin_id = "dubizzle:333"
|
||||
self.persistence.upsert_car(record)
|
||||
self.persistence.mark_sold_not_in_listing_by_urls({"https://www.dubizzle.com/VehicleDetail/other~US"})
|
||||
|
||||
with self.persistence.session_scope() as session:
|
||||
sold_car = session.execute(
|
||||
select(Car).where(Car.origin_id == "dubizzle:333")
|
||||
).scalar_one()
|
||||
self.assertTrue(sold_car.is_sold)
|
||||
|
||||
revived = self._record("333", price=1500)
|
||||
revived.origin_id = "dubizzle:333"
|
||||
self.persistence.upsert_car(revived)
|
||||
|
||||
with self.persistence.session_scope() as session:
|
||||
active_car = session.execute(
|
||||
select(Car).where(Car.origin_id == "dubizzle:333")
|
||||
).scalar_one()
|
||||
|
||||
self.assertFalse(active_car.is_sold)
|
||||
self.assertEqual(active_car.price, 1500)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
||||
@@ -269,6 +269,69 @@ class TestScraperSync(unittest.TestCase):
|
||||
self.assertEqual(result["cars_upserted"], 1)
|
||||
self.assertEqual(result["total"], 1)
|
||||
|
||||
def test_sync_listing_marks_sold_after_full_scan(self) -> None:
|
||||
scraper = self._make_scraper()
|
||||
scraper.persistence.create_tables = MagicMock()
|
||||
scraper.persistence.start_sync_run = MagicMock(return_value=7)
|
||||
scraper.persistence.finish_sync_run = MagicMock()
|
||||
scraper.persistence.mark_sold_not_in_listing = MagicMock(return_value=1)
|
||||
scraper.persistence.mark_sold_not_in_listing_by_urls = MagicMock(return_value=0)
|
||||
scraper._sync_listing_streaming = MagicMock(return_value={
|
||||
"listing": {
|
||||
"vehicles_collected": 2,
|
||||
"early_stopped": False,
|
||||
"truncated_by_time_budget": False,
|
||||
},
|
||||
"total": 2,
|
||||
"skipped_existing": 0,
|
||||
"cars_upserted": 2,
|
||||
"cars_failed": 0,
|
||||
"images_upserted": 4,
|
||||
"failures": [],
|
||||
"all_listing_origin_urls": {
|
||||
"https://www.dubizzle.com/VehicleDetail/111~US",
|
||||
"https://www.dubizzle.com/VehicleDetail/222~US",
|
||||
},
|
||||
"all_listing_origin_ids": {
|
||||
"dubizzle:111",
|
||||
"dubizzle:222",
|
||||
},
|
||||
})
|
||||
|
||||
result = scraper.sync_listing(only_new=False, limit=None)
|
||||
|
||||
self.assertEqual(result["status"], "success")
|
||||
scraper.persistence.mark_sold_not_in_listing.assert_called_once()
|
||||
scraper.persistence.mark_sold_not_in_listing_by_urls.assert_not_called()
|
||||
|
||||
def test_sync_listing_skips_mark_sold_for_partial_scan(self) -> None:
|
||||
scraper = self._make_scraper()
|
||||
scraper.persistence.create_tables = MagicMock()
|
||||
scraper.persistence.start_sync_run = MagicMock(return_value=8)
|
||||
scraper.persistence.finish_sync_run = MagicMock()
|
||||
scraper.persistence.mark_sold_not_in_listing_by_urls = MagicMock()
|
||||
scraper._sync_listing_streaming = MagicMock(return_value={
|
||||
"listing": {
|
||||
"vehicles_collected": 1,
|
||||
"early_stopped": False,
|
||||
"truncated_by_time_budget": False,
|
||||
},
|
||||
"total": 1,
|
||||
"skipped_existing": 0,
|
||||
"cars_upserted": 1,
|
||||
"cars_failed": 0,
|
||||
"images_upserted": 2,
|
||||
"failures": [],
|
||||
"all_listing_origin_urls": {
|
||||
"https://www.dubizzle.com/VehicleDetail/111~US",
|
||||
},
|
||||
})
|
||||
|
||||
result = scraper.sync_listing(only_new=True)
|
||||
|
||||
self.assertEqual(result["status"], "success")
|
||||
scraper.persistence.mark_sold_not_in_listing_by_urls.assert_not_called()
|
||||
|
||||
def test_sync_batch_timeout_does_not_duplicate_already_processed_urls(self) -> None:
|
||||
scraper = self._make_scraper()
|
||||
scraper.persistence.upsert_cars_batch = MagicMock(return_value={
|
||||
|
||||
Reference in New Issue
Block a user