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("dubizzle_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 engine_kwargs["pool_timeout"] = 30 # Таймауты запросов и блокировок. engine_kwargs["connect_args"] = { "options": "-c statement_timeout=120000 -c lock_timeout=30000" } 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. 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 def get_existing_origin_urls(self, origin_urls: list[str]) -> set[str]: if not origin_urls: return set() with self.session_scope() as session: rows = session.execute(select(Car.origin_url).where(Car.origin_url.in_(origin_urls))).all() return {str(row[0]) for row in rows if row and row[0]} def get_existing_origin_ids(self, origin_ids: list[str]) -> set[str]: if not origin_ids: return set() with self.session_scope() as session: rows = session.execute(select(Car.origin_id).where(Car.origin_id.in_(origin_ids))).all() return {str(row[0]) for row in rows if row and row[0]} def get_existing_urls_and_ids( self, origin_urls: list[str], origin_ids: list[str], ) -> tuple[set[str], set[str]]: # Загрузка URL и origin_id. if not origin_urls and not origin_ids: return set(), set() urls: set[str] = set() ids: set[str] = set() with self.session_scope() as session: max_len = max(len(origin_urls), len(origin_ids), 1) for i in range(0, max_len, _IN_CHUNK_SIZE): url_chunk = origin_urls[i:i + _IN_CHUNK_SIZE] id_chunk = origin_ids[i:i + _IN_CHUNK_SIZE] conditions = [] if url_chunk: conditions.append(Car.origin_url.in_(url_chunk)) if id_chunk: conditions.append(Car.origin_id.in_(id_chunk)) if not conditions: continue rows = session.execute( select(Car.origin_url, Car.origin_id).where(or_(*conditions)) ).all() for r in rows: if r[0]: urls.add(str(r[0])) if r[1]: ids.add(str(r[1])) return urls, ids @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): # Вставка или обновление авто. 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: # Загружаем существующие записи. 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)) # Готовим map car_id -> 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]: # Запасной поштучный 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 = "dubizzle") -> int: """Помечает авто как проданные, если их нет в активном листинге (по origin_id).""" if not active_origin_ids: return 0 with self.session_scope() as session: stmt = ( update(Car) .where(Car.origin_id.notin_(active_origin_ids)) .where(Car.is_sold == False) # noqa: E712 .where(Car.origin_id.like("dubizzle:%")) .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 def mark_sold_not_in_listing_by_urls(self, active_origin_urls: set[str], lane: str = "dubizzle") -> int: """Помечает авто как проданные, если их URL нет в активном листинге. Для PostgreSQL использует временную таблицу + LEFT JOIN вместо NOT IN, что кардинально быстрее при больших объёмах (100K+ URLs). Встроенная защита: нормализация URL (тильда/дефис) + safety-check на аномальный процент. """ if not active_origin_urls: return 0 # Нормализация URL. def _norm(url: str) -> str: return url.replace("~", "-") normalized_urls = {_norm(u) for u in active_origin_urls} is_postgres = "postgresql" in self.settings.database.url with self.session_scope() as session: if is_postgres: # Считаем кандидатов на sold. total_active = session.execute( text("SELECT count(*) FROM cars WHERE is_sold = FALSE AND origin_id LIKE 'dubizzle:%%'") ).scalar() or 0 if total_active == 0: return 0 # Временная таблица URL. session.execute(text("CREATE TEMP TABLE IF NOT EXISTS _active_urls (url TEXT NOT NULL) ON COMMIT DROP")) session.execute(text("TRUNCATE _active_urls")) # Вставляем URL чанками. url_list = list(normalized_urls) for i in range(0, len(url_list), _IN_CHUNK_SIZE): chunk = url_list[i:i + _IN_CHUNK_SIZE] values = ",".join(f"(:{f'u{j}'})" for j in range(len(chunk))) params = {f"u{j}": url for j, url in enumerate(chunk)} session.execute(text(f"INSERT INTO _active_urls (url) VALUES {values}"), params) # Индекс для JOIN. session.execute(text("CREATE INDEX IF NOT EXISTS _ix_active_urls ON _active_urls (url)")) # Считаем будущие sold. would_mark = session.execute(text(""" SELECT count(*) FROM cars c LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url WHERE a.url IS NULL AND c.is_sold = FALSE AND c.origin_id LIKE 'dubizzle:%%' """)).scalar() or 0 # Защита от аномалии. if total_active > 100 and would_mark > total_active * 0.8: logger.error( "mark_sold safety abort: would mark %d/%d (%.0f%%) as sold — likely URL format mismatch", would_mark, total_active, would_mark / total_active * 100, ) return 0 # Массовая пометка sold. result = session.execute(text(""" UPDATE cars SET is_sold = TRUE FROM ( SELECT c.id FROM cars c LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url WHERE a.url IS NULL AND c.is_sold = FALSE AND c.origin_id LIKE 'dubizzle:%%' ) sub WHERE cars.id = sub.id """)) count = result.rowcount or 0 else: # Упрощённый путь для SQLite. stmt = ( update(Car) .where(Car.is_sold == False) # noqa: E712 .where(Car.origin_id.like("dubizzle:%")) .values(is_sold=True) ) # Загружаем active URL. all_active = session.execute( select(Car.id, Car.origin_url).where( Car.is_sold == False, Car.origin_id.like("dubizzle:%") # noqa: E712 ) ).all() mark_ids = [row[0] for row in all_active if _norm(row[1]) not in normalized_urls] if not mark_ids: return 0 total_active = len(all_active) if total_active > 100 and len(mark_ids) > total_active * 0.8: logger.error( "mark_sold safety abort: would mark %d/%d (%.0f%%) as sold — likely URL format mismatch", len(mark_ids), total_active, len(mark_ids) / total_active * 100, ) return 0 for i in range(0, len(mark_ids), _IN_CHUNK_SIZE): chunk = mark_ids[i:i + _IN_CHUNK_SIZE] session.execute(update(Car).where(Car.id.in_(chunk)).values(is_sold=True)) count = len(mark_ids) if count: logger.info("Marked %d cars as sold by URL (no longer in listing)", count) return count def get_all_origin_ids_for_lane(self, prefix: str = "dubizzle:") -> set[str]: """Возвращает все известные origin_id для заданного префикса. Использует yield_per для потоковой загрузки при большом количестве записей. """ with self.session_scope() as session: result = session.execute( select(Car.origin_id).where(Car.origin_id.like(f"{prefix}%")).execution_options(yield_per=10000) ) return {str(row[0]) for row in result if row and row[0]} def get_all_active_origin_urls_for_lane(self, prefix: str = "dubizzle:") -> set[str]: """Возвращает origin_url всех активных (не проданных) авто для заданного lane-префикса.""" with self.session_scope() as session: result = session.execute( select(Car.origin_url).where( Car.origin_id.like(f"{prefix}%"), Car.is_sold == False, # noqa: E712 ).execution_options(yield_per=10000) ) return {str(row[0]) for row in result if row and row[0]} def count_active_cars_for_lane(self, prefix: str = "dubizzle:") -> int: """Возвращает количество активных (не проданных) авто для заданного lane-префикса.""" from sqlalchemy import func as sa_func with self.session_scope() as session: result = session.execute( select(sa_func.count()).select_from(Car).where( Car.origin_id.like(f"{prefix}%"), Car.is_sold == False, # noqa: E712 ) ) return int(result.scalar() or 0) def get_active_origin_urls_batch_for_refresh( self, prefix: str = "dubizzle:", offset: int = 0, limit: int = 500, ) -> list[str]: """Возвращает батч origin_url активных авто для rolling refresh. Сортировка по last_seen_at ASC — давно не обновлённые идут первыми. """ with self.session_scope() as session: result = session.execute( select(Car.origin_url).where( Car.origin_id.like(f"{prefix}%"), Car.is_sold == False, # noqa: E712 ).order_by(Car.last_seen_at.asc()).offset(offset).limit(limit) ) return [str(row[0]) for row in result if row and row[0]]