Schedule stale gallery enrichment
This commit is contained in:
@@ -12,7 +12,7 @@ from ..deps import get_persistence
|
||||
from ...storage.db import PersistenceService
|
||||
from ...storage.models import SyncRun
|
||||
from ...worker.celery_app import celery_app
|
||||
from ...worker.constants import MOBILEDE_SYNC_QUEUE
|
||||
from ...worker.constants import MOBILEDE_IMAGES_QUEUE, MOBILEDE_SYNC_QUEUE
|
||||
from ...worker.tasks import (
|
||||
mobilede_enrich_images_batch_task,
|
||||
mobilede_sync_detail_task,
|
||||
@@ -45,9 +45,10 @@ class MobileDeSyncDetailRequest(BaseModel):
|
||||
|
||||
|
||||
class MobileDeEnrichImagesRequest(BaseModel):
|
||||
limit: int = 50
|
||||
lane: str = "mobile_de_cars"
|
||||
batch_size: int = 10
|
||||
max_existing_images: int = 1
|
||||
staleness_hours: int = 168
|
||||
delay_seconds: float = 1.0
|
||||
|
||||
|
||||
@@ -89,12 +90,12 @@ def start_mobilede_sync_detail(body: MobileDeSyncDetailRequest):
|
||||
def start_mobilede_enrich_images(body: MobileDeEnrichImagesRequest):
|
||||
result = mobilede_enrich_images_batch_task.apply_async(
|
||||
kwargs=body.model_dump(),
|
||||
queue=MOBILEDE_SYNC_QUEUE,
|
||||
queue=MOBILEDE_IMAGES_QUEUE,
|
||||
)
|
||||
return {
|
||||
"task_id": result.id,
|
||||
"status": "queued",
|
||||
"queue": MOBILEDE_SYNC_QUEUE,
|
||||
"queue": MOBILEDE_IMAGES_QUEUE,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -345,9 +345,38 @@ class RuntimeMobileDeSegment:
|
||||
return " / ".join(parts) or self.make_id or "mobilede-segment"
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class RuntimeMobileDeEnrichmentConfig:
|
||||
enabled: bool = True
|
||||
batch_size: int = 10
|
||||
max_existing_images: int = 1
|
||||
staleness_hours: int = 168
|
||||
delay_seconds: float = 1.0
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeEnrichmentConfig":
|
||||
data = data or {}
|
||||
enabled = _optional_bool(data.get("enabled"))
|
||||
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"))
|
||||
try:
|
||||
delay_seconds = float(data.get("delay_seconds", 1.0))
|
||||
except (TypeError, ValueError):
|
||||
delay_seconds = 1.0
|
||||
return cls(
|
||||
enabled=True if enabled is None else enabled,
|
||||
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),
|
||||
delay_seconds=max(0.0, delay_seconds),
|
||||
)
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class RuntimeMobileDeConfig:
|
||||
segments: tuple[RuntimeMobileDeSegment, ...] = ()
|
||||
enrichment: RuntimeMobileDeEnrichmentConfig = field(default_factory=RuntimeMobileDeEnrichmentConfig)
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, data: dict[str, Any] | None) -> "RuntimeMobileDeConfig":
|
||||
@@ -361,7 +390,10 @@ class RuntimeMobileDeConfig:
|
||||
segment = RuntimeMobileDeSegment.from_dict(item)
|
||||
if segment is not None:
|
||||
segments.append(segment)
|
||||
return cls(segments=tuple(segments))
|
||||
return cls(
|
||||
segments=tuple(segments),
|
||||
enrichment=RuntimeMobileDeEnrichmentConfig.from_dict(data.get("enrichment")),
|
||||
)
|
||||
|
||||
def is_empty(self) -> bool:
|
||||
return not self.segments
|
||||
|
||||
@@ -5,7 +5,7 @@ from contextlib import contextmanager
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any, Iterator
|
||||
|
||||
from sqlalchemy import create_engine, delete, func, or_, select, text, update
|
||||
from sqlalchemy import case, create_engine, delete, func, or_, select, text, update
|
||||
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
||||
from sqlalchemy.orm import Session, sessionmaker
|
||||
|
||||
@@ -185,6 +185,8 @@ class PersistenceService:
|
||||
@staticmethod
|
||||
def _car_payload(record: CarRecord) -> dict[str, object]:
|
||||
payload = record.model_dump(mode="python")
|
||||
if record.details_confirmed:
|
||||
payload["details_fetched_at"] = datetime.now(timezone.utc)
|
||||
return {key: value for key, value in payload.items() if key in CAR_DB_FIELDS}
|
||||
|
||||
@staticmethod
|
||||
@@ -194,6 +196,8 @@ class PersistenceService:
|
||||
record: CarRecord,
|
||||
) -> dict[str, object]:
|
||||
prepared = dict(payload)
|
||||
if not record.details_confirmed:
|
||||
prepared["details_fetched_at"] = car.details_fetched_at
|
||||
if record.preserve_existing_details and not record.details_confirmed:
|
||||
for key, fallback_values in DETAIL_VALUE_FALLBACKS.items():
|
||||
if prepared.get(key) in fallback_values:
|
||||
@@ -907,24 +911,36 @@ class PersistenceService:
|
||||
self,
|
||||
prefix: str | tuple[str, ...] = MOBILEDE_ORIGIN_PREFIXES,
|
||||
*,
|
||||
limit: int = 50,
|
||||
limit: int | None = None,
|
||||
max_existing_images: int = 1,
|
||||
stale_before: datetime | None = None,
|
||||
) -> list[tuple[int, str, str, int]]:
|
||||
if limit is not None and int(limit) <= 0:
|
||||
return []
|
||||
prefixes = (prefix,) if isinstance(prefix, str) else tuple(prefix)
|
||||
with self.session_scope() as session:
|
||||
image_count = func.count(Image.id)
|
||||
result = session.execute(
|
||||
conditions = [
|
||||
_origin_prefix_filter(Car.origin_id, prefixes),
|
||||
Car.is_sold == False, # noqa: E712
|
||||
]
|
||||
if stale_before is not None:
|
||||
conditions.append(or_(Car.details_fetched_at.is_(None), Car.details_fetched_at <= stale_before))
|
||||
stmt = (
|
||||
select(Car.id, Car.origin_id, Car.origin_url, image_count.label("image_count"))
|
||||
.outerjoin(Image, Image.car_id == Car.id)
|
||||
.where(
|
||||
_origin_prefix_filter(Car.origin_id, prefixes),
|
||||
Car.is_sold == False, # noqa: E712
|
||||
)
|
||||
.group_by(Car.id, Car.origin_id, Car.origin_url)
|
||||
.where(*conditions)
|
||||
.group_by(Car.id, Car.origin_id, Car.origin_url, Car.details_fetched_at)
|
||||
.having(image_count <= max(0, int(max_existing_images)))
|
||||
.order_by(Car.last_seen_at.desc())
|
||||
.limit(max(1, int(limit)))
|
||||
.order_by(
|
||||
case((Car.details_fetched_at.is_(None), 0), else_=1).asc(),
|
||||
Car.details_fetched_at.asc(),
|
||||
Car.id.asc(),
|
||||
)
|
||||
)
|
||||
if limit is not None:
|
||||
stmt = stmt.limit(int(limit))
|
||||
result = session.execute(stmt)
|
||||
return [
|
||||
(int(row[0]), str(row[1]), str(row[2]), int(row[3] or 0))
|
||||
for row in result
|
||||
|
||||
@@ -56,6 +56,7 @@ class Car(Base):
|
||||
first_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
|
||||
last_seen_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now(), index=True)
|
||||
sold_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True)
|
||||
details_fetched_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True, index=True)
|
||||
images: Mapped[list["Image"]] = relationship("Image", back_populates="car", cascade="all, delete-orphan")
|
||||
|
||||
|
||||
|
||||
@@ -40,6 +40,7 @@ class CarRecord(BaseModel):
|
||||
first_seen_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
|
||||
last_seen_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
|
||||
sold_at: datetime | None = None
|
||||
details_fetched_at: datetime | None = None
|
||||
skip_image_sync: bool = False
|
||||
details_confirmed: bool = False
|
||||
images_confirmed: bool = False
|
||||
@@ -91,5 +92,6 @@ class CarRead(BaseModel):
|
||||
first_seen_at: datetime | None = None
|
||||
last_seen_at: datetime | None = None
|
||||
sold_at: datetime | None = None
|
||||
details_fetched_at: datetime | None = None
|
||||
images: list[ImageRead] = Field(default_factory=list)
|
||||
|
||||
|
||||
@@ -15,6 +15,7 @@ from .constants import (
|
||||
GLOBAL_DB_PROGRESS_TS_KEY,
|
||||
GLOBAL_PROGRESS_TS_KEY,
|
||||
MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||||
MOBILEDE_IMAGES_QUEUE,
|
||||
MOBILEDE_SYNC_QUEUE,
|
||||
MOBILEDE_SYNC_TASK_NAME,
|
||||
)
|
||||
@@ -35,6 +36,7 @@ def _env_bool(name: str, default: bool) -> bool:
|
||||
|
||||
|
||||
MOBILEDE_BEAT_SYNC_ENABLED = _env_bool("MOBILEDE_BEAT_SYNC_ENABLED", True)
|
||||
MOBILEDE_BEAT_ENRICH_ENABLED = _env_bool("MOBILEDE_BEAT_ENRICH_ENABLED", True)
|
||||
|
||||
|
||||
def _runtime_sync_kwargs() -> dict[str, bool | float]:
|
||||
@@ -49,6 +51,10 @@ def _runtime_sync_expires_seconds() -> float:
|
||||
return settings.celery.beat_sync_interval_minutes * 60.0
|
||||
|
||||
|
||||
def _runtime_enrich_interval_seconds() -> float:
|
||||
return max(60.0, float(os.getenv("MOBILEDE_ENRICH_INTERVAL_SECONDS", "900")))
|
||||
|
||||
|
||||
def _has_fresh_active_progress(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
@@ -155,17 +161,26 @@ if _hard > _max_hard:
|
||||
|
||||
beat_schedule = {}
|
||||
if MOBILEDE_BEAT_SYNC_ENABLED:
|
||||
beat_schedule = {
|
||||
"periodic-mobilede-sync-search": {
|
||||
"task": MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||||
"schedule": _runtime_sync_expires_seconds(),
|
||||
"args": (),
|
||||
"kwargs": _runtime_sync_kwargs(),
|
||||
"options": {
|
||||
"queue": MOBILEDE_SYNC_QUEUE,
|
||||
"expires": _runtime_sync_expires_seconds(),
|
||||
},
|
||||
}
|
||||
beat_schedule["periodic-mobilede-sync-search"] = {
|
||||
"task": MOBILEDE_RUNTIME_SEGMENTS_TASK,
|
||||
"schedule": _runtime_sync_expires_seconds(),
|
||||
"args": (),
|
||||
"kwargs": _runtime_sync_kwargs(),
|
||||
"options": {
|
||||
"queue": MOBILEDE_SYNC_QUEUE,
|
||||
"expires": _runtime_sync_expires_seconds(),
|
||||
},
|
||||
}
|
||||
if MOBILEDE_BEAT_ENRICH_ENABLED:
|
||||
beat_schedule["periodic-mobilede-enrich-images"] = {
|
||||
"task": MOBILEDE_ENRICH_IMAGES_TASK,
|
||||
"schedule": _runtime_enrich_interval_seconds(),
|
||||
"args": (),
|
||||
"kwargs": {},
|
||||
"options": {
|
||||
"queue": MOBILEDE_IMAGES_QUEUE,
|
||||
"expires": _runtime_enrich_interval_seconds(),
|
||||
},
|
||||
}
|
||||
|
||||
celery_app.conf.update(
|
||||
@@ -195,7 +210,7 @@ celery_app.conf.update(
|
||||
MOBILEDE_RUNTIME_SEGMENTS_TASK: {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
MOBILEDE_SYNC_TASK_NAME: {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
MOBILEDE_SYNC_DETAIL_TASK: {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
MOBILEDE_ENRICH_IMAGES_TASK: {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
MOBILEDE_ENRICH_IMAGES_TASK: {"queue": MOBILEDE_IMAGES_QUEUE},
|
||||
"mobilede_scraper.worker.tasks.*": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
},
|
||||
)
|
||||
|
||||
@@ -3,6 +3,8 @@ from __future__ import annotations
|
||||
import os
|
||||
|
||||
MOBILEDE_SYNC_QUEUE = "mobilede_sync"
|
||||
MOBILEDE_IMAGES_QUEUE = "mobilede_images"
|
||||
MOBILEDE_ENRICH_IMAGES_LOCK_KEY = "mobilede:locks:enrich_images"
|
||||
MOBILEDE_SEARCH_CURSOR_KEY = "mobilede:state:search_next_page"
|
||||
MOBILEDE_SEGMENT_CURSOR_KEY_FMT = "mobilede:state:search_next_page:{segment_key}"
|
||||
MOBILEDE_RUNTIME_SEGMENT_INDEX_KEY = "mobilede:state:runtime_segment_index"
|
||||
|
||||
@@ -2119,6 +2119,13 @@ def _expand_mobilede_search_url_segment(
|
||||
if query_keys & existing_range_keys:
|
||||
return [segment]
|
||||
|
||||
if not allow_adaptive_planning:
|
||||
logger.info(
|
||||
"mobile.de URL segment fast-start enabled: using search_url as-is for %s",
|
||||
_mobilede_short_segment_label(segment),
|
||||
)
|
||||
return [segment]
|
||||
|
||||
max_pages = int(segment.get("max_pages") or MOBILEDE_MAX_PAGE_NUMBER)
|
||||
root_segment = _mobilede_try_keep_root_segment_unsplit(
|
||||
segment,
|
||||
@@ -2128,13 +2135,6 @@ def _expand_mobilede_search_url_segment(
|
||||
if root_segment is not None:
|
||||
return root_segment
|
||||
|
||||
if not allow_adaptive_planning:
|
||||
logger.info(
|
||||
"mobile.de URL segment fast-start enabled: using search_url as-is for %s",
|
||||
_mobilede_short_segment_label(segment),
|
||||
)
|
||||
return [segment]
|
||||
|
||||
if allow_probe_planning:
|
||||
probe_planned = _expand_mobilede_search_url_segment_by_probe(segment)
|
||||
if probe_planned is not None:
|
||||
|
||||
@@ -7,7 +7,7 @@ import signal
|
||||
import hashlib
|
||||
import re
|
||||
import unicodedata
|
||||
from datetime import datetime, timezone
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from threading import Event, Thread
|
||||
import time
|
||||
import uuid
|
||||
@@ -2431,7 +2431,7 @@ def mobilede_sync_detail_task(self, listing_id: str, lane: str = "mobile_de_cars
|
||||
|
||||
@shared_task(
|
||||
name="mobilede.enrich_images_batch",
|
||||
queue=MOBILEDE_SYNC_QUEUE,
|
||||
queue=MOBILEDE_IMAGES_QUEUE,
|
||||
bind=True,
|
||||
max_retries=1,
|
||||
default_retry_delay=120,
|
||||
@@ -2439,53 +2439,100 @@ def mobilede_sync_detail_task(self, listing_id: str, lane: str = "mobile_de_cars
|
||||
)
|
||||
def mobilede_enrich_images_batch_task(
|
||||
self,
|
||||
limit: int = 50,
|
||||
lane: str = "mobile_de_cars",
|
||||
max_existing_images: int = 1,
|
||||
delay_seconds: float = 1.0,
|
||||
batch_size: int | None = None,
|
||||
max_existing_images: int | None = None,
|
||||
staleness_hours: int | None = None,
|
||||
delay_seconds: float | None = None,
|
||||
):
|
||||
runtime_config = RuntimeConfig.from_file(Settings().runtime_config_file)
|
||||
enrichment = runtime_config.mobilede.enrichment
|
||||
if not enrichment.enabled:
|
||||
return {
|
||||
"status": "disabled",
|
||||
"candidates": 0,
|
||||
"enriched": 0,
|
||||
"failed": 0,
|
||||
"skipped": 0,
|
||||
}
|
||||
|
||||
effective_batch_size = enrichment.batch_size if batch_size is None else max(1, int(batch_size))
|
||||
effective_max_images = (
|
||||
enrichment.max_existing_images
|
||||
if max_existing_images is None
|
||||
else max(0, int(max_existing_images))
|
||||
)
|
||||
effective_staleness_hours = (
|
||||
enrichment.staleness_hours
|
||||
if staleness_hours is None
|
||||
else max(0, int(staleness_hours))
|
||||
)
|
||||
effective_delay = enrichment.delay_seconds if delay_seconds is None else max(0.0, float(delay_seconds))
|
||||
|
||||
redis_client = _get_redis()
|
||||
owner_token = str(getattr(self.request, "id", None) or uuid.uuid4().hex)
|
||||
lock_ttl = max(300, int(float(os.getenv("MOBILEDE_ENRICH_LOCK_TTL_SECONDS", "3600"))))
|
||||
if not _acquire_lock(redis_client, MOBILEDE_ENRICH_IMAGES_LOCK_KEY, owner_token, lock_ttl):
|
||||
return {
|
||||
"status": "locked",
|
||||
"candidates": 0,
|
||||
"enriched": 0,
|
||||
"failed": 0,
|
||||
"skipped": 0,
|
||||
}
|
||||
|
||||
persistence = _get_persistence()
|
||||
scraper = MobileDeScraper(persistence=persistence)
|
||||
candidates = persistence.get_active_cars_batch_for_image_enrich(
|
||||
limit=limit,
|
||||
max_existing_images=max_existing_images,
|
||||
)
|
||||
enriched = 0
|
||||
failed = 0
|
||||
skipped = 0
|
||||
for _car_id, origin_id, origin_url, image_count in candidates:
|
||||
listing_id = _mobilede_extract_listing_id(origin_url) or str(origin_id).rsplit(":", 1)[-1]
|
||||
if not listing_id:
|
||||
skipped += 1
|
||||
continue
|
||||
try:
|
||||
result = scraper.sync_detail(str(listing_id), lane=lane)
|
||||
enriched += 1
|
||||
logger.info(
|
||||
"mobile.de image enrich completed: listing_id=%s origin_id=%s old_images=%s result=%s",
|
||||
listing_id,
|
||||
origin_id,
|
||||
image_count,
|
||||
result.get("upsert", {}),
|
||||
)
|
||||
except Exception as exc:
|
||||
failed += 1
|
||||
logger.warning(
|
||||
"mobile.de image enrich failed: listing_id=%s origin_id=%s error=%s",
|
||||
listing_id,
|
||||
origin_id,
|
||||
exc,
|
||||
exc_info=True,
|
||||
)
|
||||
if delay_seconds:
|
||||
time.sleep(max(0.0, float(delay_seconds)))
|
||||
return {
|
||||
"status": "success",
|
||||
"candidates": len(candidates),
|
||||
"enriched": enriched,
|
||||
"failed": failed,
|
||||
"skipped": skipped,
|
||||
}
|
||||
try:
|
||||
stale_before = datetime.now(timezone.utc) - timedelta(hours=effective_staleness_hours)
|
||||
candidates = persistence.get_active_cars_batch_for_image_enrich(
|
||||
max_existing_images=effective_max_images,
|
||||
stale_before=stale_before,
|
||||
)
|
||||
enriched = 0
|
||||
failed = 0
|
||||
skipped = 0
|
||||
for batch_start in range(0, len(candidates), effective_batch_size):
|
||||
batch = candidates[batch_start:batch_start + effective_batch_size]
|
||||
for _car_id, origin_id, origin_url, image_count in batch:
|
||||
listing_id = _mobilede_extract_listing_id(origin_url) or str(origin_id).rsplit(":", 1)[-1]
|
||||
if not listing_id:
|
||||
skipped += 1
|
||||
continue
|
||||
try:
|
||||
result = scraper.sync_detail(str(listing_id), lane=lane)
|
||||
enriched += 1
|
||||
logger.info(
|
||||
"mobile.de image enrich completed: listing_id=%s origin_id=%s old_images=%s result=%s",
|
||||
listing_id,
|
||||
origin_id,
|
||||
image_count,
|
||||
result.get("upsert", {}),
|
||||
)
|
||||
except Exception as exc:
|
||||
failed += 1
|
||||
logger.warning(
|
||||
"mobile.de image enrich failed: listing_id=%s origin_id=%s error=%s",
|
||||
listing_id,
|
||||
origin_id,
|
||||
exc,
|
||||
exc_info=True,
|
||||
)
|
||||
if effective_delay:
|
||||
time.sleep(effective_delay)
|
||||
return {
|
||||
"status": "success",
|
||||
"candidates": len(candidates),
|
||||
"enriched": enriched,
|
||||
"failed": failed,
|
||||
"skipped": skipped,
|
||||
}
|
||||
finally:
|
||||
_release_lock_if_owner(
|
||||
redis_client,
|
||||
MOBILEDE_ENRICH_IMAGES_LOCK_KEY,
|
||||
owner_token,
|
||||
)
|
||||
|
||||
|
||||
@shared_task(
|
||||
|
||||
Reference in New Issue
Block a user