Refactor to new ingestion pipeline architecture with discovery, fetch, enrichment services
This commit is contained in:
@@ -1,7 +1,10 @@
|
||||
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
|
||||
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import signal
|
||||
from threading import Event, Thread
|
||||
import time
|
||||
import uuid
|
||||
@@ -13,6 +16,7 @@ from redis import Redis
|
||||
from ..core.config import Settings, parse_listing_segments
|
||||
from ..scraper import IAAIScraper
|
||||
from ..storage.db import PersistenceService
|
||||
from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
|
||||
|
||||
logger = logging.getLogger("iaai_scraper.worker.tasks")
|
||||
|
||||
@@ -25,6 +29,13 @@ SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
|
||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
|
||||
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
||||
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
|
||||
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
|
||||
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
|
||||
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
|
||||
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
|
||||
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
|
||||
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
|
||||
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
|
||||
|
||||
|
||||
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||
@@ -91,6 +102,283 @@ def _sync_listing_lock_ttl_seconds() -> int:
|
||||
return max(effective_hard + 120, 300)
|
||||
|
||||
|
||||
def _sync_segment_lock_ttl_seconds() -> int:
|
||||
return 300
|
||||
|
||||
|
||||
def _task_progress_key(task_id: str) -> str:
|
||||
return TASK_PROGRESS_KEY_FMT.format(task_id=task_id)
|
||||
|
||||
|
||||
def _update_task_progress(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
task_id: str,
|
||||
stage: str,
|
||||
ttl_seconds: int,
|
||||
**payload,
|
||||
) -> None:
|
||||
try:
|
||||
data = {
|
||||
"task_id": task_id,
|
||||
"stage": stage,
|
||||
"ts": int(time.time()),
|
||||
**payload,
|
||||
}
|
||||
redis_client.set(
|
||||
_task_progress_key(task_id),
|
||||
json.dumps(data, ensure_ascii=False),
|
||||
ex=max(60, int(ttl_seconds)),
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to update task progress for %s", task_id, exc_info=True)
|
||||
|
||||
|
||||
def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
|
||||
try:
|
||||
redis_client.delete(_task_progress_key(task_id))
|
||||
except Exception:
|
||||
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
|
||||
|
||||
|
||||
def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None) -> dict[str, object]:
|
||||
del limit, only_new
|
||||
settings = Settings()
|
||||
persistence = _get_persistence()
|
||||
persistence.create_tables()
|
||||
|
||||
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
|
||||
discovered_urls = discovery_result.vehicle_urls
|
||||
active_urls = set(discovered_urls)
|
||||
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
|
||||
|
||||
new_urls = [url for url in discovered_urls if url not in existing_urls]
|
||||
sold_count = 0
|
||||
if active_urls:
|
||||
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
|
||||
|
||||
cars_upserted = 0
|
||||
cars_failed = 0
|
||||
images_upserted = 0
|
||||
failures: list[dict[str, str]] = []
|
||||
|
||||
if new_urls:
|
||||
with IAAIScraper(settings) as scraper:
|
||||
batch_size = settings.celery.batch_size
|
||||
for batch_start in range(0, len(new_urls), batch_size):
|
||||
batch_urls = new_urls[batch_start:batch_start + batch_size]
|
||||
batch_result = scraper.sync_batch(batch_urls, lane=lane)
|
||||
cars_upserted += int(batch_result.get("cars_upserted", 0))
|
||||
cars_failed += int(batch_result.get("cars_failed", 0))
|
||||
images_upserted += int(batch_result.get("images_upserted", 0))
|
||||
failures.extend(batch_result.get("failures", []))
|
||||
|
||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||
return {
|
||||
"status": status,
|
||||
"run_id": None,
|
||||
"cars_upserted": cars_upserted,
|
||||
"cars_failed": cars_failed,
|
||||
"images_upserted": images_upserted,
|
||||
"skipped_existing": len(discovered_urls) - len(new_urls),
|
||||
"elapsed_seconds": None,
|
||||
"failures": failures,
|
||||
"full_scan_completed": True,
|
||||
"hourly_mode": "sitemap_diff",
|
||||
"discovered_urls": len(discovered_urls),
|
||||
"new_urls": len(new_urls),
|
||||
"sold_marked": sold_count,
|
||||
"transport": discovery_result.stats.transport,
|
||||
}
|
||||
|
||||
|
||||
def _hourly_sitemap_full_refresh_sync(*, lane: str, limit: int | None, only_new: bool | None) -> dict[str, object]:
|
||||
del limit, only_new
|
||||
settings = Settings()
|
||||
persistence = _get_persistence()
|
||||
persistence.create_tables()
|
||||
|
||||
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
|
||||
discovered_urls = discovery_result.vehicle_urls
|
||||
active_urls = set(discovered_urls)
|
||||
sold_count = 0
|
||||
if active_urls:
|
||||
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
|
||||
|
||||
cars_upserted = 0
|
||||
cars_failed = 0
|
||||
images_upserted = 0
|
||||
failures: list[dict[str, str]] = []
|
||||
|
||||
with IAAIScraper(settings) as scraper:
|
||||
batch_size = settings.celery.batch_size
|
||||
for batch_start in range(0, len(discovered_urls), batch_size):
|
||||
batch_urls = discovered_urls[batch_start:batch_start + batch_size]
|
||||
batch_result = scraper.sync_batch(batch_urls, lane=lane)
|
||||
cars_upserted += int(batch_result.get("cars_upserted", 0))
|
||||
cars_failed += int(batch_result.get("cars_failed", 0))
|
||||
images_upserted += int(batch_result.get("images_upserted", 0))
|
||||
failures.extend(batch_result.get("failures", []))
|
||||
|
||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||
return {
|
||||
"status": status,
|
||||
"run_id": None,
|
||||
"cars_upserted": cars_upserted,
|
||||
"cars_failed": cars_failed,
|
||||
"images_upserted": images_upserted,
|
||||
"skipped_existing": 0,
|
||||
"elapsed_seconds": None,
|
||||
"failures": failures,
|
||||
"full_scan_completed": True,
|
||||
"hourly_mode": "sitemap_full_refresh",
|
||||
"discovered_urls": len(discovered_urls),
|
||||
"new_urls": None,
|
||||
"sold_marked": sold_count,
|
||||
"transport": discovery_result.stats.transport,
|
||||
}
|
||||
|
||||
|
||||
def _hourly_sitemap_rolling_refresh_sync(
|
||||
*,
|
||||
redis_client: Redis,
|
||||
lane: str,
|
||||
limit: int | None,
|
||||
only_new: bool | None,
|
||||
) -> dict[str, object]:
|
||||
del limit, only_new
|
||||
settings = Settings()
|
||||
persistence = _get_persistence()
|
||||
persistence.create_tables()
|
||||
|
||||
discovery_result = discover_vehicle_urls_from_sitemap_with_stats(settings=settings)
|
||||
discovered_urls = discovery_result.vehicle_urls
|
||||
active_urls = set(discovered_urls)
|
||||
sold_count = 0
|
||||
if active_urls:
|
||||
sold_count = persistence.mark_sold_not_in_listing_by_urls(active_urls, lane="iaai")
|
||||
|
||||
existing_urls = persistence.get_all_active_origin_urls_for_lane("iaai:")
|
||||
new_urls = [url for url in discovered_urls if url not in existing_urls]
|
||||
|
||||
total_active = persistence.count_active_cars_for_lane("iaai:")
|
||||
batch_size = max(1, int(settings.discovery.hourly_refresh_batch_size))
|
||||
try:
|
||||
offset = int(redis_client.get(SITEMAP_HOURLY_REFRESH_OFFSET_KEY) or 0)
|
||||
except Exception:
|
||||
offset = 0
|
||||
|
||||
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
|
||||
prefix="iaai:",
|
||||
offset=offset,
|
||||
limit=batch_size,
|
||||
)
|
||||
if not refresh_urls and total_active > 0:
|
||||
offset = 0
|
||||
refresh_urls = persistence.get_active_origin_urls_batch_for_refresh(
|
||||
prefix="iaai:",
|
||||
offset=0,
|
||||
limit=batch_size,
|
||||
)
|
||||
|
||||
next_offset = 0
|
||||
if total_active > 0:
|
||||
next_offset = offset + len(refresh_urls)
|
||||
if next_offset >= total_active:
|
||||
next_offset = 0
|
||||
|
||||
try:
|
||||
redis_client.set(SITEMAP_HOURLY_REFRESH_OFFSET_KEY, str(next_offset))
|
||||
except Exception:
|
||||
logger.warning("Failed to persist hourly rolling refresh offset", exc_info=True)
|
||||
|
||||
seen: set[str] = set()
|
||||
target_urls: list[str] = []
|
||||
for url in new_urls + refresh_urls:
|
||||
if url in seen:
|
||||
continue
|
||||
seen.add(url)
|
||||
target_urls.append(url)
|
||||
|
||||
cars_upserted = 0
|
||||
cars_failed = 0
|
||||
images_upserted = 0
|
||||
failures: list[dict[str, str]] = []
|
||||
|
||||
if target_urls:
|
||||
with IAAIScraper(settings) as scraper:
|
||||
worker_batch_size = settings.celery.batch_size
|
||||
for batch_start in range(0, len(target_urls), worker_batch_size):
|
||||
batch_urls = target_urls[batch_start:batch_start + worker_batch_size]
|
||||
batch_result = scraper.sync_batch(batch_urls, lane=lane)
|
||||
cars_upserted += int(batch_result.get("cars_upserted", 0))
|
||||
cars_failed += int(batch_result.get("cars_failed", 0))
|
||||
images_upserted += int(batch_result.get("images_upserted", 0))
|
||||
failures.extend(batch_result.get("failures", []))
|
||||
|
||||
status = "success" if not failures else ("partial_success" if cars_upserted else "failed")
|
||||
return {
|
||||
"status": status,
|
||||
"run_id": None,
|
||||
"cars_upserted": cars_upserted,
|
||||
"cars_failed": cars_failed,
|
||||
"images_upserted": images_upserted,
|
||||
"skipped_existing": max(0, len(discovered_urls) - len(new_urls)),
|
||||
"elapsed_seconds": None,
|
||||
"failures": failures,
|
||||
"full_scan_completed": True,
|
||||
"hourly_mode": "sitemap_rolling_refresh",
|
||||
"discovered_urls": len(discovered_urls),
|
||||
"new_urls": len(new_urls),
|
||||
"refresh_urls": len(refresh_urls),
|
||||
"sold_marked": sold_count,
|
||||
"transport": discovery_result.stats.transport,
|
||||
"refresh_offset": offset,
|
||||
"refresh_next_offset": next_offset,
|
||||
"active_total": total_active,
|
||||
}
|
||||
|
||||
|
||||
def _start_stall_watchdog(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
task_id: str,
|
||||
stall_timeout_seconds: int,
|
||||
) -> tuple[Event, Thread]:
|
||||
stop_event = Event()
|
||||
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
|
||||
|
||||
def _watchdog() -> None:
|
||||
key = _task_progress_key(task_id)
|
||||
while not stop_event.wait(interval_seconds):
|
||||
try:
|
||||
raw = redis_client.get(key)
|
||||
if not raw:
|
||||
continue
|
||||
data = json.loads(raw)
|
||||
last_ts = int(data.get("ts") or 0)
|
||||
if not last_ts:
|
||||
continue
|
||||
age = int(time.time()) - last_ts
|
||||
if age < stall_timeout_seconds:
|
||||
continue
|
||||
logger.error(
|
||||
"Task %s stalled for %ss at stage=%s payload=%s; killing worker process for redelivery",
|
||||
task_id,
|
||||
age,
|
||||
data.get("stage"),
|
||||
data,
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
||||
continue
|
||||
os.kill(os.getpid(), signal.SIGKILL)
|
||||
|
||||
thread = Thread(target=_watchdog, name=f"task-stall-watchdog-{task_id[:8]}", daemon=True)
|
||||
thread.start()
|
||||
return stop_event, thread
|
||||
|
||||
|
||||
def _get_persistence() -> PersistenceService:
|
||||
settings = Settings()
|
||||
persistence = PersistenceService(settings)
|
||||
@@ -360,6 +648,204 @@ def _start_lock_heartbeat(
|
||||
return stop_event, thread
|
||||
|
||||
|
||||
def _reset_segments_progress(redis_client: Redis, total: int) -> None:
|
||||
try:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.delete(SYNC_SEGMENTS_PROGRESS_KEY)
|
||||
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
|
||||
pipe.execute()
|
||||
except Exception:
|
||||
logger.warning("Failed to reset segments progress", exc_info=True)
|
||||
|
||||
|
||||
def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
|
||||
"""Помечает сегмент завершённым. Возвращает (completed_count, total)."""
|
||||
try:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index)))
|
||||
pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
|
||||
pipe.scard(SYNC_SEGMENTS_PROGRESS_KEY)
|
||||
pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
|
||||
results = pipe.execute()
|
||||
completed = int(results[2] or 0)
|
||||
total = int(results[3] or 0) if results[3] else 0
|
||||
return completed, total
|
||||
except Exception:
|
||||
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True)
|
||||
return 0, 0
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_segment_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=60,
|
||||
acks_late=True,
|
||||
)
|
||||
def sync_segment_task(
|
||||
self,
|
||||
segment_index: int,
|
||||
segment: dict,
|
||||
lane: str = "iaai_cars",
|
||||
only_new: bool | None = None,
|
||||
is_bootstrap: bool = False,
|
||||
):
|
||||
"""Обработка одного сегмента листинга. Запускается параллельно несколькими воркерами."""
|
||||
persistence = _get_persistence()
|
||||
persistence.create_tables()
|
||||
task_id = self.request.id or "unknown"
|
||||
redis_client = _get_redis()
|
||||
settings = Settings()
|
||||
|
||||
# Per-segment lock — защита от случайного дубля
|
||||
seg_lock_key = SYNC_SEGMENT_LOCK_KEY_FMT.format(idx=int(segment_index))
|
||||
owner_token = f"{task_id}:{uuid.uuid4().hex}"
|
||||
lock_ttl = _sync_segment_lock_ttl_seconds()
|
||||
lock_acquired = _acquire_lock(redis_client, seg_lock_key, owner_token, lock_ttl)
|
||||
if not lock_acquired:
|
||||
logger.info("sync_segment_task[%d] skipped: already running", segment_index)
|
||||
return {"status": "skipped", "segment_index": segment_index, "reason": "duplicate"}
|
||||
|
||||
seg_make = segment.get("make")
|
||||
seg_year_min = segment.get("year_min")
|
||||
seg_year_max = segment.get("year_max")
|
||||
seg_label = f"{seg_make or 'ALL'}"
|
||||
if seg_year_min is not None or seg_year_max is not None:
|
||||
seg_label += f" ({seg_year_min}-{seg_year_max})"
|
||||
|
||||
heartbeat_stop: Event | None = None
|
||||
heartbeat_thread: Thread | None = None
|
||||
watchdog_stop: Event | None = None
|
||||
watchdog_thread: Thread | None = None
|
||||
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
|
||||
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
|
||||
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage="segment_task_started",
|
||||
ttl_seconds=progress_ttl,
|
||||
segment_index=segment_index,
|
||||
segment_label=seg_label,
|
||||
)
|
||||
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(redis_client, seg_lock_key, owner_token, lock_ttl)
|
||||
watchdog_stop, watchdog_thread = _start_stall_watchdog(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stall_timeout_seconds=stall_timeout,
|
||||
)
|
||||
|
||||
try:
|
||||
def _job():
|
||||
with IAAIScraper() as scraper:
|
||||
scraper.set_progress_callback(
|
||||
lambda stage, meta: _update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage=stage,
|
||||
ttl_seconds=progress_ttl,
|
||||
segment_index=segment_index,
|
||||
segment_label=seg_label,
|
||||
**meta,
|
||||
)
|
||||
)
|
||||
base_url = scraper.settings.listing.cars_url
|
||||
seg_url = scraper._build_segment_listing_url(base_url, seg_make) if seg_make else None
|
||||
return scraper.sync_listing(
|
||||
make=None if seg_url else seg_make,
|
||||
model=None,
|
||||
lane=lane,
|
||||
only_new=only_new,
|
||||
listing_url=seg_url,
|
||||
year_min=seg_year_min,
|
||||
year_max=seg_year_max,
|
||||
skip_mark_sold=True,
|
||||
)
|
||||
|
||||
result = _run_browser_job(_job)
|
||||
segment_done = bool(result.get("full_scan_completed", False))
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage="segment_task_completed",
|
||||
ttl_seconds=progress_ttl,
|
||||
segment_index=segment_index,
|
||||
segment_label=seg_label,
|
||||
status=result.get("status", "success"),
|
||||
cars_upserted=result.get("cars_upserted", 0),
|
||||
cars_failed=result.get("cars_failed", 0),
|
||||
)
|
||||
|
||||
if is_bootstrap and segment_done:
|
||||
completed, total = _mark_segment_completed(redis_client, segment_index)
|
||||
logger.warning(
|
||||
"Segment %d (%s) bootstrap done: %d/%d completed",
|
||||
segment_index, seg_label, completed, total,
|
||||
)
|
||||
if total > 0 and completed >= total:
|
||||
_set_full_scan_done(redis_client, True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_bootstrap_failure_streak(redis_client)
|
||||
logger.warning("All %d segments completed; bootstrap full scan done", total)
|
||||
|
||||
return {
|
||||
"status": result.get("status", "success"),
|
||||
"segment_index": segment_index,
|
||||
"segment_label": seg_label,
|
||||
"cars_upserted": result.get("cars_upserted", 0),
|
||||
"cars_failed": result.get("cars_failed", 0),
|
||||
"vehicles_collected": result.get("listing", {}).get("vehicles_collected", 0),
|
||||
"full_scan_completed": segment_done,
|
||||
}
|
||||
|
||||
except SoftTimeLimitExceeded:
|
||||
logger.warning("sync_segment_task[%d] soft timeout — partial progress saved", segment_index)
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage="segment_task_soft_timeout",
|
||||
ttl_seconds=progress_ttl,
|
||||
segment_index=segment_index,
|
||||
segment_label=seg_label,
|
||||
)
|
||||
return {
|
||||
"status": "timed_out",
|
||||
"segment_index": segment_index,
|
||||
"segment_label": seg_label,
|
||||
}
|
||||
except Exception as exc:
|
||||
logger.error("sync_segment_task[%d] failed: %s — %s", segment_index, seg_label, exc, exc_info=True)
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage="segment_task_failed",
|
||||
ttl_seconds=progress_ttl,
|
||||
segment_index=segment_index,
|
||||
segment_label=seg_label,
|
||||
error=str(exc),
|
||||
)
|
||||
try:
|
||||
raise self.retry(exc=exc)
|
||||
except self.MaxRetriesExceededError:
|
||||
return {
|
||||
"status": "failed",
|
||||
"segment_index": segment_index,
|
||||
"segment_label": seg_label,
|
||||
"error": str(exc),
|
||||
}
|
||||
finally:
|
||||
if watchdog_stop is not None:
|
||||
watchdog_stop.set()
|
||||
if watchdog_thread is not None:
|
||||
watchdog_thread.join(timeout=1)
|
||||
if heartbeat_stop is not None:
|
||||
heartbeat_stop.set()
|
||||
if heartbeat_thread is not None:
|
||||
heartbeat_thread.join(timeout=1)
|
||||
_clear_task_progress(redis_client, task_id)
|
||||
_release_lock_if_owner(redis_client, seg_lock_key, owner_token)
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_vehicle_task",
|
||||
bind=True,
|
||||
@@ -477,10 +963,20 @@ def sync_listing_task(
|
||||
}
|
||||
|
||||
try:
|
||||
settings = Settings()
|
||||
full_scan_done_before_run = _is_full_scan_done(redis_client)
|
||||
force_bootstrap_full_scan = not full_scan_done_before_run
|
||||
effective_limit = None if force_bootstrap_full_scan else limit
|
||||
effective_only_new = False if force_bootstrap_full_scan else only_new
|
||||
hourly_mode = settings.discovery.hourly_mode.strip().lower()
|
||||
discovery_mode = settings.discovery.mode.strip().lower()
|
||||
prefer_sitemap_mainline = (
|
||||
make is None
|
||||
and model is None
|
||||
and effective_limit is None
|
||||
and not effective_only_new
|
||||
)
|
||||
use_hourly_sitemap_sync = full_scan_done_before_run and prefer_sitemap_mainline
|
||||
|
||||
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
|
||||
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
|
||||
@@ -503,6 +999,58 @@ def sync_listing_task(
|
||||
effective_limit,
|
||||
)
|
||||
|
||||
if use_hourly_sitemap_sync:
|
||||
if hourly_mode == "diff":
|
||||
logger.info("Hourly mode: running sitemap diff sync instead of full listing traversal")
|
||||
result = _hourly_sitemap_diff_sync(
|
||||
lane=lane,
|
||||
limit=effective_limit,
|
||||
only_new=effective_only_new,
|
||||
)
|
||||
elif hourly_mode == "full_refresh":
|
||||
logger.info("Hourly mode: running sitemap full refresh of all active vehicles")
|
||||
result = _hourly_sitemap_full_refresh_sync(
|
||||
lane=lane,
|
||||
limit=effective_limit,
|
||||
only_new=effective_only_new,
|
||||
)
|
||||
else:
|
||||
logger.info("Hourly mode: running sitemap rolling refresh of active vehicles")
|
||||
result = _hourly_sitemap_rolling_refresh_sync(
|
||||
redis_client=redis_client,
|
||||
lane=lane,
|
||||
limit=effective_limit,
|
||||
only_new=effective_only_new,
|
||||
)
|
||||
try:
|
||||
redis_client.set(SITEMAP_HOURLY_LAST_COUNT_KEY, str(result.get("discovered_urls", 0)))
|
||||
except Exception:
|
||||
logger.warning("Failed to persist hourly sitemap count", exc_info=True)
|
||||
|
||||
summary = {
|
||||
"task_id": task_id,
|
||||
"run_id": result.get("run_id"),
|
||||
"status": result.get("status", "success"),
|
||||
"cars_upserted": result.get("cars_upserted", 0),
|
||||
"cars_failed": result.get("cars_failed", 0),
|
||||
"images_upserted": result.get("images_upserted", 0),
|
||||
"skipped_existing": result.get("skipped_existing", 0),
|
||||
"elapsed_seconds": result.get("elapsed_seconds"),
|
||||
"failures_count": len(result.get("failures") or []),
|
||||
"hourly_mode": result.get("hourly_mode"),
|
||||
"discovered_urls": result.get("discovered_urls", 0),
|
||||
"new_urls": result.get("new_urls", 0),
|
||||
"sold_marked": result.get("sold_marked", 0),
|
||||
}
|
||||
logger.info(
|
||||
"sync_listing_task hourly diff completed: status=%s, new=%d, sold=%d, skipped=%d",
|
||||
summary["status"],
|
||||
summary["new_urls"],
|
||||
summary["sold_marked"],
|
||||
summary["skipped_existing"],
|
||||
)
|
||||
return summary
|
||||
|
||||
_clear_followup_pending(redis_client)
|
||||
|
||||
heartbeat_stop, heartbeat_thread = _start_lock_heartbeat(
|
||||
@@ -514,9 +1062,91 @@ def sync_listing_task(
|
||||
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
|
||||
|
||||
# Определяем сегменты из конфига.
|
||||
settings = Settings()
|
||||
segments = parse_listing_segments(settings.listing.listing_segments_json)
|
||||
use_segmented = bool(segments) and make is None and model is None
|
||||
use_segmented = (
|
||||
bool(segments)
|
||||
and make is None
|
||||
and model is None
|
||||
and not prefer_sitemap_mainline
|
||||
and discovery_mode != "sitemap"
|
||||
)
|
||||
|
||||
if prefer_sitemap_mainline:
|
||||
if discovery_mode != "sitemap":
|
||||
logger.warning(
|
||||
"Unfiltered full scan forcing sitemap discovery despite IAAI_DISCOVERY_MODE=%s",
|
||||
discovery_mode or "unset",
|
||||
)
|
||||
if segments:
|
||||
logger.info("Ignoring configured listing segments for unfiltered sitemap full scan")
|
||||
|
||||
# --- Параллельный диспатч сегментов: dispatch & exit ---
|
||||
if use_segmented and settings.celery.parallel_segments:
|
||||
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
|
||||
already_completed: set[int] = set()
|
||||
if force_bootstrap_full_scan:
|
||||
try:
|
||||
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
|
||||
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()}
|
||||
except Exception:
|
||||
already_completed = set()
|
||||
|
||||
pending = [
|
||||
(idx, seg) for idx, seg in enumerate(segments)
|
||||
if idx not in already_completed
|
||||
]
|
||||
|
||||
if not pending:
|
||||
# Всё уже сделано — фиксируем bootstrap done.
|
||||
if force_bootstrap_full_scan:
|
||||
_set_full_scan_done(redis_client, True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_bootstrap_failure_streak(redis_client)
|
||||
logger.warning("Parallel segments: nothing to dispatch (all completed)")
|
||||
return {
|
||||
"status": "success",
|
||||
"task_id": task_id,
|
||||
"mode": "parallel_segments",
|
||||
"segments_total": len(segments),
|
||||
"segments_dispatched": 0,
|
||||
"segments_already_completed": len(already_completed),
|
||||
}
|
||||
|
||||
# При первом запуске bootstrap фиксируем total, чтобы знать когда остановиться.
|
||||
if force_bootstrap_full_scan and not already_completed:
|
||||
_reset_segments_progress(redis_client, len(segments))
|
||||
|
||||
dispatched = 0
|
||||
for idx, seg in pending:
|
||||
try:
|
||||
self.app.send_task(
|
||||
"iaai_scraper.worker.tasks.sync_segment_task",
|
||||
kwargs={
|
||||
"segment_index": idx,
|
||||
"segment": seg,
|
||||
"lane": lane,
|
||||
"only_new": effective_only_new,
|
||||
"is_bootstrap": force_bootstrap_full_scan,
|
||||
},
|
||||
queue="scraping",
|
||||
)
|
||||
dispatched += 1
|
||||
except Exception:
|
||||
logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
|
||||
|
||||
logger.warning(
|
||||
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)",
|
||||
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan,
|
||||
)
|
||||
return {
|
||||
"status": "success",
|
||||
"task_id": task_id,
|
||||
"mode": "parallel_segments",
|
||||
"segments_total": len(segments),
|
||||
"segments_dispatched": dispatched,
|
||||
"segments_already_completed": len(already_completed),
|
||||
}
|
||||
# --- конец параллельной ветки ---
|
||||
|
||||
resume_from_segment = 0
|
||||
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
|
||||
@@ -608,10 +1238,7 @@ def sync_listing_task(
|
||||
)
|
||||
if force_bootstrap_full_scan:
|
||||
_set_full_scan_done(redis_client, False)
|
||||
# SoftTimeLimit считаем failure для circuit breaker:
|
||||
# если сегмент систематически не укладывается в time limit,
|
||||
# нельзя крутить followup бесконечно.
|
||||
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=True)
|
||||
_enqueue_bootstrap_followup("soft_time_limit_exceeded", count_as_failure=False)
|
||||
# Partial progress уже записан в БД через finish_sync_run.
|
||||
# Не retry — следующий запуск продолжит обработку по расписанию.
|
||||
return {
|
||||
@@ -646,3 +1273,89 @@ def sync_listing_task(
|
||||
heartbeat_thread.join(timeout=max(1.0, min(5.0, lock_ttl / 10)))
|
||||
if lock_acquired:
|
||||
_release_lock_if_owner(redis_client, SYNC_LISTING_LOCK_KEY, owner_token)
|
||||
|
||||
|
||||
# Новые задачи для ingestion pipeline
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.discover_vehicles_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=60,
|
||||
acks_late=True,
|
||||
)
|
||||
def discover_vehicles_task(self, max_urls: int | None = None):
|
||||
"""Задача для обнаружения новых URL автомобилей."""
|
||||
from ..discovery_service import DiscoveryService
|
||||
|
||||
try:
|
||||
discovery = DiscoveryService()
|
||||
count = discovery.discover_new_vehicles(max_urls=max_urls)
|
||||
logger.info("discover_vehicles_task completed: %d candidates added", count)
|
||||
return {"status": "success", "candidates_added": count}
|
||||
except Exception as exc:
|
||||
logger.error("discover_vehicles_task failed: %s", exc, exc_info=True)
|
||||
raise self.retry(exc=exc)
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.fetch_pending_candidates_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=30,
|
||||
acks_late=True,
|
||||
)
|
||||
def fetch_pending_candidates_task(self, limit: int = 10):
|
||||
"""Задача для захвата данных ожидающих кандидатов."""
|
||||
from ..fetch_service import FetchService
|
||||
|
||||
try:
|
||||
fetch = FetchService()
|
||||
count = fetch.process_pending_candidates(limit=limit)
|
||||
logger.info("fetch_pending_candidates_task completed: %d candidates processed", count)
|
||||
return {"status": "success", "candidates_processed": count}
|
||||
except Exception as exc:
|
||||
logger.error("fetch_pending_candidates_task failed: %s", exc, exc_info=True)
|
||||
raise self.retry(exc=exc)
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.enrich_snapshots_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=30,
|
||||
acks_late=True,
|
||||
)
|
||||
def enrich_snapshots_task(self, limit: int = 10):
|
||||
"""Задача для парсинга и обогащения snapshots."""
|
||||
from ..enrichment_service import EnrichmentService
|
||||
|
||||
try:
|
||||
enrichment = EnrichmentService()
|
||||
count = enrichment.process_unparsed_snapshots(limit=limit)
|
||||
logger.info("enrich_snapshots_task completed: %d snapshots enriched", count)
|
||||
return {"status": "success", "snapshots_enriched": count}
|
||||
except Exception as exc:
|
||||
logger.error("enrich_snapshots_task failed: %s", exc, exc_info=True)
|
||||
raise self.retry(exc=exc)
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.run_ingestion_pipeline_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=120,
|
||||
acks_late=True,
|
||||
)
|
||||
def run_ingestion_pipeline_task(self):
|
||||
"""Задача для запуска полного ingestion pipeline."""
|
||||
from ..scheduler_service import SchedulerService
|
||||
|
||||
try:
|
||||
scheduler = SchedulerService()
|
||||
scheduler.run_full_pipeline()
|
||||
logger.info("run_ingestion_pipeline_task completed")
|
||||
return {"status": "success"}
|
||||
except Exception as exc:
|
||||
logger.error("run_ingestion_pipeline_task failed: %s", exc, exc_info=True)
|
||||
raise self.retry(exc=exc)
|
||||
|
||||
Reference in New Issue
Block a user