diff --git a/.gitignore b/.gitignore index 0c27e35..8f86660 100644 --- a/.gitignore +++ b/.gitignore @@ -19,4 +19,6 @@ build/ artifacts/ celerybeat-schedule* tokens_data/ -tokens.json \ No newline at end of file +tokens.json +iaai_deploy.tar.gz +.tmp_iaai_fast_ref/ \ No newline at end of file diff --git a/docker-compose.yml b/docker-compose.yml index 87ea783..596f342 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -7,38 +7,53 @@ x-app-env: &app-env IAAI_DATABASE_POOL_RECYCLE_SECONDS: ${IAAI_DATABASE_POOL_RECYCLE_SECONDS:-1800} IAAI_DATABASE_POOL_SIZE: ${IAAI_DATABASE_POOL_SIZE:-20} IAAI_DATABASE_MAX_OVERFLOW: ${IAAI_DATABASE_MAX_OVERFLOW:-40} - CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-900} - CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-1200} - CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-2400} + # Для больших Search-выборок (десятки тысяч лотов) run должен жить дольше одного батча. + CELERY_TASK_SOFT_TIME_LIMIT: ${CELERY_TASK_SOFT_TIME_LIMIT:-7200} + CELERY_TASK_TIME_LIMIT: ${CELERY_TASK_TIME_LIMIT:-7500} + CELERY_BROKER_VISIBILITY_TIMEOUT: ${CELERY_BROKER_VISIBILITY_TIMEOUT:-10800} IAAI_PROFILE: ${IAAI_PROFILE:-fast} IAAI_SCRAPING_PROFILE: ${IAAI_SCRAPING_PROFILE:-${IAAI_PROFILE:-fast}} IAAI_HTTP_FIRST: ${IAAI_HTTP_FIRST:-true} IAAI_BROWSER_FALLBACK_ENABLED: ${IAAI_BROWSER_FALLBACK_ENABLED:-false} IAAI_ANONYMOUS_BOOTSTRAP_ENABLED: ${IAAI_ANONYMOUS_BOOTSTRAP_ENABLED:-false} IAAI_CHALLENGE_REFRESH_ATTEMPTS: ${IAAI_CHALLENGE_REFRESH_ATTEMPTS:-1} - IAAI_LISTING_POST_ATTEMPTS: ${IAAI_LISTING_POST_ATTEMPTS:-2} - IAAI_REQUEST_JITTER_MAX_S: ${IAAI_REQUEST_JITTER_MAX_S:-0} - IAAI_DETAIL_RETRIES: ${IAAI_DETAIL_RETRIES:-1} + IAAI_LISTING_POST_ATTEMPTS: ${IAAI_LISTING_POST_ATTEMPTS:-1} + IAAI_REQUEST_JITTER_MAX_S: ${IAAI_REQUEST_JITTER_MAX_S:-0.03} + IAAI_DETAIL_RETRIES: ${IAAI_DETAIL_RETRIES:-2} IAAI_LISTING_RETRIES: ${IAAI_LISTING_RETRIES:-1} # 0 = без лимита. CELERY_BEAT_SYNC_LIMIT: ${CELERY_BEAT_SYNC_LIMIT:-0} + # Следующий плановый прогон — не раньше чем через час после предыдущего tick beat. + CELERY_BEAT_SYNC_INTERVAL_MINUTES: ${CELERY_BEAT_SYNC_INTERVAL_MINUTES:-60} + # Любой внутренний follow-up после неполного bootstrap тоже не стартует сразу. + IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS: ${IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS:-3600} CELERY_WORKER_CONCURRENCY: ${CELERY_WORKER_CONCURRENCY:-4} CELERY_WORKER_POOL: ${CELERY_WORKER_POOL:-prefork} CELERY_WORKER_MAX_TASKS_PER_CHILD: ${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} - CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-1000} + CELERY_BATCH_SIZE: ${CELERY_BATCH_SIZE:-1500} IAAI_INTER_BATCH_DELAY_SECONDS: ${IAAI_INTER_BATCH_DELAY_SECONDS:-0} + # Стабильный fast-профиль: не душим IAAI слишком большим числом HTTPS detail-соединений. + IAAI_FETCH_CONCURRENCY: ${IAAI_FETCH_CONCURRENCY:-96} + IAAI_MAX_RETRIES: ${IAAI_MAX_RETRIES:-2} + IAAI_RETRY_DELAY_SECONDS: ${IAAI_RETRY_DELAY_SECONDS:-0.35} + CELERY_PARALLEL_SEGMENTS: ${CELERY_PARALLEL_SEGMENTS:-false} IAAI_FAIL_RATE_THRESHOLD: ${IAAI_FAIL_RATE_THRESHOLD:-0.9} IAAI_PARALLEL_TABS: ${IAAI_PARALLEL_TABS:-8} - IAAI_FETCH_CONCURRENCY: ${IAAI_FETCH_CONCURRENCY:-32} IAAI_BLOCK_RESOURCES: ${IAAI_BLOCK_RESOURCES:-true} IAAI_FILTERED_SEARCH_URL: ${IAAI_FILTERED_SEARCH_URL:-} IAAI_FILTERED_SEARCH_URLS: ${IAAI_FILTERED_SEARCH_URLS:-} - IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-runtime} + IAAI_FAST_PATH_TIMEOUT_MS: ${IAAI_FAST_PATH_TIMEOUT_MS:-10000} + IAAI_LISTING_SEGMENTS: ${IAAI_LISTING_SEGMENTS:-[]} + IAAI_DISCOVERY_MODE: ${IAAI_DISCOVERY_MODE:-listing} + IAAI_HOURLY_MODE: ${IAAI_HOURLY_MODE:-rolling_refresh} + IAAI_HOURLY_REFRESH_BATCH_SIZE: ${IAAI_HOURLY_REFRESH_BATCH_SIZE:-500} IAAI_MAX_PAGES_PER_RUN: ${IAAI_MAX_PAGES_PER_RUN:-9999} IAAI_MAX_VEHICLES_PER_RUN: ${IAAI_MAX_VEHICLES_PER_RUN:-50000} - IAAI_ALWAYS_FULL_SCAN: ${IAAI_ALWAYS_FULL_SCAN:-true} - IAAI_HUMAN_PACE_ENABLED: ${IAAI_HUMAN_PACE_ENABLED:-true} - IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/home/app/tokens.json} + IAAI_ALWAYS_FULL_SCAN: ${IAAI_ALWAYS_FULL_SCAN:-false} + IAAI_HUMAN_PACE_ENABLED: ${IAAI_HUMAN_PACE_ENABLED:-false} + # Первый sync после clean restart стартует сразу; следующий запуск — только через hourly guard/beat. + IAAI_STARTUP_SYNC_ENABLED: ${IAAI_STARTUP_SYNC_ENABLED:-true} + IAAI_TOKENS_FILE: ${IAAI_TOKENS_FILE:-/data/tokens.json} IAAI_RUNTIME_CONFIG_FILE: ${IAAI_RUNTIME_CONFIG_FILE:-/app/runtime_config.json} IAAI_BROWSER_ENGINE: chromium # Self-heal для worker. @@ -168,7 +183,9 @@ services: stop_grace_period: 60s command: > celery -A iaai_scraper.worker.celery_app worker - --loglevel=info --concurrency=${CELERY_WORKER_CONCURRENCY:-4} --pool=${CELERY_WORKER_POOL:-prefork} + --loglevel=info + --concurrency=${CELERY_WORKER_CONCURRENCY:-4} + --pool=${CELERY_WORKER_POOL:-prefork} --pidfile=/tmp/celery-worker.pid -Q iaai_sync --max-tasks-per-child=${CELERY_WORKER_MAX_TASKS_PER_CHILD:-5} healthcheck: diff --git a/iaai_scraper/browser/fast_client.py b/iaai_scraper/browser/fast_client.py index 205857e..1910596 100644 --- a/iaai_scraper/browser/fast_client.py +++ b/iaai_scraper/browser/fast_client.py @@ -519,7 +519,17 @@ class IAAIFastClient: if max_pages is not None and max_pages > 0: total_pages = min(total_pages, max_pages) gbp_search_query = first_page.gbp_search_query + logger.info( + "Fast listing first page parsed: scope=%s result_count=%s page_size=%s total_pages=%s first_page_vehicles=%s", + scope_path, + first_page.result_count, + page_size, + total_pages, + len(first_page.vehicles), + ) for page_number in range(2, total_pages + 1): + if page_number == 2 or page_number % 25 == 0 or page_number == total_pages: + logger.info("Fast listing page fetch progress: page=%s/%s", page_number, total_pages) page_html = self._fetch_listing_page(scope_path, gbp_search_query, page_number, page_size) parsed_page = parse_listing_page(page_html) gbp_search_query = parsed_page.gbp_search_query diff --git a/iaai_scraper/core/logs.py b/iaai_scraper/core/logs.py index b1dfbff..bcefd9d 100644 --- a/iaai_scraper/core/logs.py +++ b/iaai_scraper/core/logs.py @@ -18,8 +18,9 @@ def set_trace_id(trace_id: str) -> None: def setup_logging(level: str = "INFO", log_file: str | None = None) -> None: - # stderr — Docker и Celery prefork корректно его подхватывают. - handlers: list[logging.Handler] = [logging.StreamHandler(sys.stderr)] + # stdout — основной поток для `docker logs`, Docker Desktop и compose logs. + # stderr в PowerShell часто выглядит как NativeCommandError, хотя это обычные логи. + handlers: list[logging.Handler] = [logging.StreamHandler(sys.stdout)] if log_file: handlers.append(logging.FileHandler(log_file, encoding="utf-8")) trace_filter = TraceIdFilter() diff --git a/iaai_scraper/fast_sync.py b/iaai_scraper/fast_sync.py index 8f0281c..19eced0 100644 --- a/iaai_scraper/fast_sync.py +++ b/iaai_scraper/fast_sync.py @@ -183,6 +183,7 @@ class FastSyncEngine: prepared_rows: list[CarRecord] = [] db_processed = 0 retry_candidates: list[FastListingVehicle] = [] + transient_failure_log_count = 0 started_details = time.perf_counter() with concurrent.futures.ThreadPoolExecutor(max_workers=min(self.fetch_concurrency, len(candidates))) as executor: future_to_vehicle = { @@ -232,7 +233,13 @@ class FastSyncEngine: stats.protection_events += 1 if self._is_transient_detail_error(exc): retry_candidates.append(vehicle) - logger.info("Fast detail transient failure queued for retry inventory_id=%s: %s", vehicle.inventory_id, exc) + transient_failure_log_count += 1 + if transient_failure_log_count % 100 == 0: + logger.warning( + "Fast detail transient failures queued for retry: count=%d latest_inventory_id=%s", + transient_failure_log_count, + vehicle.inventory_id, + ) else: logger.exception("Fast detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc) diff --git a/iaai_scraper/scraper.py b/iaai_scraper/scraper.py index cf7a8a4..e4718b0 100644 --- a/iaai_scraper/scraper.py +++ b/iaai_scraper/scraper.py @@ -2234,6 +2234,8 @@ class IAAIScraper: only_new = rc.only_new if lane == "iaai_cars" and rc.lane is not None: lane = rc.lane + if listing_url is None and rc.name: + listing_url = rc.name trace_id = self._new_trace_id("sync-listing") started_at = time.perf_counter() @@ -2260,6 +2262,8 @@ class IAAIScraper: year_min, year_max, ) + if listing_url: + logger.info("sync_listing uses explicit listing_url from runtime/args: %s", listing_url) fast_engine = FastSyncEngine( client=self.fast_client, mapper=self.fast_mapper, diff --git a/iaai_scraper/worker/celery_app.py b/iaai_scraper/worker/celery_app.py index 227063e..eda34ab 100644 --- a/iaai_scraper/worker/celery_app.py +++ b/iaai_scraper/worker/celery_app.py @@ -16,6 +16,9 @@ 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:" +SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing" +SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done" +SYNC_LAST_COMPLETED_AT_KEY = "iaai:state:sync_listing_last_completed_at" def _env_bool(name: str, default: bool) -> bool: @@ -43,6 +46,26 @@ def _has_fresh_active_progress(redis_client: Redis, *, max_age_seconds: int = 18 return False +def _seconds_until_next_allowed_sync(redis_client: Redis) -> int: + """Запрещает новый автозапуск раньше чем через интервал beat после завершения полного run.""" + min_interval = max(0, int(settings.celery.beat_sync_interval_minutes * 60)) + if min_interval <= 0: + return 0 + try: + full_done = str(redis_client.get(SYNC_FULL_SCAN_DONE_KEY) or "").strip().lower() in {"1", "true", "yes", "on"} + if not full_done: + return 0 + completed_raw = redis_client.get(SYNC_LAST_COMPLETED_AT_KEY) + if not completed_raw: + return 0 + completed_at = int(float(completed_raw)) + except Exception: + logger.warning("Failed to inspect last sync completion timestamp", exc_info=True) + return 0 + elapsed = int(time.time()) - completed_at + return max(0, min_interval - elapsed) + + @celery_setup_logging.connect def _configure_logging(loglevel=None, **kwargs): # Перехватываем логирование Celery и пишем только в stderr (Docker logs). @@ -147,7 +170,15 @@ def _on_worker_ready(**kwargs): ) has_fresh_progress = _has_fresh_active_progress(redis_client) - for stale_key in ("iaai:locks:sync_listing",): + next_allowed_delay = _seconds_until_next_allowed_sync(redis_client) + if next_allowed_delay > 0: + logger.info( + "Worker ready: last full sync finished recently; next auto sync allowed in %ss, skip startup dispatch", + next_allowed_delay, + ) + return + + for stale_key in (SYNC_LISTING_LOCK_KEY,): try: ttl = redis_client.ttl(stale_key) if ttl is not None and ttl != -2 and not has_fresh_progress: diff --git a/iaai_scraper/worker/fast_sync.py b/iaai_scraper/worker/fast_sync.py index 92911ce..f421e8a 100644 --- a/iaai_scraper/worker/fast_sync.py +++ b/iaai_scraper/worker/fast_sync.py @@ -183,6 +183,7 @@ class FastSyncEngine: prepared_rows: list[CarRecord] = [] db_processed = 0 retry_candidates: list[FastListingVehicle] = [] + transient_failure_log_count = 0 started_details = time.perf_counter() with concurrent.futures.ThreadPoolExecutor(max_workers=min(self.fetch_concurrency, len(candidates))) as executor: future_to_vehicle = { @@ -232,7 +233,13 @@ class FastSyncEngine: stats.protection_events += 1 if self._is_transient_detail_error(exc): retry_candidates.append(vehicle) - logger.info("Fast detail transient failure queued for retry inventory_id=%s: %s", vehicle.inventory_id, exc) + transient_failure_log_count += 1 + if transient_failure_log_count % 100 == 0: + logger.warning( + "Fast detail transient failures queued for retry: count=%d latest_inventory_id=%s", + transient_failure_log_count, + vehicle.inventory_id, + ) else: logger.exception("Fast detail parse failed inventory_id=%s: %s", vehicle.inventory_id, exc) diff --git a/iaai_scraper/worker/tasks.py b/iaai_scraper/worker/tasks.py index 2dd5b46..740de2a 100644 --- a/iaai_scraper/worker/tasks.py +++ b/iaai_scraper/worker/tasks.py @@ -48,6 +48,7 @@ HOURLY_FAILURE_STREAK_LIMIT = 3 HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending" SYNC_LISTING_TASK_NAME = "iaai.sync_cars_feed" +SYNC_LAST_COMPLETED_AT_KEY = "iaai:state:sync_listing_last_completed_at" 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" @@ -867,6 +868,31 @@ def _clear_followup_pending(redis_client: Redis) -> None: logger.warning("Failed to clear follow-up pending flag", exc_info=True) +def _mark_sync_completed(redis_client: Redis) -> None: + try: + redis_client.set(SYNC_LAST_COMPLETED_AT_KEY, str(int(time.time())), ex=7 * 24 * 60 * 60) + except Exception: + logger.warning("Failed to mark sync completion timestamp", exc_info=True) + + +def _seconds_until_next_allowed_sync(redis_client: Redis, settings: Settings) -> int: + min_interval = max(0, int(settings.celery.beat_sync_interval_minutes * 60)) + if min_interval <= 0: + return 0 + try: + if not _is_full_scan_done(redis_client): + return 0 + completed_raw = redis_client.get(SYNC_LAST_COMPLETED_AT_KEY) + if not completed_raw: + return 0 + completed_at = int(float(completed_raw)) + except Exception: + logger.warning("Failed to inspect next allowed sync time", exc_info=True) + return 0 + elapsed = int(time.time()) - completed_at + return max(0, min_interval - elapsed) + + def _bump_bootstrap_failure_streak( redis_client: Redis, *, @@ -1310,6 +1336,11 @@ def sync_listing_task( watchdog_stop: Event | None = None watchdog_thread: Thread | None = None force_bootstrap_full_scan = False + settings = Settings() + followup_min_delay_seconds = max( + 0, + int(os.getenv("IAAI_SYNC_FOLLOWUP_MIN_DELAY_SECONDS", str(settings.celery.beat_sync_interval_minutes * 60))), + ) def _enqueue_bootstrap_followup( reason: str, @@ -1317,6 +1348,7 @@ def sync_listing_task( *, count_as_failure: bool = False, ) -> None: + delay_seconds = max(int(delay_seconds), followup_min_delay_seconds) flag_ttl = max(lock_ttl, delay_seconds + 300) if not _try_set_followup_pending(redis_client, ttl_seconds=flag_ttl): logger.info( @@ -1391,7 +1423,18 @@ def sync_listing_task( } try: - settings = Settings() + next_allowed_delay = _seconds_until_next_allowed_sync(redis_client, settings) + if next_allowed_delay > 0: + logger.info( + "sync_listing_task skipped: previous full run finished recently; next run allowed in %ss", + next_allowed_delay, + ) + return { + "status": "skipped", + "reason": "next_sync_not_due_yet", + "retry_after_seconds": next_allowed_delay, + "task_id": task_id, + } # Любой реально стартовавший sync_listing снимает pending-флаг followup, # чтобы watchdog/continuation могли корректно планировать следующий run @@ -1832,6 +1875,7 @@ def sync_listing_task( summary["cars_failed"], summary["failures_count"], ) + _mark_sync_completed(redis_client) return summary except SoftTimeLimitExceeded: diff --git a/runtime_config.json b/runtime_config.json index bb43bcb..7d09878 100644 --- a/runtime_config.json +++ b/runtime_config.json @@ -1,6 +1,6 @@ { "sync": { - "name": null, + "name": "https://www.iaai.com/Search?url=B%2bU066dM8%2flZtRvnzTwIEi8Pib9%2fM2fLpCTQDIgQRm4%3d", "ids_initial_size": null, "ids_next_size": null, "ids_max_pages": null,