2 Commits
Author SHA1 Message Date
qananasikq d3edce29bd Обновить README.md 2026-04-08 22:38:26 +02:00
qananasikq 50d107cfba Initial project version 2026-04-08 23:36:19 +03:00
14 changed files with 106 additions and 137 deletions
+1
View File
@@ -61,3 +61,4 @@ CELERY_TASK_SOFT_TIME_LIMIT=600
CELERY_TASK_TIME_LIMIT=900
CELERY_WORKER_CONCURRENCY=1
CELERY_BEAT_SYNC_INTERVAL_MINUTES=60
CELERY_BEAT_SYNC_LIMIT=26
+38 -38
View File
@@ -1,17 +1,17 @@
# IAAI Scraper
Парсер аукционных автомобилей с [iaai.com](https://www.iaai.com). Ходит по листингу, собирает карточки машин, вытаскивает данные из DOM и перехваченных XHR-ответов, складывает всё в PostgreSQL. Работает через Playwright (headless Chromium), крутится в Docker.
Небольшой сервис для сбора автомобилей с [iaai.com](https://www.iaai.com).
Собирает ссылки из листинга, открывает карточки машин, достаёт данные из HTML и XHR, сохраняет всё в PostgreSQL.
## Что важно по сайту
## Как устроен сайт и его защита
У IAAI стоит **Imperva Incapsula** — внешний anti-bot и WAF.
IAAI не использует Cloudflare или Akamai, но у них своя защита:
- headless Chromium без прокси часто режется
- часть данных грузится через XHR, часть есть сразу в HTML
- есть cookie banner
- **Детект автоматизации** — проверяют `navigator.webdriver`, `chrome.runtime`, WebGL-рендерер, `hardwareConcurrency` и прочие browser fingerprint параметры. Если видят Playwright/Puppeteer — блокируют.
- **Rate limiting** — после нескольких быстрых запросов подряд начинают отдавать пустые страницы или редиректить. Нет явного 429, просто перестают отдавать данные.
- **Geo-блокировка** — часть контента доступна только с US/CA IP. С европейских адресов листинг может быть пустым.
- **Динамическая подгрузка** — карточки машин подгружаются через XHR (`/Search`, `/VehicleDetail`), часть данных приходит в JSON, часть рендерится на сервере. Нельзя просто дёрнуть HTML нужен полноценный браузер с JS.
- **Cookie consent** — при первом заходе показывают баннер, без принятия кук часть функционала не работает.
Поэтому в проекте используется Playwright, паузы между действиями и прокси.
## Запуск
@@ -28,7 +28,8 @@ Swagger-документация: `http://localhost:8000/docs`
## API
**Здоровье и статистика:**
- `GET /health` — статус сервиса и подключения к БД
- `GET /health` -
статус сервиса и подключения к БД
- `GET /api/v1/stats` — сколько машин/картинок в базе, топ брендов
**Машины:**
@@ -43,7 +44,7 @@ Swagger-документация: `http://localhost:8000/docs`
- `GET /api/v1/tasks` — все задачи
- `GET /api/v1/sync-runs` — история запусков
Пример — скрапнуть конкретную машину:
Пример:
```bash
curl -X POST http://localhost:8000/api/v1/tasks/sync-vehicle
-H "Content-Type: application/json" \
@@ -52,7 +53,7 @@ curl -X POST http://localhost:8000/api/v1/tasks/sync-vehicle
## CLI
Для отладки без API/Celery:
Для локального запуска без API и Celery:
```bash
python main.py init-db
python main.py collect-listing --make Toyota
@@ -64,51 +65,50 @@ python main.py sync-listing --limit 10
```
iaai_scraper/
scraper.py — оркестратор: связывает browser → parser → storage
cli.py — CLI-команды (init-db, sync-vehicle, sync-listing, ...)
proxy_bridge.py — HTTP→SOCKS5 мост (Chromium не умеет SOCKS5 с авторизацией)
scraper.py — основной код скрапера
cli.py — CLI-команды
proxy_bridge.py — HTTP→SOCKS5 мост
browser/
factory.py — создание браузера, stealth-инъекции, fingerprint
listing.py — сбор ссылок на машины с листинга, пагинация
network.py — перехват XHR/fetch ответов через Playwright events
pace.py — рандомные паузы и движение мыши
factory.py — запуск браузера и контекста
listing.py — сбор ссылок из листинга
network.py — перехват XHR/fetch
pace.py — паузы между действиями
parsing/
parser.py — VehicleParser: 3 канала (DOM + XHR JSON + embedded JSON)
mapper.py — CarMapper: нормализация → CarRecord, content hash
parser.py — разбор карточки машины
mapper.py — нормализация данных
storage/
models.py — SQLAlchemy: Car (30+ полей), Image, SyncRun, ScrapeTask
schemas.py — Pydantic: CarRecord, ImageRecord, CarRead
enums.py — допустимые значения (drive, gearbox, body_type, ...)
db.py — PersistenceService: upsert (insert/update/skip), sync runs
models.py — SQLAlchemy модели
schemas.py — Pydantic схемы
enums.py — справочники значений
db.py — работа с БД
api/
app.py — FastAPI factory, lifespan, роутеры
deps.py — dependency injection (Settings, PersistenceService)
app.py — FastAPI приложение
deps.py — зависимости
routes/
health.py — GET /health
cars.py — CRUD по машинам + GET /stats
tasks.py — управление Celery-задачами + sync-runs
cars.py — машины и статистика
tasks.py — задачи и история запусков
worker/
celery_app.py — конфиг Celery, beat-расписание
tasks.py — sync_vehicle_task, sync_listing_task
celery_app.py — конфиг Celery
tasks.py — фоновые задачи
core/
config.py — Settings (dataclass), все env-переменные
logs.py — логирование с trace_id (ContextVar)
retry.py — декоратор @retryable с exponential backoff
utils.py — VIN_RE, deep_find_key, вспомогательные функции
config.py — настройки
logs.py — логирование
retry.py — retry-логика
utils.py — утилиты
alembic/ — миграции БД
tests/ — 30 тестов (SQLite in-memory)
```
## Конфигурация
Всё через env-переменные (полный список в `.env.example`):
Основные переменные лежат в `.env.example`:
**БД:** `IAAI_DATABASE_URL`, `IAAI_DATABASE_POOL_SIZE`
**Redis:** `IAAI_REDIS_URL`
@@ -119,7 +119,7 @@ tests/ — 30 тестов (SQLite in-memory)
## Миграции
В Docker миграции накатываются автоматически при старте API-контейнера (`alembic upgrade head` в entrypoint).
В Docker миграции запускаются автоматически при старте API-контейнера.
Вручную:
```bash
@@ -133,4 +133,4 @@ alembic revision --autogenerate -m "add_column_x"
pytest -q
```
30 тестов, SQLite in-memory, без внешних зависимостей.
Тесты идут на SQLite in-memory, без внешних сервисов.
+1 -4
View File
@@ -1,5 +1,3 @@
"""FastAPI app."""
from contextlib import asynccontextmanager
from fastapi import FastAPI
@@ -11,7 +9,6 @@ from .routes import cars, health, tasks
@asynccontextmanager
async def lifespan(app: FastAPI):
"""Жизненный цикл."""
persistence: PersistenceService = app.state.persistence
persistence.create_tables()
yield
@@ -37,5 +34,5 @@ def create_app(settings: Settings | None = None) -> FastAPI:
return app
# Для запуска через uvicorn.
# uvicorn entrypoint
app = create_app()
-2
View File
@@ -1,5 +1,3 @@
"""FastAPI deps."""
from fastapi import Request
from ..storage.db import PersistenceService
+1 -8
View File
@@ -1,5 +1,3 @@
"""Cars endpoints."""
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy import func, select
@@ -21,7 +19,6 @@ def list_cars(
is_sold: bool | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
"""Список автомобилей с пагинацией и фильтрами."""
with persistence.session_scope() as session:
query = select(Car)
@@ -40,7 +37,7 @@ def list_cars(
count_query = select(func.count()).select_from(query.subquery())
total = session.execute(count_query).scalar() or 0
# Пагинация.
# пагинация
offset = (page - 1) * per_page
cars = session.execute(
query.order_by(Car.last_seen_at.desc()).offset(offset).limit(per_page)
@@ -60,7 +57,6 @@ def get_car(
car_id: int,
persistence: PersistenceService = Depends(get_persistence),
):
"""Детальная информация об автомобиле с изображениями."""
with persistence.session_scope() as session:
car = session.get(Car, car_id)
if car is None:
@@ -73,7 +69,6 @@ def get_car_by_origin(
origin_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
"""Поиск автомобиля по origin_id."""
with persistence.session_scope() as session:
car = session.execute(
select(Car).where(Car.origin_id == origin_id)
@@ -85,7 +80,6 @@ def get_car_by_origin(
@router.get("/stats")
def get_stats(persistence: PersistenceService = Depends(get_persistence)):
"""Общая статистика по БД."""
with persistence.session_scope() as session:
total_cars = session.execute(select(func.count(Car.id))).scalar() or 0
total_images = session.execute(select(func.count(Image.id))).scalar() or 0
@@ -104,7 +98,6 @@ def get_stats(persistence: PersistenceService = Depends(get_persistence)):
def _car_to_dict(car: Car, include_images: bool = False) -> dict:
"""Сериализация Car в dict."""
result = {
"id": car.id,
"parser_id": car.parser_id,
-3
View File
@@ -1,5 +1,3 @@
"""Health endpoint."""
from fastapi import APIRouter, Depends
from sqlalchemy import text
@@ -11,7 +9,6 @@ router = APIRouter()
@router.get("/health")
def health_check(persistence: PersistenceService = Depends(get_persistence)):
"""Проверка API и БД."""
db_ok = False
try:
with persistence.session_scope() as session:
+17 -17
View File
@@ -1,5 +1,3 @@
"""Task endpoints."""
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException, Query
@@ -14,7 +12,7 @@ from ...worker.tasks import sync_vehicle_task, sync_listing_task
router = APIRouter()
# --- Схемы.
# --- схемы
class SyncVehicleRequest(BaseModel):
vehicle_url: str
@@ -29,17 +27,20 @@ class SyncListingRequest(BaseModel):
only_new: bool | None = None
# --- Эндпоинты.
# --- эндпоинты
@router.post("/tasks/sync-vehicle")
def start_sync_vehicle(
body: SyncVehicleRequest,
persistence: PersistenceService = Depends(get_persistence),
):
"""Запустить задачу скрапинга одного автомобиля через Celery."""
result = sync_vehicle_task.delay(body.vehicle_url, lane=body.lane)
result = sync_vehicle_task.apply_async(
args=(body.vehicle_url,),
kwargs={"lane": body.lane},
queue="scraping",
)
# Регистрируем задачу в БД.
# регистрируем в бд
persistence.create_scrape_task(
celery_task_id=result.id,
task_type="sync_vehicle",
@@ -58,13 +59,15 @@ def start_sync_listing(
body: SyncListingRequest,
persistence: PersistenceService = Depends(get_persistence),
):
"""Запустить задачу полного цикла листинга через Celery."""
result = sync_listing_task.delay(
make=body.make,
model=body.model,
lane=body.lane,
limit=body.limit,
only_new=body.only_new,
result = sync_listing_task.apply_async(
kwargs={
"make": body.make,
"model": body.model,
"lane": body.lane,
"limit": body.limit,
"only_new": body.only_new,
},
queue="scraping",
)
persistence.create_scrape_task(
@@ -83,7 +86,6 @@ def get_task_status(
task_id: str,
persistence: PersistenceService = Depends(get_persistence),
):
"""Статус Celery-задачи."""
task_info = persistence.get_scrape_task(task_id)
if task_info is None:
raise HTTPException(status_code=404, detail="Task not found")
@@ -98,7 +100,6 @@ def list_tasks(
task_type: str | None = None,
persistence: PersistenceService = Depends(get_persistence),
):
"""Список задач с пагинацией."""
with persistence.session_scope() as session:
query = select(ScrapeTask)
if status:
@@ -143,7 +144,6 @@ def list_sync_runs(
per_page: int = Query(20, ge=1, le=100),
persistence: PersistenceService = Depends(get_persistence),
):
"""История запусков синхронизации."""
with persistence.session_scope() as session:
total = session.execute(select(func.count(SyncRun.id))).scalar() or 0
+10 -10
View File
@@ -16,7 +16,7 @@ VEHICLE_HREF_RE = re.compile(r"/VehicleDetail/\d+(?:~[A-Z]{2})?", re.IGNORECASE)
@dataclass(slots=True)
class ListingVehicleLink:
# Ссылка на карточку.
# Ссылка на карточку
href: str
title: str = ""
lot_number: str | None = None
@@ -24,7 +24,7 @@ class ListingVehicleLink:
@dataclass(slots=True)
class ListingPageResult:
# Результат страницы листинга.
# Результат страницы листинга
source_url: str
page_number: int
vehicle_links: list[ListingVehicleLink] = field(default_factory=list)
@@ -34,12 +34,12 @@ class ListingPageResult:
class ListingCollector:
def __init__(self, settings: Settings, pacer: HumanPacer) -> None:
# Сборщик ссылок листинга.
# Сборщик ссылок листинга
self.settings = settings
self.pacer = pacer
def open_cars_listing(self, page: Page) -> None:
# Открываем листинг и ждем ссылки.
# Открываем листинг и ждём ссылки
logger.info("Opening cars listing page: %s", self.settings.listing.cars_url)
page.goto(self.settings.listing.cars_url, wait_until="domcontentloaded")
try:
@@ -53,7 +53,7 @@ class ListingCollector:
self.pacer.after_listing_open()
def apply_filters(self, page: Page, make: str | None = None, model: str | None = None) -> dict[str, str | None]:
# Пробуем make/model фильтры.
# Пробуем make/model фильтры
applied = {"make": None, "model": None}
if make and self._try_fill_filter_input(page, ["input[placeholder*='Make']", "input[aria-label*='Make']"], make):
applied["make"] = make
@@ -64,7 +64,7 @@ class ListingCollector:
return applied
def collect_current_page(self, page: Page, page_number: int = 1) -> ListingPageResult:
# Собираем ссылки текущей страницы.
# Собираем ссылки текущей страницы
anchors = page.locator("a[href*='/VehicleDetail/']")
total = min(anchors.count(), self.settings.listing.page_link_limit)
links: list[ListingVehicleLink] = []
@@ -87,7 +87,7 @@ class ListingCollector:
return ListingPageResult(source_url=page.url, page_number=page_number, vehicle_links=links, pagination_available=next_page_detected, next_page_detected=next_page_detected)
def go_to_next_page(self, page: Page) -> bool:
# Переход на следующую страницу.
# Переход на следующую страницу
selectors = ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]
for selector in selectors:
locator = page.locator(selector).first
@@ -109,7 +109,7 @@ class ListingCollector:
return False
def collect_listing_links(self, page: Page, *, make: str | None = None, model: str | None = None) -> dict[str, Any]:
# Последовательный сбор с лимитами.
# Последовательный сбор с лимитами
self.open_cars_listing(page)
applied_filters = self.apply_filters(page, make=make, model=model)
pages: list[dict[str, object]] = []
@@ -130,7 +130,7 @@ class ListingCollector:
@staticmethod
def _try_fill_filter_input(page: Page, selectors: list[str], value: str) -> bool:
# Несколько селекторов для фильтра.
# Несколько селекторов для фильтра
for selector in selectors:
locator = page.locator(selector).first
if locator.count() == 0:
@@ -150,7 +150,7 @@ class ListingCollector:
@staticmethod
def _has_next_page(page: Page) -> bool:
# Проверка кнопки next.
# Проверка кнопки next
for selector in ["a[aria-label*='Next']", "button[aria-label*='Next']", "a.pagination-next", "button.pagination-next", "a:has-text('Next')", "button:has-text('Next')"]:
if page.locator(selector).count() > 0:
return True
+6 -6
View File
@@ -23,7 +23,7 @@ class NetworkCapture:
_origin: str | None = None
def attach(self, page: Page, origin_url: str | None = None) -> None:
# Подписка на сетевые события.
# Подписка на сетевые события
try:
self._origin = urlparse(origin_url or page.url).netloc.lower() or None
except Exception:
@@ -32,14 +32,14 @@ class NetworkCapture:
page.on("response", self._on_response)
def _is_same_origin(self, url: str) -> bool:
# Фильтр по домену.
# Фильтр по домену
if not self.settings.capture.capture_same_origin_only or not self._origin:
return True
netloc = urlparse(url).netloc.lower()
return netloc == self._origin or netloc.endswith(".iaai.com")
def _on_request(self, request: Request) -> None:
# Берем только xhr/fetch.
# Берём только xhr/fetch
if request.resource_type not in {"xhr", "fetch"}:
return
if not self._is_same_origin(request.url):
@@ -47,7 +47,7 @@ class NetworkCapture:
if len(self.requests) >= self.settings.capture.max_requests:
return
key = f"{request.method}:{request.url}:{request.post_data or ''}"
# Убираем дубли.
# Убираем дубли
if key in self._seen_req:
return
self._seen_req.add(key)
@@ -57,7 +57,7 @@ class NetworkCapture:
})
def _on_response(self, response: Response) -> None:
# Сохраняем JSON.
# Сохраняем JSON
request = response.request
if request.resource_type not in {"xhr", "fetch"}:
return
@@ -95,7 +95,7 @@ class NetworkCapture:
@staticmethod
def _categorize(url: str) -> str:
# Простая категория URL.
# Простая категория URL
low = url.lower()
mapping = {
"images": ["image", "media", "photos", "gallery"],
+1 -11
View File
@@ -8,7 +8,6 @@ load_dotenv()
@dataclass(slots=True)
class FingerprintConfig:
# Браузерный отпечаток.
user_agent: str = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) "
@@ -32,7 +31,6 @@ class FingerprintConfig:
@dataclass(slots=True)
class CaptureConfig:
# Параметры capture.
capture_same_origin_only: bool = os.getenv("IAAI_CAPTURE_SAME_ORIGIN_ONLY", "true").strip().lower() in {"1", "true", "yes", "on"}
max_requests: int = int(os.getenv("IAAI_MAX_CAPTURED_REQUESTS", "40"))
max_json_responses: int = int(os.getenv("IAAI_MAX_CAPTURED_JSON_RESPONSES", "30"))
@@ -40,7 +38,6 @@ class CaptureConfig:
@dataclass(slots=True)
class HumanPaceConfig:
# Паузы между действиями.
enabled: bool = os.getenv("IAAI_HUMAN_PACE_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}
after_listing_open_min_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MIN_S", "2.5"))
after_listing_open_max_s: float = float(os.getenv("IAAI_AFTER_LISTING_OPEN_MAX_S", "4.5"))
@@ -58,7 +55,6 @@ class HumanPaceConfig:
@dataclass(slots=True)
class ListingConfig:
# Лимиты листинга.
cars_url: str = os.getenv("IAAI_CARS_LISTING_URL", "https://www.iaai.com/Vehiclelisting/Cars")
max_pages_per_run: int = int(os.getenv("IAAI_MAX_PAGES_PER_RUN", "5"))
max_vehicles_per_run: int = int(os.getenv("IAAI_MAX_VEHICLES_PER_RUN", "100"))
@@ -69,7 +65,6 @@ class ListingConfig:
@dataclass(slots=True)
class DatabaseConfig:
# Подключение к БД.
url: str = os.getenv("IAAI_DATABASE_URL", "postgresql+psycopg2://iaai:iaai@localhost:5432/iaai_scraper")
echo: bool = os.getenv("IAAI_DATABASE_ECHO", "false").strip().lower() in {"1", "true", "yes", "on"}
pool_size: int = int(os.getenv("IAAI_DATABASE_POOL_SIZE", "5"))
@@ -78,24 +73,22 @@ class DatabaseConfig:
@dataclass(slots=True)
class RedisConfig:
# Подключение к Redis.
url: str = os.getenv("IAAI_REDIS_URL", "redis://localhost:6379/0")
@dataclass(slots=True)
class CeleryConfig:
# Настройки Celery.
broker_url: str = os.getenv("CELERY_BROKER_URL", "") or ""
result_backend: str = os.getenv("CELERY_RESULT_BACKEND", "") or ""
task_soft_time_limit: int = int(os.getenv("CELERY_TASK_SOFT_TIME_LIMIT", "600"))
task_time_limit: int = int(os.getenv("CELERY_TASK_TIME_LIMIT", "900"))
worker_concurrency: int = int(os.getenv("CELERY_WORKER_CONCURRENCY", "1"))
beat_sync_interval_minutes: int = int(os.getenv("CELERY_BEAT_SYNC_INTERVAL_MINUTES", "60"))
beat_sync_limit: int = int(os.getenv("CELERY_BEAT_SYNC_LIMIT", "26"))
@dataclass(slots=True)
class ProxyConfig:
# Настройки прокси.
server: str | None = (os.getenv("IAAI_PROXY_SERVER") or "").strip() or None
username: str | None = (os.getenv("IAAI_PROXY_USERNAME") or "").strip() or None
password: str | None = (os.getenv("IAAI_PROXY_PASSWORD") or "").strip() or None
@@ -105,7 +98,6 @@ class ProxyConfig:
return bool(self.server)
def to_playwright_dict(self) -> dict[str, str] | None:
# Формат для Playwright.
if not self.server:
return None
result: dict[str, str] = {"server": self.server}
@@ -118,7 +110,6 @@ class ProxyConfig:
@dataclass(slots=True)
class Settings:
# Общие настройки.
home_url: str = "https://www.iaai.com/"
default_timeout_ms: int = int(os.getenv("IAAI_TIMEOUT_MS", "45000"))
network_settle_ms: int = int(os.getenv("IAAI_NETWORK_SETTLE_MS", "800"))
@@ -143,5 +134,4 @@ class Settings:
proxy: ProxyConfig = field(default_factory=ProxyConfig)
# Настройки по умолчанию.
settings = Settings()
+6 -1
View File
@@ -6,16 +6,19 @@ from pathlib import Path
from typing import Any, Iterable
# Файловые утилиты
def save_to_json(data: Any, filename: str | Path) -> None:
path = Path(filename)
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8")
# Небольшая случайная пауза между запросами
def short_sleep(a: float = 0.10, b: float = 0.35) -> None:
time.sleep(random.uniform(a, b))
# Маскирование персональных данных
def mask_email(email: str) -> str:
if "@" not in email:
return "***"
@@ -24,6 +27,7 @@ def mask_email(email: str) -> str:
return f"{safe_local}@{domain}"
# Возвращает первое непустое значение
def first_non_empty(values: Iterable[Any]) -> Any | None:
for value in values:
if value not in (None, "", [], {}, ()):
@@ -31,12 +35,13 @@ def first_non_empty(values: Iterable[Any]) -> Any | None:
return None
# Regex для VIN, lot, price.
# Регулярные выражения для VIN, lot и price
VIN_RE = re.compile(r"\b([A-HJ-NPR-Z0-9]{17})\b", re.IGNORECASE)
LOT_RE = re.compile(r"\b(\d{7,10})\b")
PRICE_RE = re.compile(r"\$\s?([\d,]+(?:\.\d{1,2})?)")
# Глубокий рекурсивный поиск значений по ключам
def deep_find_key(obj, target_keys: set[str], max_depth: int = 64, _depth: int = 0) -> list:
found = []
if _depth >= max_depth:
+22 -22
View File
@@ -33,7 +33,7 @@ class IAAIScraper:
return match.group(1)
def __init__(self, runtime_settings: Settings | None = None) -> None:
# Базовые зависимости и сервисы скрапера.
# Базовые зависимости и сервисы скрапера
self.settings = runtime_settings or settings
setup_logging(self.settings.log_level, self.settings.log_file)
self.trace_id = uuid.uuid4().hex[:12]
@@ -49,10 +49,10 @@ class IAAIScraper:
self.persistence = PersistenceService(self.settings)
self._shutdown_requested = False
# browser lifecycle
# Жизненный цикл браузера
def _new_trace_id(self, prefix: str) -> str:
# Новый trace id для отдельной операции.
# Новый trace id для отдельной операции
trace_id = f"{prefix}-{uuid.uuid4().hex[:8]}"
self.trace_id = trace_id
set_trace_id(trace_id)
@@ -69,7 +69,7 @@ class IAAIScraper:
self.close()
def close(self) -> None:
# Закрываем ресурсы аккуратно, даже если Playwright уже в ошибке.
# Закрываем ресурсы аккуратно даже при ошибках Playwright
if self.context is not None:
try:
self.context.close()
@@ -93,7 +93,7 @@ class IAAIScraper:
self.playwright = None
def _new_context(self, storage_state: str | None = None) -> BrowserContext:
# Контекст пересоздаётся, чтобы не тянуть старое состояние страницы.
# Контекст пересоздаётся, чтобы не тянуть старое состояние страницы
if self.browser is None:
self.__enter__()
if self.context:
@@ -109,10 +109,10 @@ class IAAIScraper:
context = self._new_context()
return context.new_page()
# listing
# Сбор листинга
def collect_listing(self, make: str | None = None, model: str | None = None):
# Собираем ссылки карточек с публичного листинга.
# Собираем ссылки карточек с публичного листинга
page = self._get_unauthenticated_page()
try:
@@ -130,11 +130,11 @@ class IAAIScraper:
self._new_context()
return self.context.new_page()
# scrape
# Открытие и скрейп карточки
@retryable(max_attempts=3, jitter_seconds=0.25)
def open_vehicle_page(self, vehicle_url: str):
# Лёгкое открытие страницы без сетевого дампа и сохранения в БД.
# Лёгкое открытие страницы без сетевого дампа и сохранения в БД
trace_id = self._new_trace_id("open")
started_at = time.perf_counter()
page = self._get_page()
@@ -174,7 +174,7 @@ class IAAIScraper:
page.close()
def _scrape_on_page(self, page: Page, vehicle_url: str):
# Полный проход по карточке с network capture и маппингом в DB-модель.
# Полный проход по карточке с network capture и маппингом в DB-модель
trace_id = self._new_trace_id("scrape")
started_at = time.perf_counter()
capture = NetworkCapture(self.settings)
@@ -215,11 +215,11 @@ class IAAIScraper:
save_to_json(network_dump, self.settings.raw_output_json)
return result
# db sync
# Синхронизация с БД
def sync_vehicle(self, vehicle_url: str, lane: str = "iaai"):
"""Scrape + upsert одного авто."""
# Один URL -> один scrape -> один upsert.
# Один URL -> один scrape -> один upsert
trace_id = self._new_trace_id("sync-vehicle")
started_at = time.perf_counter()
self.persistence.create_tables()
@@ -275,7 +275,7 @@ class IAAIScraper:
only_new: bool | None = None,
):
"""Листинг + sync всех найденных машин."""
# Массовая синхронизация с общим run_id и сбором ошибок.
# Массовая синхронизация с общим run_id и сбором ошибок
trace_id = self._new_trace_id("sync-listing")
started_at = time.perf_counter()
self.persistence.create_tables()
@@ -351,7 +351,7 @@ class IAAIScraper:
if index < total:
self.pacer.between_vehicles()
except Exception as exc:
# collect_listing itself failed — count as total failure
# Если падает collect_listing, считаем это общей ошибкой запуска
if not failures:
failures.append({"vehicle_url": "collect_listing", "error": str(exc)})
logger.error("sync_listing failed: %s", exc)
@@ -385,11 +385,11 @@ class IAAIScraper:
"failures": failures,
}
# scheduler
# Встроенный планировщик
def run_scheduled(self) -> None:
# Простой бесконечный цикл без внешнего планировщика.
# Graceful shutdown по SIGINT/SIGTERM.
# Простой бесконечный цикл без внешнего планировщика
# Graceful shutdown по SIGINT/SIGTERM
def _handle_shutdown(signum, frame):
logger.info("Received signal %s, shutting down gracefully...", signum)
self._shutdown_requested = True
@@ -408,7 +408,7 @@ class IAAIScraper:
logger.info("=== Scheduler cycle #%d starting ===", cycle)
start = time.time()
try:
# Сбрасываем контекст перед каждым циклом.
# Сбрасываем контекст перед каждым циклом
if self.context:
try:
self.context.close()
@@ -416,7 +416,7 @@ class IAAIScraper:
pass
self.context = None
# Проверяем, что browser/playwright живы; пересоздаём при необходимости.
# Проверяем, что browser/playwright живы, и пересоздаём при необходимости
if self.browser is None or self.playwright is None:
logger.info("Browser/Playwright not available, re-initializing...")
self.close()
@@ -433,12 +433,12 @@ class IAAIScraper:
except Exception as exc:
elapsed = time.time() - start
logger.error("Cycle #%d failed after %.1fs: %s", cycle, elapsed, exc)
# Полный сброс при любой ошибке цикла — следующий цикл пересоздаст всё.
# Полный сброс при любой ошибке цикла, следующий цикл пересоздаст всё
try:
self.close()
except Exception:
pass
# Пересоздаём browser для следующего цикла.
# Пересоздаём browser для следующего цикла
try:
self.__enter__()
except Exception as reinit_exc:
@@ -450,7 +450,7 @@ class IAAIScraper:
sleep_time = max(0, interval - (time.time() - start))
if sleep_time > 0:
logger.info("Sleeping %.0f seconds until next cycle...", sleep_time)
# Прерываемый sleep — проверяем shutdown каждые 5 секунд.
# Прерываемый sleep, проверяем shutdown каждые 5 секунд
slept = 0.0
while slept < sleep_time and not self._shutdown_requested:
chunk = min(5.0, sleep_time - slept)
+2 -10
View File
@@ -1,5 +1,3 @@
"""Celery app factory."""
from celery import Celery
from ..core.config import settings
@@ -20,40 +18,34 @@ celery_app = Celery(
)
celery_app.conf.update(
# Сериализация.
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
# Лимиты.
task_soft_time_limit=settings.celery.task_soft_time_limit,
task_time_limit=settings.celery.task_time_limit,
worker_concurrency=settings.celery.worker_concurrency,
# Playwright sync API: prefork + 1 процесс.
# prefork + 1 процесс, тк playwright sync api
worker_pool="prefork",
worker_prefetch_multiplier=1,
# Retry policy брокера.
broker_connection_retry_on_startup=True,
# Beat schedule.
beat_schedule={
"periodic-sync-listing": {
"task": "iaai_scraper.worker.tasks.sync_listing_task",
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
"args": (),
"kwargs": {"limit": settings.celery.beat_sync_limit},
"options": {"queue": "scraping"},
},
},
# Роутинг задач.
task_routes={
"iaai_scraper.worker.tasks.*": {"queue": "scraping"},
},
)
# Автообнаружение задач.
celery_app.autodiscover_tasks(["iaai_scraper.worker"])
+1 -5
View File
@@ -1,5 +1,3 @@
"""Celery tasks для IAAI."""
import json
import logging
from datetime import datetime, timezone
@@ -25,12 +23,11 @@ def _get_persistence() -> PersistenceService:
acks_late=True,
)
def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
"""Скрапинг и upsert одного автомобиля."""
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()
# Регистрируем запуск.
# регистрируем запуск
persistence.update_scrape_task(
task_id,
status="running",
@@ -86,7 +83,6 @@ def sync_listing_task(
limit: int | None = None,
only_new: bool | None = None,
):
"""Полный цикл: листинг + sync всех найденных машин."""
task_id = self.request.id
persistence = _get_persistence()
persistence.create_tables()