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