fast parser
This commit is contained in:
@@ -1,7 +1,9 @@
|
||||
# Инициализация Celery-приложения и периодических задач.
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
|
||||
from celery import Celery
|
||||
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
|
||||
@@ -13,6 +15,7 @@ from ..core.logs import setup_logging
|
||||
logger = logging.getLogger("iaai_scraper.worker.celery_app")
|
||||
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
|
||||
IAAI_SYNC_QUEUE = "iaai_sync"
|
||||
PROGRESS_KEY_PREFIX = "iaai:state:task_progress:"
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
@@ -20,6 +23,26 @@ def _env_bool(name: str, default: bool) -> bool:
|
||||
return raw in {"1", "true", "yes", "on"}
|
||||
|
||||
|
||||
def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 180) -> bool:
|
||||
now = int(time.time())
|
||||
try:
|
||||
for raw_key in redis_client.scan_iter(f"{PROGRESS_KEY_PREFIX}*"):
|
||||
payload = redis_client.get(raw_key)
|
||||
if not payload:
|
||||
continue
|
||||
try:
|
||||
progress = json.loads(payload)
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
ts = int(progress.get("ts") or 0)
|
||||
stage = str(progress.get("stage") or "")
|
||||
if ts > 0 and now - ts <= max_age_seconds and stage not in {"segment_done", "sync_done", "failed"}:
|
||||
return True
|
||||
except Exception:
|
||||
logger.warning("Failed to inspect startup progress keys", exc_info=True)
|
||||
return False
|
||||
|
||||
|
||||
@celery_setup_logging.connect
|
||||
def _configure_logging(loglevel=None, **kwargs):
|
||||
# Перехватываем логирование Celery и пишем только в stderr (Docker logs).
|
||||
@@ -55,7 +78,7 @@ _soft = settings.celery.task_soft_time_limit
|
||||
_hard = settings.celery.task_time_limit
|
||||
_max_hard = _soft + 120 if _soft else _hard
|
||||
if _hard > _max_hard:
|
||||
logger.warning(
|
||||
logger.info(
|
||||
"CELERY_TASK_TIME_LIMIT=%d too far from CELERY_TASK_SOFT_TIME_LIMIT=%d; "
|
||||
"clamping hard limit to %d",
|
||||
_hard, _soft, _max_hard,
|
||||
@@ -75,7 +98,7 @@ celery_app.conf.update(
|
||||
task_track_started=True,
|
||||
worker_concurrency=settings.celery.worker_concurrency,
|
||||
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
|
||||
worker_pool="solo",
|
||||
worker_pool=settings.celery.worker_pool,
|
||||
worker_prefetch_multiplier=1,
|
||||
broker_connection_retry_on_startup=True,
|
||||
broker_transport_options={
|
||||
@@ -123,14 +146,17 @@ def _on_worker_ready(**kwargs):
|
||||
retry_on_timeout=True,
|
||||
)
|
||||
|
||||
has_fresh_progress = _has_fresh_active_progress(redis_client)
|
||||
for stale_key in ("iaai:locks:sync_listing",):
|
||||
try:
|
||||
ttl = redis_client.ttl(stale_key)
|
||||
if ttl is not None and ttl != -2:
|
||||
if ttl is not None and ttl != -2 and not has_fresh_progress:
|
||||
redis_client.delete(stale_key)
|
||||
logger.warning("Cleared stale lock on startup: %s (ttl was %s)", stale_key, ttl)
|
||||
logger.info("Cleared stale lock on startup: %s (ttl was %s)", stale_key, ttl)
|
||||
elif ttl is not None and ttl != -2:
|
||||
logger.info("Keeping sync lock on startup because fresh active progress exists: %s (ttl=%s)", stale_key, ttl)
|
||||
except Exception:
|
||||
logger.warning("Failed to clear stale lock %s on startup", stale_key, exc_info=True)
|
||||
logger.warning("Failed to inspect stale lock %s on startup", stale_key, exc_info=True)
|
||||
|
||||
try:
|
||||
queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) or 0)
|
||||
@@ -141,6 +167,11 @@ def _on_worker_ready(**kwargs):
|
||||
return
|
||||
|
||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
if not should_dispatch and not has_fresh_progress:
|
||||
redis_client.delete(STARTUP_SYNC_DISPATCH_KEY)
|
||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
if should_dispatch:
|
||||
logger.info("Worker ready: stale startup dedupe key ignored because queue is empty and no fresh active progress exists")
|
||||
except Exception:
|
||||
logger.warning("Worker ready startup sync dedupe check failed; skipping immediate dispatch", exc_info=True)
|
||||
return
|
||||
|
||||
@@ -12,8 +12,12 @@ logger = logging.getLogger("iaai_scraper.worker.self_heal")
|
||||
|
||||
IAAI_SYNC_QUEUE = "iaai_sync"
|
||||
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||
GLOBAL_DB_PROGRESS_TS_KEY = "iaai:state:last_db_progress_ts"
|
||||
SELF_HEAL_RESTART_LOCK_KEY = "iaai:state:self_heal_restart_in_progress"
|
||||
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
||||
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
||||
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
|
||||
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
||||
SIGKILL_FALLBACK = getattr(signal, "SIGKILL", signal.SIGTERM)
|
||||
|
||||
|
||||
@@ -79,6 +83,41 @@ def _read_last_progress_ts(redis_client: Redis) -> int | None:
|
||||
return max_ts or None
|
||||
|
||||
|
||||
def _read_last_db_progress_ts(redis_client: Redis) -> int | None:
|
||||
raw = redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY)
|
||||
if raw:
|
||||
ts = _safe_int(raw)
|
||||
if ts > 0:
|
||||
return ts
|
||||
|
||||
max_ts = 0
|
||||
for key in redis_client.scan_iter(match="iaai:state:task_progress:*"):
|
||||
try:
|
||||
payload = redis_client.get(key)
|
||||
if not payload:
|
||||
continue
|
||||
data = json.loads(payload)
|
||||
if str(data.get("stage") or "") == "fast_db_progress":
|
||||
ts = _safe_int(data.get("ts"), 0)
|
||||
else:
|
||||
ts = _safe_int(data.get("last_db_progress_ts"), 0)
|
||||
if ts > max_ts:
|
||||
max_ts = ts
|
||||
except Exception:
|
||||
continue
|
||||
return max_ts or None
|
||||
|
||||
|
||||
def _reset_bootstrap_checkpoint_for_db_idle(redis_client: Redis) -> None:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.delete(SYNC_LISTING_CHECKPOINT_KEY)
|
||||
pipe.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY)
|
||||
pipe.delete(GLOBAL_PROGRESS_TS_KEY)
|
||||
pipe.delete(GLOBAL_DB_PROGRESS_TS_KEY)
|
||||
pipe.set(SYNC_FULL_SCAN_DONE_KEY, "0")
|
||||
pipe.execute()
|
||||
|
||||
|
||||
def _has_inflight_work(redis_client: Redis, queue_name: str) -> tuple[bool, dict[str, int]]:
|
||||
"""Есть ли признаки активной/зависшей работы, даже если очередь пуста."""
|
||||
queue_len = _safe_int(redis_client.llen(queue_name), 0)
|
||||
@@ -135,14 +174,16 @@ def main() -> None:
|
||||
queue_name = os.getenv("IAAI_CELERY_QUEUE", IAAI_SYNC_QUEUE)
|
||||
check_interval = max(5, _env_int("IAAI_SELF_HEAL_CHECK_INTERVAL_SECONDS", 30))
|
||||
stall_seconds = max(180, _env_int("IAAI_SELF_HEAL_STALL_SECONDS", 720))
|
||||
db_idle_seconds = max(60, _env_int("IAAI_DB_IDLE_RESTART_SECONDS", 3600))
|
||||
startup_grace = max(30, _env_int("IAAI_SELF_HEAL_STARTUP_GRACE_SECONDS", 300))
|
||||
restart_cooldown = max(60, _env_int("IAAI_SELF_HEAL_RESTART_COOLDOWN_SECONDS", 300))
|
||||
|
||||
logger.warning(
|
||||
"Self-heal watchdog enabled: queue=%s check_interval=%ss stall=%ss startup_grace=%ss cooldown=%ss",
|
||||
logger.info(
|
||||
"Self-heal watchdog enabled: queue=%s check_interval=%ss stall=%ss db_idle=%ss startup_grace=%ss cooldown=%ss",
|
||||
queue_name,
|
||||
check_interval,
|
||||
stall_seconds,
|
||||
db_idle_seconds,
|
||||
startup_grace,
|
||||
restart_cooldown,
|
||||
)
|
||||
@@ -175,9 +216,22 @@ def main() -> None:
|
||||
)
|
||||
age = stall_seconds + 1
|
||||
|
||||
restart_reason = f"progress_age={age}s > {stall_seconds}s"
|
||||
db_idle_restart = False
|
||||
if age <= stall_seconds:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
last_db_ts = _read_last_db_progress_ts(redis_client)
|
||||
db_age = None if last_db_ts is None else max(0, now_ts - int(last_db_ts))
|
||||
if db_age is None:
|
||||
first_allowed_ts = int(started_at) + max(startup_grace, db_idle_seconds)
|
||||
if now_ts < first_allowed_ts:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
db_age = db_idle_seconds + 1
|
||||
if db_age <= db_idle_seconds:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
db_idle_restart = True
|
||||
restart_reason = f"db_idle_age={db_age}s > {db_idle_seconds}s"
|
||||
|
||||
# Глобальный anti-storm lock: чтобы много воркеров не рестартились одновременно.
|
||||
acquired = bool(
|
||||
@@ -193,11 +247,13 @@ def main() -> None:
|
||||
continue
|
||||
|
||||
logger.error(
|
||||
"Self-heal: detected global stall (inflight=%s, progress_age=%ss > %ss). Restarting worker process...",
|
||||
"Self-heal: detected stall (inflight=%s, %s). Restarting worker process...",
|
||||
inflight,
|
||||
age,
|
||||
stall_seconds,
|
||||
restart_reason,
|
||||
)
|
||||
if db_idle_restart:
|
||||
logger.error("Self-heal: no DB writes for too long; clearing checkpoint to restart from segment 1")
|
||||
_reset_bootstrap_checkpoint_for_db_idle(redis_client)
|
||||
# Небольшой джиттер, чтобы при одинаковом событии у разных контейнеров
|
||||
# перезапуск был не строго одновременно.
|
||||
time.sleep(random.uniform(0.3, 2.0))
|
||||
|
||||
@@ -12,7 +12,8 @@ from billiard.exceptions import SoftTimeLimitExceeded
|
||||
from celery import shared_task
|
||||
from redis import Redis
|
||||
|
||||
from ..core.config import Settings, parse_listing_segments
|
||||
from ..core.config import Settings, build_fast_listing_segments_for_makes, build_listing_segments_for_makes, parse_listing_segments
|
||||
from ..core.runtime_config import RuntimeConfig
|
||||
from ..scraper import IAAIScraper
|
||||
from ..storage.db import PersistenceService
|
||||
from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitemap_with_stats
|
||||
@@ -53,16 +54,35 @@ 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}"
|
||||
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||
GLOBAL_DB_PROGRESS_TS_KEY = "iaai:state:last_db_progress_ts"
|
||||
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
|
||||
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
|
||||
STALL_WATCHDOG_NAVIGATION_STAGES = {
|
||||
"listing_next_page_started",
|
||||
"listing_resume_progress",
|
||||
}
|
||||
STALL_WATCHDOG_LONG_RUNNING_STAGES = {
|
||||
"fast_listing_collected",
|
||||
"fast_detail_progress",
|
||||
}
|
||||
STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max(
|
||||
300,
|
||||
int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")),
|
||||
)
|
||||
STALL_WATCHDOG_DETAIL_GRACE_SECONDS = max(
|
||||
600,
|
||||
int(os.getenv("STALL_WATCHDOG_DETAIL_GRACE_SECONDS", "1200")),
|
||||
)
|
||||
DB_IDLE_RESTART_SECONDS = max(60, int(os.getenv("IAAI_DB_IDLE_RESTART_SECONDS", "3600")))
|
||||
TERMINAL_PROGRESS_STAGES = {
|
||||
"segment_done",
|
||||
"segment_failed",
|
||||
"segment_task_completed",
|
||||
"segment_task_failed",
|
||||
"segment_task_soft_timeout",
|
||||
"sync_done",
|
||||
"failed",
|
||||
}
|
||||
|
||||
|
||||
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||
@@ -123,6 +143,19 @@ def _update_task_progress(
|
||||
) -> None:
|
||||
try:
|
||||
now_ts = int(time.time())
|
||||
existing_task_started_ts: int | None = None
|
||||
try:
|
||||
existing_raw = redis_client.get(_task_progress_key(task_id))
|
||||
if existing_raw:
|
||||
existing_payload = json.loads(existing_raw)
|
||||
existing_task_started_ts = _safe_int(existing_payload.get("task_started_ts"))
|
||||
except Exception:
|
||||
existing_task_started_ts = None
|
||||
payload.setdefault("task_started_ts", existing_task_started_ts or now_ts)
|
||||
if stage != "fast_db_progress" and "last_db_progress_ts" not in payload:
|
||||
last_db_progress_ts = _safe_int(redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY))
|
||||
if last_db_progress_ts is not None:
|
||||
payload["last_db_progress_ts"] = last_db_progress_ts
|
||||
data = {
|
||||
"task_id": task_id,
|
||||
"stage": stage,
|
||||
@@ -140,6 +173,8 @@ def _update_task_progress(
|
||||
# Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера
|
||||
# (когда PID жив, но прогресс по задачам не двигается).
|
||||
pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
|
||||
if stage == "fast_db_progress":
|
||||
pipe.set(GLOBAL_DB_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
|
||||
pipe.execute()
|
||||
except Exception:
|
||||
logger.warning("Failed to update task progress for %s", task_id, exc_info=True)
|
||||
@@ -155,6 +190,8 @@ def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
|
||||
def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int:
|
||||
if stage in STALL_WATCHDOG_NAVIGATION_STAGES:
|
||||
return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS)
|
||||
if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES:
|
||||
return max(int(default_timeout), STALL_WATCHDOG_DETAIL_GRACE_SECONDS)
|
||||
return int(default_timeout)
|
||||
|
||||
|
||||
@@ -407,6 +444,7 @@ def _start_stall_watchdog(
|
||||
stall_timeout_seconds: int,
|
||||
lock_key: str | None = None,
|
||||
lock_owner: str | None = None,
|
||||
db_idle_restart_seconds: int | None = None,
|
||||
) -> tuple[Event, Thread]:
|
||||
stop_event = Event()
|
||||
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
|
||||
@@ -418,6 +456,7 @@ def _start_stall_watchdog(
|
||||
watchdog_born = time.monotonic()
|
||||
absolute_deadline = stall_timeout_seconds * 3
|
||||
while not stop_event.wait(interval_seconds):
|
||||
db_idle_restart = False
|
||||
try:
|
||||
raw = redis_client.get(key)
|
||||
if not raw:
|
||||
@@ -443,18 +482,32 @@ def _start_stall_watchdog(
|
||||
last_ts = int(data.get("ts") or 0)
|
||||
if not last_ts:
|
||||
continue
|
||||
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
|
||||
age = int(time.time()) - last_ts
|
||||
if age < effective_stall_timeout:
|
||||
watchdog_born = time.monotonic() # reset absolute deadline on real progress
|
||||
continue
|
||||
logger.error(
|
||||
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
||||
task_id,
|
||||
age,
|
||||
stage,
|
||||
data,
|
||||
db_idle_restart = bool(
|
||||
db_idle_restart_seconds
|
||||
and _should_restart_for_db_idle(data, db_idle_restart_seconds)
|
||||
)
|
||||
if db_idle_restart:
|
||||
logger.error(
|
||||
"Task %s has no DB writes for >%ss at segment=%s/%s stage=%s; full restart required",
|
||||
task_id,
|
||||
db_idle_restart_seconds,
|
||||
data.get("segment_index"),
|
||||
data.get("segments_total"),
|
||||
stage,
|
||||
)
|
||||
else:
|
||||
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
|
||||
age = int(time.time()) - last_ts
|
||||
if age < effective_stall_timeout:
|
||||
watchdog_born = time.monotonic() # reset absolute deadline on real progress
|
||||
continue
|
||||
logger.error(
|
||||
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
||||
task_id,
|
||||
age,
|
||||
stage,
|
||||
data,
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
||||
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
|
||||
@@ -463,6 +516,12 @@ def _start_stall_watchdog(
|
||||
else:
|
||||
continue
|
||||
|
||||
if db_idle_restart:
|
||||
_restart_bootstrap_from_first_segment(
|
||||
redis_client,
|
||||
reason=f"no DB writes for >{db_idle_restart_seconds}s",
|
||||
)
|
||||
|
||||
# ── Pre-SIGTERM cleanup: release lock so next task can run ──
|
||||
if lock_key and lock_owner:
|
||||
try:
|
||||
@@ -515,6 +574,51 @@ def _start_stall_watchdog(
|
||||
return stop_event, thread
|
||||
|
||||
|
||||
def _should_restart_for_db_idle(progress: dict, db_idle_restart_seconds: int) -> bool:
|
||||
stage = str(progress.get("stage") or "")
|
||||
if stage in TERMINAL_PROGRESS_STAGES:
|
||||
return False
|
||||
if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES:
|
||||
timeout = _stall_timeout_for_progress(stage, db_idle_restart_seconds)
|
||||
progress_ts = _safe_int(progress.get("ts")) or 0
|
||||
return progress_ts > 0 and int(time.time()) - progress_ts >= timeout
|
||||
|
||||
segments_total = _safe_int(progress.get("segments_total"))
|
||||
segment_index = _safe_int(progress.get("segment_index"))
|
||||
if segments_total is None or segment_index is None:
|
||||
return False
|
||||
if segments_total <= 0 or segment_index >= segments_total - 1:
|
||||
return False
|
||||
|
||||
now_ts = int(time.time())
|
||||
progress_ts = _safe_int(progress.get("ts")) or 0
|
||||
if progress_ts <= 0:
|
||||
return False
|
||||
|
||||
db_progress_ts = _safe_int(progress.get("last_db_progress_ts"))
|
||||
if db_progress_ts is None:
|
||||
db_progress_ts = _read_global_db_progress_ts()
|
||||
task_started_ts = _safe_int(progress.get("task_started_ts")) or progress_ts
|
||||
last_db_or_start_ts = max(db_progress_ts or 0, task_started_ts)
|
||||
return now_ts - last_db_or_start_ts >= int(db_idle_restart_seconds)
|
||||
|
||||
|
||||
def _safe_int(value) -> int | None:
|
||||
try:
|
||||
return int(value)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _read_global_db_progress_ts() -> int | None:
|
||||
try:
|
||||
redis_client = _get_redis()
|
||||
raw = redis_client.get(GLOBAL_DB_PROGRESS_TS_KEY)
|
||||
return _safe_int(raw)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
_persistence_instance: PersistenceService | None = None
|
||||
|
||||
|
||||
@@ -722,6 +826,21 @@ def _save_last_completed_segment(redis_client: Redis, segment_index: int) -> Non
|
||||
logger.warning("Failed to save sync listing checkpoint", exc_info=True)
|
||||
|
||||
|
||||
def _restart_bootstrap_from_first_segment(redis_client: Redis, *, reason: str) -> None:
|
||||
"""Clear bootstrap checkpoint so the next full scan starts from segment 1."""
|
||||
try:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.delete(SYNC_LISTING_CHECKPOINT_KEY)
|
||||
pipe.delete(SYNC_LISTING_FOLLOWUP_PENDING_KEY)
|
||||
pipe.delete(GLOBAL_PROGRESS_TS_KEY)
|
||||
pipe.delete(GLOBAL_DB_PROGRESS_TS_KEY)
|
||||
pipe.set(SYNC_FULL_SCAN_DONE_KEY, "0")
|
||||
pipe.execute()
|
||||
logger.error("Bootstrap restart requested from segment 1: %s", reason)
|
||||
except Exception:
|
||||
logger.warning("Failed to reset bootstrap checkpoint for DB-idle restart", exc_info=True)
|
||||
|
||||
|
||||
def _clear_sync_checkpoint(redis_client: Redis) -> None:
|
||||
try:
|
||||
redis_client.delete(SYNC_LISTING_CHECKPOINT_KEY)
|
||||
@@ -852,6 +971,7 @@ def _start_lock_heartbeat(
|
||||
key: str,
|
||||
owner_token: str,
|
||||
ttl_seconds: int,
|
||||
stop_on_lost: bool = True,
|
||||
) -> tuple[Event, Thread]:
|
||||
stop_event = Event()
|
||||
interval_seconds = max(5.0, min(30.0, ttl_seconds / 3))
|
||||
@@ -862,7 +982,10 @@ def _start_lock_heartbeat(
|
||||
refreshed = _refresh_lock_if_owner(redis_client, key, owner_token, ttl_seconds)
|
||||
if refreshed is False:
|
||||
logger.warning("Lost sync_listing lock ownership for %s", owner_token)
|
||||
return
|
||||
if stop_on_lost:
|
||||
return
|
||||
consecutive_failures += 1
|
||||
continue
|
||||
if refreshed is None:
|
||||
consecutive_failures += 1
|
||||
if consecutive_failures >= 5:
|
||||
@@ -903,6 +1026,33 @@ def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[in
|
||||
return 0, 0
|
||||
|
||||
|
||||
def _build_listing_segments(settings: Settings) -> list[dict[str, str | int | None]]:
|
||||
"""Возвращает сегменты листинга с учётом runtime_config.
|
||||
|
||||
При IAAI_LISTING_SEGMENTS=runtime сегменты строятся из filters.brands / filters.include.brands.
|
||||
Остальные значения IAAI_LISTING_SEGMENTS сохраняют прежнее поведение: auto или JSON.
|
||||
"""
|
||||
raw_segments = settings.listing.listing_segments_json.strip()
|
||||
if raw_segments.casefold() != "runtime":
|
||||
return parse_listing_segments(raw_segments)
|
||||
|
||||
parsed_override = parse_listing_segments(raw_segments)
|
||||
if parsed_override:
|
||||
return parsed_override
|
||||
|
||||
runtime_config = RuntimeConfig.from_file(settings.runtime_config_file)
|
||||
brands = runtime_config.filters.include.brands
|
||||
if settings.scraping_profile.http_first and settings.listing.fast_segment_year_splits:
|
||||
segments = build_fast_listing_segments_for_makes(list(brands))
|
||||
else:
|
||||
segments = build_listing_segments_for_makes(list(brands))
|
||||
if not segments:
|
||||
logger.warning(
|
||||
"IAAI_LISTING_SEGMENTS=runtime, but runtime_config filters.brands is empty; segmented listing disabled"
|
||||
)
|
||||
return segments
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_segment_task",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
@@ -980,14 +1130,11 @@ def sync_segment_task(
|
||||
**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,
|
||||
make=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,
|
||||
@@ -1234,16 +1381,10 @@ def sync_listing_task(
|
||||
SYNC_LISTING_LOCK_KEY,
|
||||
owner_token,
|
||||
lock_ttl,
|
||||
stop_on_lost=False,
|
||||
)
|
||||
stall_timeout = max(120, int(settings.celery.task_stall_timeout_seconds))
|
||||
progress_ttl = max(lock_ttl + 120, stall_timeout + 120)
|
||||
watchdog_stop, watchdog_thread = _start_stall_watchdog(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stall_timeout_seconds=stall_timeout,
|
||||
lock_key=SYNC_LISTING_LOCK_KEY,
|
||||
lock_owner=owner_token,
|
||||
)
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
@@ -1253,7 +1394,17 @@ def sync_listing_task(
|
||||
|
||||
full_scan_done_before_run = _is_full_scan_done(redis_client)
|
||||
always_full_scan = bool(settings.discovery.always_full_scan)
|
||||
force_bootstrap_full_scan = always_full_scan or (not full_scan_done_before_run)
|
||||
explicit_filtered_run = bool(
|
||||
make is not None
|
||||
or model is not None
|
||||
or (limit is not None and int(limit or 0) > 0)
|
||||
or only_new is True
|
||||
)
|
||||
force_bootstrap_full_scan = (always_full_scan or (not full_scan_done_before_run)) and not explicit_filtered_run
|
||||
if explicit_filtered_run and not full_scan_done_before_run:
|
||||
logger.info(
|
||||
"Explicit sync_listing request detected; honoring make/model/limit/only_new before bootstrap full scan is complete",
|
||||
)
|
||||
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()
|
||||
@@ -1270,13 +1421,29 @@ def sync_listing_task(
|
||||
and not effective_only_new
|
||||
and not always_full_scan
|
||||
)
|
||||
use_hourly_sitemap_sync = (not always_full_scan) and full_scan_done_before_run and prefer_sitemap_mainline
|
||||
# Определяем сегменты из конфига/env или из runtime_config при IAAI_LISTING_SEGMENTS=runtime.
|
||||
segments = _build_listing_segments(settings)
|
||||
|
||||
watchdog_stop, watchdog_thread = _start_stall_watchdog(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stall_timeout_seconds=stall_timeout,
|
||||
lock_key=SYNC_LISTING_LOCK_KEY,
|
||||
lock_owner=owner_token,
|
||||
db_idle_restart_seconds=(DB_IDLE_RESTART_SECONDS if segments else None),
|
||||
)
|
||||
use_hourly_sitemap_sync = (
|
||||
(not always_full_scan)
|
||||
and full_scan_done_before_run
|
||||
and prefer_sitemap_mainline
|
||||
and not segments
|
||||
)
|
||||
|
||||
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
|
||||
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
|
||||
# Никаких page-level resume — внутри сегмента всегда стартуем с page 1.
|
||||
last_completed_segment: int | None = None
|
||||
if force_bootstrap_full_scan:
|
||||
if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run):
|
||||
last_completed_segment = _load_last_completed_segment(redis_client)
|
||||
else:
|
||||
# После завершения bootstrap чекпоинт не нужен никогда.
|
||||
@@ -1418,14 +1585,11 @@ def sync_listing_task(
|
||||
|
||||
self.update_state(state="STARTED", meta={"stage": "sync_listing_started", "task_id": task_id})
|
||||
|
||||
# Определяем сегменты из конфига.
|
||||
segments = parse_listing_segments(settings.listing.listing_segments_json)
|
||||
use_segmented = (
|
||||
bool(segments)
|
||||
and make is None
|
||||
and model is None
|
||||
and not prefer_sitemap_mainline
|
||||
and discovery_mode != "sitemap"
|
||||
and not use_hourly_sitemap_sync
|
||||
)
|
||||
|
||||
if prefer_sitemap_mainline:
|
||||
@@ -1543,7 +1707,7 @@ def sync_listing_task(
|
||||
start_page=1,
|
||||
progress_callback=(
|
||||
(lambda seg_idx: _save_last_completed_segment(redis_client, seg_idx))
|
||||
if force_bootstrap_full_scan else None
|
||||
if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run) else None
|
||||
),
|
||||
)
|
||||
return scraper.sync_listing(
|
||||
|
||||
Reference in New Issue
Block a user