Clean runtime code
This commit is contained in:
@@ -0,0 +1,85 @@
|
|||||||
|
POSTGRES_USER=mobilede
|
||||||
|
POSTGRES_PASSWORD=mobilede
|
||||||
|
POSTGRES_DB=mobilede_scraper
|
||||||
|
MOBILEDE_HOST_API_PORT=8003
|
||||||
|
MOBILEDE_HOST_POSTGRES_PORT=5435
|
||||||
|
MOBILEDE_HOST_REDIS_PORT=6382
|
||||||
|
MOBILEDE_DATABASE_URL=postgresql+psycopg2://mobilede:mobilede@postgres:5432/mobilede_scraper
|
||||||
|
MOBILEDE_REDIS_URL=redis://redis:6379/0
|
||||||
|
CELERY_BROKER_URL=redis://redis:6379/0
|
||||||
|
CELERY_RESULT_BACKEND=redis://redis:6379/0
|
||||||
|
CELERY_TASK_SOFT_TIME_LIMIT=86400
|
||||||
|
CELERY_TASK_TIME_LIMIT=86520
|
||||||
|
CELERY_BROKER_VISIBILITY_TIMEOUT=90000
|
||||||
|
CELERY_WORKER_CONCURRENCY=1
|
||||||
|
CELERY_BEAT_SYNC_LIMIT=0
|
||||||
|
MOBILEDE_STARTUP_SYNC_ENABLED=true
|
||||||
|
MOBILEDE_STARTUP_MAX_PAGES=100
|
||||||
|
MOBILEDE_BEAT_MAX_PAGES=100
|
||||||
|
MOBILEDE_REQUEST_DELAY_SECONDS=0.25
|
||||||
|
MOBILEDE_CONTINUOUS_SYNC_ENABLED=true
|
||||||
|
MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS=3600
|
||||||
|
MOBILEDE_FULL_PASS_REPEAT_DELAY_SECONDS=3600
|
||||||
|
MOBILEDE_BOOTSTRAP_CONTINUATION_DELAY_SECONDS=0
|
||||||
|
MOBILEDE_BOOTSTRAP_FULL_SCAN_ENABLED=true
|
||||||
|
MOBILEDE_INCREMENTAL_AFTER_BOOTSTRAP=false
|
||||||
|
MOBILEDE_INCREMENTAL_PAGE_WINDOW=10
|
||||||
|
MOBILEDE_PROGRESS_LOG_EVERY_PAGES=20
|
||||||
|
MOBILEDE_RUNTIME_INITIAL_TASKS=2
|
||||||
|
MOBILEDE_SKIP_LATE_OVERFLOW_CHILDREN_DURING_BOOTSTRAP=true
|
||||||
|
MOBILEDE_CONCURRENT_PAGES=1
|
||||||
|
MOBILEDE_DETAIL_QUEUE_ENABLED=true
|
||||||
|
MOBILEDE_DETAIL_QUEUE_MAX_PENDING=10000
|
||||||
|
MOBILEDE_DETAIL_LOCK_TTL_SECONDS=180
|
||||||
|
MOBILEDE_DETAIL_WORKER_CONCURRENCY=10
|
||||||
|
MOBILEDE_DETAIL_REQUEST_DELAY_SECONDS=0.08
|
||||||
|
MOBILEDE_DETAIL_RATE_LIMIT_SLOTS=1
|
||||||
|
MOBILEDE_DETAIL_MAX_TASKS_PER_CHILD=500
|
||||||
|
MOBILEDE_HTTP_POOL_SIZE=32
|
||||||
|
MOBILEDE_SELECTIVE_DETAIL_ENRICH_ENABLED=false
|
||||||
|
MOBILEDE_DETAIL_ENRICH_IMAGES_ENABLED=false
|
||||||
|
MOBILEDE_DETAIL_TIMEOUT_SECONDS=10
|
||||||
|
MOBILEDE_HTTP_MAX_RETRIES=2
|
||||||
|
MOBILEDE_SEARCH_NETWORK_RETRY_DELAY_SECONDS=30
|
||||||
|
MOBILEDE_HUMAN_PACE_ENABLED=false
|
||||||
|
MOBILEDE_COMPACT_SEGMENTS=true
|
||||||
|
MOBILEDE_PREPLAN_SEGMENT_PROBES=true
|
||||||
|
MOBILEDE_PREPLAN_MAX_PROBES=40
|
||||||
|
MOBILEDE_PREPLAN_MAX_SEGMENTS=2400
|
||||||
|
MOBILEDE_PREPLAN_SPLIT_THRESHOLD_RATIO=1.0
|
||||||
|
MOBILEDE_PREPLAN_MAX_SPLIT_DEPTH=2
|
||||||
|
MOBILEDE_PLAN_PROBE_TIMEOUT_SECONDS=8
|
||||||
|
MOBILEDE_OVERFLOW_MAX_CHILD_SEGMENTS=3
|
||||||
|
MOBILEDE_SEGMENT_TARGET_RESULTS=950
|
||||||
|
MOBILEDE_SEGMENT_TARGET_MIN_RATIO=0.80
|
||||||
|
MOBILEDE_SEGMENT_TARGET_MAX_RATIO=0.98
|
||||||
|
MOBILEDE_SEGMENT_TINY_RATIO=0.45
|
||||||
|
MOBILEDE_SPLIT_SEGMENTS_BY_MILEAGE=true
|
||||||
|
MOBILEDE_RUNTIME_BRAND_SEGMENTS_ENABLED=false
|
||||||
|
MOBILEDE_RUNTIME_BRAND_ROOT_PROBES_ENABLED=false
|
||||||
|
MOBILEDE_FILTERED_SEARCH_URL=https://www.mobile.de/ru/%D1%82%D1%80%D0%B0%D0%BD%D1%81%D0%BF%D0%BE%D1%80%D1%82%D0%BD%D1%8B%D0%B5-%D1%81%D1%80%D0%B5%D0%B4%D1%81%D1%82%D0%B2%D0%B0/%D0%BF%D0%BE%D0%B8%D1%81%D0%BA.html?isSearchRequest=true&s=Car&vc=Car&ms=3500&ms=20100&ms=15200&ms=14400&od=up&sb=rel&ref=dsp
|
||||||
|
MOBILEDE_FILTERED_SEARCH_URLS=
|
||||||
|
MOBILEDE_COMPLETE_SCOPE_BRANDS=BMW,Volvo,Ferrari,Lexus
|
||||||
|
MOBILEDE_DETAIL_REFRESH_STALE_ENABLED=false
|
||||||
|
MOBILEDE_ENRICH_INTERVAL_SECONDS=30
|
||||||
|
MOBILEDE_PROXY_SERVER=
|
||||||
|
MOBILEDE_PROXY_USERNAME=
|
||||||
|
MOBILEDE_PROXY_PASSWORD=
|
||||||
|
SOCKS5_PROXY_HOST=
|
||||||
|
SOCKS5_PROXY_PORT=
|
||||||
|
SOCKS5_PROXY_USER=
|
||||||
|
SOCKS5_PROXY_PASS=
|
||||||
|
HTTP_PROXY=
|
||||||
|
HTTPS_PROXY=
|
||||||
|
ALL_PROXY=
|
||||||
|
NO_PROXY=
|
||||||
|
MOBILEDE_RUNTIME_CONFIG_FILE=/app/runtime_config.json
|
||||||
|
MOBILEDE_OVERFLOW_MIN_PRICE_SPLIT_SPAN=500
|
||||||
|
MOBILEDE_ROTATE_RUNTIME_SEGMENTS=true
|
||||||
|
MOBILEDE_OVERFLOW_SPLIT_ENABLED=true
|
||||||
|
TZ=UTC
|
||||||
|
|
||||||
|
|
||||||
|
MOBILEDE_LIGHT_REFRESH_EXISTING=true
|
||||||
|
MOBILEDE_SKIP_IMAGES_FOR_UPDATED=1
|
||||||
|
MOBILEDE_BEAT_SYNC_ENABLED=false
|
||||||
@@ -1,35 +1,86 @@
|
|||||||
# mobile.de Scraper
|
# mobile.de Scraper
|
||||||
|
|
||||||
Парсер [`mobile.de`](https://www.mobile.de/ru).
|
Сервис собирает объявления из mobile.de в PostgreSQL. Поиск и загрузка галерей разделены: основной worker сохраняет данные из Search, а отдельный worker получает только фотографии.
|
||||||
|
|
||||||
Сервис для сбора объявлений авто с mobile.de, сохранения в PostgreSQL и регулярного обновления.
|
## Как работает
|
||||||
|
|
||||||
## Как идет парсинг
|
1. `worker` получает исходную ссылку из `MOBILEDE_FILTERED_SEARCH_URL`.
|
||||||
- Парсинг запускается задачами Celery в очереди `mobilede_sync`.
|
2. Если ссылка содержит несколько марок, она разделяется по маркам.
|
||||||
- Стартуем с поисковой ссылки mobile.de (URL с фильтрами: марка, цена, год, пробег и т.д.).
|
3. Каждая выдача делится по цене, году и пробегу до сегментов, которые не превышают лимит mobile.de: 50 страниц по 20 объявлений.
|
||||||
- Выдача разбивается на сегменты (по цене, году, пробегу), чтобы обходить лимиты страниц.
|
4. Данные карточек из Search сохраняются в `MOBILEDE_cars`; новые объявления ставятся в очередь галерей.
|
||||||
- По каждому сегменту читаются страницы поиска, из них достаются ID и ссылки объявлений.
|
5. `worker-images` читает `mobilede_images`, загружает галерею объявления и сохраняет изображения в `MOBILEDE_images`.
|
||||||
- По ID/ссылкам загружаются карточки авто, данные маппятся в единый формат и пишутся в БД через upsert.
|
6. Между полными проходами действует интервал `MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS`.
|
||||||
|
|
||||||
## Стек
|
План сегментов сохраняется в volume `tokens_data`. Неполный или плотный план не публикуется для обхода, чтобы не терять объявления за пределами 50-й страницы.
|
||||||
- `FastAPI` — API и служебные ручки
|
|
||||||
- `Celery + Redis` — очередь задач парсинга
|
## Сервисы
|
||||||
- `PostgreSQL` — хранение авто и изображений
|
|
||||||
- `Docker Compose` — запуск всех сервисов
|
| Сервис | Назначение |
|
||||||
|
| --- | --- |
|
||||||
|
| `postgres` | Автомобили, изображения, отметка `gallery_fetched_at` и история запусков |
|
||||||
|
| `redis` | Брокер Celery, временные Detail claims/retry, runtime-состояние и лимит запросов |
|
||||||
|
| `migrate` | Применяет Alembic-миграции перед запуском приложения |
|
||||||
|
| `worker` | Планирование сегментов и импорт Search-выдачи |
|
||||||
|
| `worker-images` | Загрузка галерей из очереди `mobilede_images` |
|
||||||
|
| `beat` | Периодические задачи Celery |
|
||||||
|
| `api` | FastAPI на порту `MOBILEDE_HOST_API_PORT` |
|
||||||
|
|
||||||
|
## Настройка
|
||||||
|
|
||||||
|
Конфигурация контейнеров находится в `.env`, правила фильтрации — в `runtime_config.json`.
|
||||||
|
Основные переменные:
|
||||||
|
|
||||||
|
- `MOBILEDE_FILTERED_SEARCH_URL` — исходная ссылка mobile.de с нужной областью поиска;
|
||||||
|
- `MOBILEDE_RESULTS_PER_PAGE=20` и `MOBILEDE_MAX_PAGE_NUMBER=50` — лимит одной выдачи;
|
||||||
|
- `MOBILEDE_SEGMENT_TARGET_RESULTS=950` — целевой размер сегмента;
|
||||||
|
- `MOBILEDE_DETAIL_WORKER_CONCURRENCY` — число процессов загрузки галерей;
|
||||||
|
- `MOBILEDE_DETAIL_REQUEST_DELAY_SECONDS` — общий интервал между запросами галерей;
|
||||||
|
- `MOBILEDE_CONTINUOUS_SYNC_DELAY_SECONDS` — пауза между полными проходами.
|
||||||
|
|
||||||
|
`worker-images` использует единый Redis-лимитер и circuit breaker: при `403` или `429` новые gallery-запросы временно откладываются, а не повторяются одновременно всеми процессами.
|
||||||
|
|
||||||
|
Detail-очередь не создаёт служебные строки в PostgreSQL. Redis хранит временный claim с TTL и счётчик ограниченных retry, а окончательным признаком успешно загруженной галереи служит `MOBILEDE_cars.gallery_fetched_at`. После потери Redis уже завершённые галереи не выбираются повторно.
|
||||||
|
|
||||||
|
## Запуск
|
||||||
|
|
||||||
|
Подготовьте локальный файл настроек:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
copy .env.example .env
|
||||||
|
```
|
||||||
|
|
||||||
|
Укажите в `.env` `MOBILEDE_FILTERED_SEARCH_URL`, затем запустите стек:
|
||||||
|
|
||||||
## Быстрый запуск
|
|
||||||
```bash
|
```bash
|
||||||
docker compose up -d --build
|
docker compose up -d --build
|
||||||
```
|
```
|
||||||
|
|
||||||
## Прокси
|
Проверить контейнеры и логи:
|
||||||
- В воркере используется прокси-мост: локальный HTTP `127.0.0.1:8899` -> внешний SOCKS5.
|
|
||||||
- Это помогает стабилизировать доступ к mobile.de и снизить блокировки.
|
|
||||||
- Параметры прокси задаются через переменные окружения (см. `.env` / `docker-compose.yml`).
|
|
||||||
|
|
||||||
## Основные сервисы
|
```bash
|
||||||
- `api` — HTTP API
|
docker compose ps
|
||||||
- `worker` — парсинг и апдейты
|
docker compose logs -f worker worker-images
|
||||||
- `beat` — планировщик задач
|
```
|
||||||
- `postgres` — база данных
|
|
||||||
- `redis` — брокер очереди
|
Миграции применяются контейнером `migrate` автоматически. Для полностью чистого тестового запуска удалите volumes PostgreSQL, Redis и `tokens_data` перед `docker compose up`.
|
||||||
|
|
||||||
|
## API
|
||||||
|
|
||||||
|
После запуска доступны:
|
||||||
|
|
||||||
|
- `GET /health` — состояние сервиса;
|
||||||
|
- `GET /api/v1/cars` и `GET /api/v1/cars/{car_id}` — автомобили;
|
||||||
|
- `GET /api/v1/stats` — агрегированная статистика;
|
||||||
|
- `POST /api/v1/mobilede/tasks/sync-runtime-segments` — запуск прохода сегментов;
|
||||||
|
- `POST /api/v1/mobilede/tasks/enrich-images` — постановка галерей в обработку;
|
||||||
|
- `GET /api/v1/sync-runs` — история запусков.
|
||||||
|
|
||||||
|
Интерактивная спецификация FastAPI доступна по `/docs`.
|
||||||
|
|
||||||
|
## Проверка
|
||||||
|
|
||||||
|
```bash
|
||||||
|
pytest
|
||||||
|
```
|
||||||
|
|
||||||
|
Для просмотра текущего плана и прогресса используйте логи `worker`, Redis-ключи `mobilede:state:bootstrap_segments_total` и `mobilede:state:bootstrap_segments_done`, а также файл `/data/mobilede_runtime_segments.json` внутри volume `tokens_data`.
|
||||||
|
|||||||
@@ -48,9 +48,7 @@ class MobileDeSyncDetailRequest(BaseModel):
|
|||||||
class MobileDeEnrichImagesRequest(BaseModel):
|
class MobileDeEnrichImagesRequest(BaseModel):
|
||||||
lane: str = "mobile_de_cars"
|
lane: str = "mobile_de_cars"
|
||||||
batch_size: int = 10
|
batch_size: int = 10
|
||||||
max_existing_images: int = 1
|
|
||||||
staleness_hours: int = 168
|
staleness_hours: int = 168
|
||||||
delay_seconds: float = 1.0
|
|
||||||
|
|
||||||
|
|
||||||
class MobileDeRuntimeSegmentsRequest(BaseModel):
|
class MobileDeRuntimeSegmentsRequest(BaseModel):
|
||||||
|
|||||||
@@ -155,7 +155,6 @@ class ProxyConfig:
|
|||||||
|
|
||||||
@dataclass(slots=True)
|
@dataclass(slots=True)
|
||||||
class Settings:
|
class Settings:
|
||||||
headless: bool = _env_bool("MOBILEDE_HEADLESS", True)
|
|
||||||
log_level: str = _env_str("MOBILEDE_LOG_LEVEL", "INFO")
|
log_level: str = _env_str("MOBILEDE_LOG_LEVEL", "INFO")
|
||||||
runtime_config_file: str | None = _env_path_str("MOBILEDE_RUNTIME_CONFIG_FILE")
|
runtime_config_file: str | None = _env_path_str("MOBILEDE_RUNTIME_CONFIG_FILE")
|
||||||
listing: ListingConfig = field(default_factory=ListingConfig)
|
listing: ListingConfig = field(default_factory=ListingConfig)
|
||||||
|
|||||||
@@ -1,20 +1,5 @@
|
|||||||
import logging
|
import logging
|
||||||
import sys
|
import sys
|
||||||
from contextvars import ContextVar
|
|
||||||
|
|
||||||
# Trace ID текущего потока или корутины.
|
|
||||||
TRACE_ID: ContextVar[str] = ContextVar("trace_id", default="-")
|
|
||||||
|
|
||||||
|
|
||||||
class TraceIdFilter(logging.Filter):
|
|
||||||
# Добавляет trace_id в запись лога.
|
|
||||||
def filter(self, record: logging.LogRecord) -> bool:
|
|
||||||
record.trace_id = TRACE_ID.get()
|
|
||||||
return True
|
|
||||||
|
|
||||||
|
|
||||||
def set_trace_id(trace_id: str) -> None:
|
|
||||||
TRACE_ID.set(trace_id)
|
|
||||||
|
|
||||||
|
|
||||||
def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
|
def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
|
||||||
@@ -22,9 +7,7 @@ def setup_logging(level: str = "INFO", log_file: str | None = None) -> None:
|
|||||||
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)]
|
handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)]
|
||||||
if log_file:
|
if log_file:
|
||||||
handlers.append(logging.FileHandler(log_file, encoding="utf-8"))
|
handlers.append(logging.FileHandler(log_file, encoding="utf-8"))
|
||||||
trace_filter = TraceIdFilter()
|
|
||||||
for handler in handlers:
|
for handler in handlers:
|
||||||
handler.addFilter(trace_filter)
|
|
||||||
handler.setLevel(getattr(logging, level.upper(), logging.INFO))
|
handler.setLevel(getattr(logging, level.upper(), logging.INFO))
|
||||||
root = logging.getLogger()
|
root = logging.getLogger()
|
||||||
root.setLevel(getattr(logging, level.upper(), logging.INFO))
|
root.setLevel(getattr(logging, level.upper(), logging.INFO))
|
||||||
|
|||||||
@@ -56,10 +56,6 @@ def _optional_bool(value: Any) -> bool | None:
|
|||||||
@dataclass(slots=True)
|
@dataclass(slots=True)
|
||||||
class RuntimeSyncConfig:
|
class RuntimeSyncConfig:
|
||||||
name: str | None = None
|
name: str | None = None
|
||||||
ids_initial_size: int | None = None
|
|
||||||
ids_next_size: int | None = None
|
|
||||||
ids_max_pages: int | None = None
|
|
||||||
condition_check_enabled: bool | None = None
|
|
||||||
lane: str | None = None
|
lane: str | None = None
|
||||||
only_new: bool | None = None
|
only_new: bool | None = None
|
||||||
limit: int | None = None
|
limit: int | None = None
|
||||||
@@ -71,10 +67,6 @@ class RuntimeSyncConfig:
|
|||||||
lane = str(data.get("lane")).strip() if data.get("lane") else None
|
lane = str(data.get("lane")).strip() if data.get("lane") else None
|
||||||
return cls(
|
return cls(
|
||||||
name=name or None,
|
name=name or None,
|
||||||
ids_initial_size=_optional_int(data.get("ids_initial_size")),
|
|
||||||
ids_next_size=_optional_int(data.get("ids_next_size")),
|
|
||||||
ids_max_pages=_optional_int(data.get("ids_max_pages")),
|
|
||||||
condition_check_enabled=_optional_bool(data.get("condition_check_enabled")),
|
|
||||||
lane=lane or None,
|
lane=lane or None,
|
||||||
only_new=_optional_bool(data.get("only_new")),
|
only_new=_optional_bool(data.get("only_new")),
|
||||||
limit=_optional_int(data.get("limit")),
|
limit=_optional_int(data.get("limit")),
|
||||||
@@ -349,27 +341,18 @@ class RuntimeMobileDeSegment:
|
|||||||
class RuntimeMobileDeEnrichmentConfig:
|
class RuntimeMobileDeEnrichmentConfig:
|
||||||
enabled: bool = True
|
enabled: bool = True
|
||||||
batch_size: int = 10
|
batch_size: int = 10
|
||||||
max_existing_images: int = 1
|
|
||||||
staleness_hours: int = 168
|
staleness_hours: int = 168
|
||||||
delay_seconds: float = 1.0
|
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeEnrichmentConfig":
|
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeEnrichmentConfig":
|
||||||
data = data or {}
|
data = data or {}
|
||||||
enabled = _optional_bool(data.get("enabled"))
|
enabled = _optional_bool(data.get("enabled"))
|
||||||
batch_size = _optional_int(data.get("batch_size"))
|
batch_size = _optional_int(data.get("batch_size"))
|
||||||
max_existing_images = _optional_int(data.get("max_existing_images"))
|
|
||||||
staleness_hours = _optional_int(data.get("staleness_hours"))
|
staleness_hours = _optional_int(data.get("staleness_hours"))
|
||||||
try:
|
|
||||||
delay_seconds = float(data.get("delay_seconds", 1.0))
|
|
||||||
except (TypeError, ValueError):
|
|
||||||
delay_seconds = 1.0
|
|
||||||
return cls(
|
return cls(
|
||||||
enabled=True if enabled is None else enabled,
|
enabled=True if enabled is None else enabled,
|
||||||
batch_size=max(1, 10 if batch_size is None else batch_size),
|
batch_size=max(1, 10 if batch_size is None else batch_size),
|
||||||
max_existing_images=max(0, 1 if max_existing_images is None else max_existing_images),
|
|
||||||
staleness_hours=max(0, 168 if staleness_hours is None else staleness_hours),
|
staleness_hours=max(0, 168 if staleness_hours is None else staleness_hours),
|
||||||
delay_seconds=max(0.0, delay_seconds),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user