494 lines
18 KiB
Python
494 lines
18 KiB
Python
import logging
|
||
from contextlib import contextmanager
|
||
from datetime import datetime, timezone
|
||
from typing import Any, Iterator
|
||
|
||
from sqlalchemy import create_engine, delete, or_, select, text, update
|
||
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
||
from sqlalchemy.orm import Session, sessionmaker
|
||
|
||
from ..core.config import Settings
|
||
from .models import Base, Car, Image, SyncRun
|
||
from .schemas import CarRecord
|
||
|
||
logger = logging.getLogger("encar_scraper.db")
|
||
|
||
|
||
CAR_DB_FIELDS = {
|
||
col.key for col in Car.__table__.columns
|
||
if col.key not in ("id",)
|
||
}
|
||
|
||
_IN_CHUNK_SIZE = 5000
|
||
|
||
|
||
class PersistenceService:
|
||
|
||
def __init__(self, settings: Settings) -> None:
|
||
self.settings = settings
|
||
engine_kwargs = {
|
||
"echo": settings.database.echo,
|
||
"future": True,
|
||
}
|
||
if "postgresql" in settings.database.url:
|
||
engine_kwargs["pool_size"] = settings.database.pool_size
|
||
engine_kwargs["max_overflow"] = settings.database.max_overflow
|
||
engine_kwargs["pool_pre_ping"] = True
|
||
engine_kwargs["pool_recycle"] = settings.database.pool_recycle_seconds
|
||
self.engine = create_engine(settings.database.url, **engine_kwargs)
|
||
self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False, future=True)
|
||
|
||
def create_tables(self) -> None:
|
||
# В тестах/локально на SQLite разрешаем create_all; для non-SQLite в проде — только через миграции.
|
||
is_sqlite = self.settings.database.url.startswith("sqlite")
|
||
if not is_sqlite and not self.settings.database.auto_create_tables:
|
||
return
|
||
try:
|
||
Base.metadata.create_all(self.engine)
|
||
except Exception:
|
||
logger.debug("create_tables skipped (schema already exists)")
|
||
|
||
@contextmanager
|
||
def session_scope(self) -> Iterator[Session]:
|
||
session = self.session_factory()
|
||
try:
|
||
yield session
|
||
session.commit()
|
||
except Exception:
|
||
session.rollback()
|
||
raise
|
||
finally:
|
||
session.close()
|
||
|
||
def start_sync_run(self, lane: str) -> int:
|
||
with self.session_scope() as session:
|
||
now = datetime.now(timezone.utc)
|
||
stale_runs = session.execute(select(SyncRun).where(SyncRun.status == "running")).scalars().all()
|
||
for stale in stale_runs:
|
||
stale.status = "failed"
|
||
stale.finished_at = now
|
||
if not stale.error_summary:
|
||
stale.error_summary = "Recovered stale running sync run before starting a new run"
|
||
|
||
run = SyncRun(status="running", lane=lane, ids_fetched=0, cars_upserted=0, cars_failed=0, images_upserted=0)
|
||
session.add(run)
|
||
session.flush()
|
||
return int(run.id)
|
||
|
||
def finish_sync_run(self, run_id: int, *, status: str, ids_fetched: int, cars_upserted: int, cars_failed: int, images_upserted: int, error_summary: str | None = None) -> None:
|
||
with self.session_scope() as session:
|
||
run = session.get(SyncRun, run_id)
|
||
if run is None:
|
||
return
|
||
run.finished_at = datetime.now(timezone.utc)
|
||
run.status = status
|
||
run.ids_fetched = ids_fetched
|
||
run.cars_upserted = cars_upserted
|
||
run.cars_failed = cars_failed
|
||
run.images_upserted = images_upserted
|
||
run.error_summary = error_summary
|
||
|
||
@staticmethod
|
||
def _add_images(session: Session, car_id: int, images: list[dict[str, object]]) -> None:
|
||
if not images:
|
||
return
|
||
session.add_all([
|
||
Image(
|
||
fullres_image=str(img["fullres_image"]),
|
||
preview_image=str(img["preview_image"]),
|
||
order_index=int(img.get("order_index", 0)),
|
||
car_id=car_id,
|
||
)
|
||
for img in images
|
||
])
|
||
|
||
@staticmethod
|
||
def _car_payload(record: CarRecord) -> dict[str, object]:
|
||
payload = record.model_dump(mode="python")
|
||
return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
|
||
|
||
def _is_postgres(self) -> bool:
|
||
return self.engine.dialect.name == "postgresql"
|
||
|
||
def _load_existing_cars(
|
||
self,
|
||
session: Session,
|
||
origin_ids: list[str],
|
||
origin_urls: list[str],
|
||
) -> tuple[dict[str, Car], dict[str, Car]]:
|
||
existing_by_id: dict[str, Car] = {}
|
||
existing_by_url: dict[str, Car] = {}
|
||
if not origin_ids and not origin_urls:
|
||
return existing_by_id, existing_by_url
|
||
|
||
existing_cars: list[Car] = []
|
||
max_len = max(len(origin_ids), len(origin_urls), 1)
|
||
for i in range(0, max_len, _IN_CHUNK_SIZE):
|
||
id_chunk = origin_ids[i:i + _IN_CHUNK_SIZE]
|
||
url_chunk = origin_urls[i:i + _IN_CHUNK_SIZE]
|
||
conditions = []
|
||
if id_chunk:
|
||
conditions.append(Car.origin_id.in_(id_chunk))
|
||
if url_chunk:
|
||
conditions.append(Car.origin_url.in_(url_chunk))
|
||
if not conditions:
|
||
continue
|
||
rows = session.execute(
|
||
select(Car).where(or_(*conditions))
|
||
).scalars().all()
|
||
existing_cars.extend(rows)
|
||
|
||
for car in existing_cars:
|
||
if car.origin_id:
|
||
existing_by_id[car.origin_id] = car
|
||
if car.origin_url:
|
||
existing_by_url[car.origin_url] = car
|
||
return existing_by_id, existing_by_url
|
||
|
||
def _load_existing_image_urls(self, session: Session, car_ids: set[int]) -> dict[int, set[str]]:
|
||
existing_images_map: dict[int, set[str]] = {}
|
||
if not car_ids:
|
||
return existing_images_map
|
||
|
||
car_id_list = list(car_ids)
|
||
for i in range(0, len(car_id_list), _IN_CHUNK_SIZE):
|
||
chunk = car_id_list[i:i + _IN_CHUNK_SIZE]
|
||
img_rows = session.execute(
|
||
select(Image.car_id, Image.fullres_image).where(Image.car_id.in_(chunk))
|
||
).all()
|
||
for cid, furl in img_rows:
|
||
existing_images_map.setdefault(int(cid), set()).add(str(furl))
|
||
return existing_images_map
|
||
|
||
@staticmethod
|
||
def _postgres_upsert_set_map(insert_stmt) -> dict[str, object]:
|
||
return {
|
||
key: getattr(insert_stmt.excluded, key)
|
||
for key in CAR_DB_FIELDS
|
||
}
|
||
|
||
def _replace_images_for_car(
|
||
self,
|
||
session: Session,
|
||
car_id: int,
|
||
images: list[dict[str, object]],
|
||
origin_id: str,
|
||
) -> int:
|
||
nested = session.begin_nested()
|
||
try:
|
||
session.execute(delete(Image).where(Image.car_id == car_id))
|
||
self._add_images(session, car_id, images)
|
||
session.flush()
|
||
nested.commit()
|
||
return len(images)
|
||
except Exception:
|
||
nested.rollback()
|
||
logger.warning("Image replacement failed for car %s, keeping old images", origin_id, exc_info=True)
|
||
return 0
|
||
|
||
def _upsert_cars_batch_postgres(self, records: list[CarRecord]) -> dict[str, int]:
|
||
inserted = 0
|
||
updated = 0
|
||
images_total = 0
|
||
|
||
with self.session_scope() as session:
|
||
origin_ids = [r.origin_id for r in records if r.origin_id]
|
||
origin_urls = [r.origin_url for r in records if r.origin_url]
|
||
existing_by_id, existing_by_url = self._load_existing_cars(session, origin_ids, origin_urls)
|
||
|
||
entries: list[dict[str, object]] = []
|
||
upsert_payloads: list[dict[str, object]] = []
|
||
for record in records:
|
||
payload = self._car_payload(record)
|
||
images = [image.model_dump(mode="python") for image in record.images]
|
||
car_by_id = existing_by_id.get(record.origin_id)
|
||
car_by_url = existing_by_url.get(record.origin_url)
|
||
entry: dict[str, object] = {"record": record, "images": images, "car_id": None}
|
||
|
||
if car_by_url is not None and car_by_url.origin_id != record.origin_id and car_by_id is None:
|
||
for key, value in payload.items():
|
||
setattr(car_by_url, key, value)
|
||
car_by_url.last_seen_at = record.last_seen_at
|
||
entry["car_id"] = int(car_by_url.id)
|
||
updated += 1
|
||
else:
|
||
upsert_payloads.append(payload)
|
||
if car_by_id is not None:
|
||
updated += 1
|
||
else:
|
||
inserted += 1
|
||
entries.append(entry)
|
||
|
||
session.flush()
|
||
|
||
if upsert_payloads:
|
||
insert_stmt = pg_insert(Car).values(upsert_payloads)
|
||
upsert_stmt = insert_stmt.on_conflict_do_update(
|
||
index_elements=[Car.origin_id],
|
||
set_=self._postgres_upsert_set_map(insert_stmt),
|
||
).returning(Car.id, Car.origin_id)
|
||
rows = session.execute(upsert_stmt).all()
|
||
car_ids_by_origin_id = {str(origin_id): int(car_id) for car_id, origin_id in rows}
|
||
for entry in entries:
|
||
if entry["car_id"] is not None:
|
||
continue
|
||
record = entry["record"]
|
||
car_id = car_ids_by_origin_id.get(record.origin_id)
|
||
if car_id is None:
|
||
raise RuntimeError(f"PostgreSQL upsert did not return car_id for {record.origin_id}")
|
||
entry["car_id"] = car_id
|
||
|
||
car_ids = {int(entry["car_id"]) for entry in entries if entry["car_id"] is not None}
|
||
existing_images_map = self._load_existing_image_urls(session, car_ids)
|
||
images_by_car_id: dict[int, list[dict[str, object]]] = {}
|
||
replace_ids: list[int] = []
|
||
|
||
for entry in entries:
|
||
car_id = int(entry["car_id"])
|
||
images = entry["images"]
|
||
new_image_urls = {
|
||
str(img.get("fullres_image", ""))
|
||
for img in images
|
||
if img.get("fullres_image")
|
||
}
|
||
old_image_urls = existing_images_map.get(car_id, set())
|
||
if new_image_urls != old_image_urls:
|
||
replace_ids.append(car_id)
|
||
images_by_car_id[car_id] = images
|
||
else:
|
||
images_total += len(old_image_urls)
|
||
|
||
if replace_ids:
|
||
for i in range(0, len(replace_ids), _IN_CHUNK_SIZE):
|
||
chunk = replace_ids[i:i + _IN_CHUNK_SIZE]
|
||
session.execute(delete(Image).where(Image.car_id.in_(chunk)))
|
||
for car_id in replace_ids:
|
||
images = images_by_car_id[car_id]
|
||
self._add_images(session, car_id, images)
|
||
images_total += len(images)
|
||
|
||
return {"inserted": inserted, "updated": updated, "images_upserted": images_total}
|
||
|
||
def upsert_car(self, record: CarRecord):
|
||
# Вставка или обновление автомобиля по origin_id/origin_url.
|
||
payload = self._car_payload(record)
|
||
images = [image.model_dump(mode="python") for image in record.images]
|
||
with self.session_scope() as session:
|
||
if self._is_postgres():
|
||
car_by_url = session.execute(
|
||
select(Car).where(Car.origin_url == record.origin_url)
|
||
).scalar_one_or_none()
|
||
if car_by_url is not None and car_by_url.origin_id != record.origin_id:
|
||
for key, value in payload.items():
|
||
setattr(car_by_url, key, value)
|
||
car_by_url.last_seen_at = record.last_seen_at
|
||
session.flush()
|
||
car_id = int(car_by_url.id)
|
||
action = "updated"
|
||
else:
|
||
existed = session.execute(
|
||
select(Car.id).where(Car.origin_id == record.origin_id)
|
||
).scalar_one_or_none() is not None
|
||
insert_stmt = pg_insert(Car).values(**payload)
|
||
upsert_stmt = insert_stmt.on_conflict_do_update(
|
||
index_elements=[Car.origin_id],
|
||
set_=self._postgres_upsert_set_map(insert_stmt),
|
||
).returning(Car.id)
|
||
car_id = int(session.execute(upsert_stmt).scalar_one())
|
||
action = "updated" if existed or car_by_url is not None else "inserted"
|
||
|
||
images_upserted = self._replace_images_for_car(
|
||
session, car_id, images, record.origin_id,
|
||
)
|
||
return {"car_id": car_id, "images_upserted": images_upserted, "action": action}
|
||
|
||
car = session.execute(
|
||
select(Car).where(or_(Car.origin_id == record.origin_id, Car.origin_url == record.origin_url))
|
||
).scalar_one_or_none()
|
||
action = "inserted"
|
||
if car is None:
|
||
car = Car(**payload)
|
||
session.add(car)
|
||
session.flush()
|
||
else:
|
||
action = "updated"
|
||
for key, value in payload.items():
|
||
setattr(car, key, value)
|
||
car.last_seen_at = record.last_seen_at
|
||
session.flush()
|
||
images_upserted = self._replace_images_for_car(session, int(car.id), images, record.origin_id)
|
||
return {"car_id": int(car.id), "images_upserted": images_upserted, "action": action}
|
||
self._add_images(session, int(car.id), images)
|
||
session.flush()
|
||
return {"car_id": int(car.id), "images_upserted": len(images), "action": action}
|
||
def upsert_cars_batch(self, records: list[CarRecord]) -> dict[str, int]:
|
||
"""Пакетный upsert нескольких автомобилей в одной транзакции.
|
||
|
||
Оптимизации:
|
||
- Дедупликация записей по origin_id перед вставкой.
|
||
- Chunked IN-queries для больших списков (обход лимита PG параметров).
|
||
- Пропуск перезаписи изображений, если набор URL не изменился.
|
||
- Один DELETE по car_id IN (...) вместо удаления по одному.
|
||
- Fallback на по-одному upsert если batch commit упал.
|
||
"""
|
||
# ── Дедупликация записей внутри батча ──
|
||
seen_ids: dict[str, int] = {}
|
||
unique_records: list[CarRecord] = []
|
||
for idx, r in enumerate(records):
|
||
key = r.origin_id or r.origin_url
|
||
if key in seen_ids:
|
||
logger.debug("Dedup: skipping duplicate record %s (idx %d vs %d)", key, idx, seen_ids[key])
|
||
continue
|
||
seen_ids[key] = idx
|
||
unique_records.append(r)
|
||
|
||
if len(unique_records) < len(records):
|
||
logger.info("Deduped batch: %d → %d records", len(records), len(unique_records))
|
||
records = unique_records
|
||
|
||
try:
|
||
return self._upsert_cars_batch_inner(records)
|
||
except Exception as exc:
|
||
logger.warning("Batch upsert failed (%s), falling back to individual upserts", exc)
|
||
return self._upsert_cars_individually(records)
|
||
|
||
def _upsert_cars_batch_inner(self, records: list[CarRecord]) -> dict[str, int]:
|
||
"""Внутренняя реализация batched upsert (одна транзакция)."""
|
||
if self._is_postgres():
|
||
return self._upsert_cars_batch_postgres(records)
|
||
|
||
inserted = 0
|
||
updated = 0
|
||
images_total = 0
|
||
|
||
with self.session_scope() as session:
|
||
# Получаем существующие записи chunked-запросами.
|
||
origin_ids = [r.origin_id for r in records if r.origin_id]
|
||
origin_urls = [r.origin_url for r in records if r.origin_url]
|
||
|
||
existing_by_id, existing_by_url = self._load_existing_cars(session, origin_ids, origin_urls)
|
||
|
||
# Предзагружаем ВСЕ изображения для обновляемых машин одним запросом.
|
||
existing_car_ids = set()
|
||
for record in records:
|
||
car = existing_by_id.get(record.origin_id) or existing_by_url.get(record.origin_url)
|
||
if car is not None:
|
||
existing_car_ids.add(int(car.id))
|
||
|
||
# Строим маппинг car_id → set(image_urls) для сравнения.
|
||
existing_images_map = self._load_existing_image_urls(session, existing_car_ids)
|
||
|
||
new_cars: list[tuple[Car, list[dict]]] = []
|
||
update_cars_needing_images: list[tuple[Car, list[dict]]] = []
|
||
|
||
for record in records:
|
||
payload = self._car_payload(record)
|
||
images = [image.model_dump(mode="python") for image in record.images]
|
||
|
||
car = existing_by_id.get(record.origin_id) or existing_by_url.get(record.origin_url)
|
||
if car is None:
|
||
car = Car(**payload)
|
||
session.add(car)
|
||
inserted += 1
|
||
new_cars.append((car, images))
|
||
else:
|
||
for key, value in payload.items():
|
||
setattr(car, key, value)
|
||
car.last_seen_at = record.last_seen_at
|
||
updated += 1
|
||
|
||
# Проверяем, изменились ли изображения.
|
||
new_image_urls = {img.get("fullres_image", "") for img in images}
|
||
old_image_urls = existing_images_map.get(int(car.id), set())
|
||
if new_image_urls != old_image_urls:
|
||
update_cars_needing_images.append((car, images))
|
||
else:
|
||
images_total += len(old_image_urls)
|
||
|
||
# Один flush для всех вставок.
|
||
session.flush()
|
||
|
||
# Добавляем изображения для новых автомобилей.
|
||
for car, images in new_cars:
|
||
self._add_images(session, int(car.id), images)
|
||
images_total += len(images)
|
||
|
||
# Массово обновляем изображения только для машин с изменёнными картинками.
|
||
if update_cars_needing_images:
|
||
update_ids = [int(car.id) for car, _ in update_cars_needing_images]
|
||
for i in range(0, len(update_ids), _IN_CHUNK_SIZE):
|
||
chunk = update_ids[i:i + _IN_CHUNK_SIZE]
|
||
session.execute(delete(Image).where(Image.car_id.in_(chunk)))
|
||
for car, images in update_cars_needing_images:
|
||
self._add_images(session, int(car.id), images)
|
||
images_total += len(images)
|
||
|
||
return {"inserted": inserted, "updated": updated, "images_upserted": images_total}
|
||
|
||
def _upsert_cars_individually(self, records: list[CarRecord]) -> dict[str, int]:
|
||
# Fallback на поштучный upsert.
|
||
inserted = 0
|
||
updated = 0
|
||
images_total = 0
|
||
for record in records:
|
||
try:
|
||
result = self.upsert_car(record)
|
||
action = result.get("action", "inserted")
|
||
if action == "inserted":
|
||
inserted += 1
|
||
else:
|
||
updated += 1
|
||
images_total += int(result.get("images_upserted", 0))
|
||
except Exception as exc:
|
||
logger.error("Individual upsert failed for %s: %s", record.origin_id, exc)
|
||
return {"inserted": inserted, "updated": updated, "images_upserted": images_total}
|
||
|
||
def mark_sold_not_in_listing(self, active_origin_ids: set[str], lane: str = "encar") -> int:
|
||
"""Помечает авто как проданные, если их нет в активном листинге (по origin_id).
|
||
|
||
Для PostgreSQL использует временную таблицу + LEFT JOIN вместо NOT IN,
|
||
что кардинально быстрее при больших объёмах (200K+ IDs).
|
||
"""
|
||
if not active_origin_ids:
|
||
return 0
|
||
with self.session_scope() as session:
|
||
if self._is_postgres():
|
||
session.execute(text("CREATE TEMP TABLE IF NOT EXISTS _active_ids (origin_id TEXT NOT NULL) ON COMMIT DROP"))
|
||
session.execute(text("TRUNCATE _active_ids"))
|
||
|
||
id_list = list(active_origin_ids)
|
||
for i in range(0, len(id_list), _IN_CHUNK_SIZE):
|
||
chunk = id_list[i:i + _IN_CHUNK_SIZE]
|
||
values = ",".join(f"(:{f'u{j}'})" for j in range(len(chunk)))
|
||
params = {f"u{j}": oid for j, oid in enumerate(chunk)}
|
||
session.execute(text(f"INSERT INTO _active_ids (origin_id) VALUES {values}"), params)
|
||
|
||
session.execute(text("CREATE INDEX IF NOT EXISTS _ix_active_ids ON _active_ids (origin_id)"))
|
||
|
||
result = session.execute(text("""
|
||
UPDATE cars
|
||
SET is_sold = TRUE
|
||
FROM (
|
||
SELECT c.id
|
||
FROM cars c
|
||
LEFT JOIN _active_ids a ON c.origin_id = a.origin_id
|
||
WHERE a.origin_id IS NULL
|
||
AND c.is_sold = FALSE
|
||
AND c.origin_id LIKE 'encar:%%'
|
||
) sub
|
||
WHERE cars.id = sub.id
|
||
"""))
|
||
count = result.rowcount or 0
|
||
else:
|
||
stmt = (
|
||
update(Car)
|
||
.where(Car.origin_id.notin_(active_origin_ids))
|
||
.where(Car.is_sold == False) # noqa: E712
|
||
.where(Car.origin_id.like("encar:%"))
|
||
.values(is_sold=True)
|
||
)
|
||
result = session.execute(stmt)
|
||
count = result.rowcount or 0
|
||
if count:
|
||
logger.info("Marked %d cars as sold (no longer in listing)", count)
|
||
return count |