fix: beat schedule should use only_new=False for full scans

This commit is contained in:
qananasikq
2026-04-25 14:55:36 +03:00
parent f8206615a9
commit 331d8d350e
5 changed files with 219 additions and 250 deletions

View File

@@ -112,7 +112,7 @@ celery_app.conf.update(
"task": "iaai.sync_cars_feed",
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
"args": (),
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True},
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": False},
"options": {
"queue": IAAI_SYNC_QUEUE,
"expires": settings.celery.beat_sync_interval_minutes * 60.0,

View File

@@ -12,8 +12,7 @@ from billiard.exceptions import SoftTimeLimitExceeded
from celery import shared_task
from redis import Redis
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 ..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
@@ -54,35 +53,16 @@ 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):
@@ -143,19 +123,6 @@ 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,
@@ -173,8 +140,6 @@ 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)
@@ -190,8 +155,6 @@ 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)
@@ -444,7 +407,6 @@ 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))
@@ -456,7 +418,6 @@ 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:
@@ -482,32 +443,18 @@ def _start_stall_watchdog(
last_ts = int(data.get("ts") or 0)
if not last_ts:
continue
db_idle_restart = bool(
db_idle_restart_seconds
and _should_restart_for_db_idle(data, db_idle_restart_seconds)
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,
)
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 тоже не отвечает дольше дедлайна — убиваем.
@@ -516,12 +463,6 @@ 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:
@@ -574,55 +515,6 @@ 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
# Listing/segment progress means the task is alive even if it writes no new
# cars. This is normal for hourly only_new runs when all vehicles are already
# present in DB. Do not kill healthy scans just because DB progress is idle.
progress_ts = _safe_int(progress.get("ts")) or 0
if progress_ts > 0 and int(time.time()) - progress_ts < int(db_idle_restart_seconds):
return False
if stage in STALL_WATCHDOG_LONG_RUNNING_STAGES:
timeout = _stall_timeout_for_progress(stage, db_idle_restart_seconds)
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())
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
@@ -830,21 +722,6 @@ 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)
@@ -975,7 +852,6 @@ 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))
@@ -986,10 +862,7 @@ 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)
if stop_on_lost:
return
consecutive_failures += 1
continue
return
if refreshed is None:
consecutive_failures += 1
if consecutive_failures >= 5:
@@ -1030,49 +903,6 @@ 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.
"""
filtered_urls = settings.listing.filtered_search_urls
if len(filtered_urls) > 1:
return [
{"make": None, "year_min": None, "year_max": None, "listing_url": url}
for url in filtered_urls
]
if len(filtered_urls) == 1:
return []
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
# Fast HTTP-first режим не использует UI year-фильтры IAAI.
# Значит runtime-сегменты должны быть "один бренд = один сегмент"
# без доп. year split, иначе получаем много лишних долгих сегментов,
# которые fast-профиль всё равно игнорирует.
if settings.scraping_profile.http_first:
segments = build_fast_listing_segments_for_makes(list(brands))
elif settings.listing.fast_segment_year_splits:
segments = build_listing_segments_for_makes(list(brands))
else:
segments = build_fast_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,
@@ -1108,10 +938,7 @@ def sync_segment_task(
seg_make = segment.get("make")
seg_year_min = segment.get("year_min")
seg_year_max = segment.get("year_max")
seg_listing_url = segment.get("listing_url")
seg_label = f"{seg_make or 'ALL'}"
if seg_listing_url:
seg_label = "FILTERED_SEARCH"
if seg_year_min is not None or seg_year_max is not None:
seg_label += f" ({seg_year_min}-{seg_year_max})"
@@ -1153,12 +980,14 @@ 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=seg_make,
make=None if seg_url else seg_make,
model=None,
lane=lane,
only_new=only_new,
listing_url=str(seg_listing_url) if seg_listing_url else None,
listing_url=seg_url,
year_min=seg_year_min,
year_max=seg_year_max,
skip_mark_sold=True,
@@ -1405,10 +1234,16 @@ 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,
@@ -1418,17 +1253,7 @@ def sync_listing_task(
full_scan_done_before_run = _is_full_scan_done(redis_client)
always_full_scan = bool(settings.discovery.always_full_scan)
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",
)
force_bootstrap_full_scan = always_full_scan or (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()
@@ -1445,31 +1270,13 @@ def sync_listing_task(
and not effective_only_new
and not always_full_scan
)
# Определяем сегменты из конфига/env или из runtime_config при IAAI_LISTING_SEGMENTS=runtime.
filtered_listing_urls = settings.listing.filtered_search_urls
filtered_listing_url = filtered_listing_urls[0] if filtered_listing_urls else None
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
)
use_hourly_sitemap_sync = (not always_full_scan) and full_scan_done_before_run and prefer_sitemap_mainline
# Segment-level checkpoint: хранит индекс последнего ПОЛНОСТЬЮ пройденного сегмента.
# Используется только во время bootstrap для пропуска уже обработанных сегментов.
# Никаких page-level resume — внутри сегмента всегда стартуем с page 1.
last_completed_segment: int | None = None
if force_bootstrap_full_scan and (always_full_scan or not full_scan_done_before_run):
if force_bootstrap_full_scan:
last_completed_segment = _load_last_completed_segment(redis_client)
else:
# После завершения bootstrap чекпоинт не нужен никогда.
@@ -1611,11 +1418,14 @@ 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 use_hourly_sitemap_sync
and not prefer_sitemap_mainline
and discovery_mode != "sitemap"
)
if prefer_sitemap_mainline:
@@ -1733,7 +1543,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 and (always_full_scan or not full_scan_done_before_run) else None
if force_bootstrap_full_scan else None
),
)
return scraper.sync_listing(
@@ -1742,7 +1552,6 @@ def sync_listing_task(
lane=lane,
limit=effective_limit,
only_new=effective_only_new,
listing_url=filtered_listing_url,
)
result = _run_browser_job(_job)