Compare commits
3 Commits
f221aebef0
...
6b69b8c002
| Author | SHA1 | Date | |
|---|---|---|---|
| 6b69b8c002 | |||
|
|
e726461e75 | ||
|
|
6aba4f0dbd |
@@ -45,6 +45,7 @@ _PAGE_COLLECT_TIMEOUT_S = 90
|
||||
_PAGE_NEXT_TIMEOUT_S = 60
|
||||
# Ротация контекста выключена по умолчанию.
|
||||
_CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0)
|
||||
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT = int(os.environ.get("IAAI_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT", "3") or 3)
|
||||
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120
|
||||
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2
|
||||
|
||||
@@ -594,6 +595,7 @@ class IAAIScraper:
|
||||
failures: list[dict[str, str]] = []
|
||||
early_stopped = False
|
||||
truncated_by_time_budget = False
|
||||
listing_recovery_attempts = 0
|
||||
|
||||
known_origin_ids: set[str] | None = None
|
||||
threshold = self.settings.listing.early_stop_threshold
|
||||
@@ -725,6 +727,21 @@ class IAAIScraper:
|
||||
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
|
||||
return True
|
||||
|
||||
def _consume_listing_recovery_attempt(*, page_number: int, reason: str) -> None:
|
||||
nonlocal listing_recovery_attempts
|
||||
listing_recovery_attempts += 1
|
||||
logger.warning(
|
||||
"Listing recovery attempt %d/%d at page %d (%s)",
|
||||
listing_recovery_attempts,
|
||||
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT,
|
||||
page_number,
|
||||
reason,
|
||||
)
|
||||
if listing_recovery_attempts > _MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT:
|
||||
raise ListingResumeError(
|
||||
f"Listing recovery exhausted at page {page_number}: {reason}"
|
||||
)
|
||||
|
||||
page = None
|
||||
try:
|
||||
page, applied_filters = self._open_listing_for_stream(
|
||||
@@ -763,6 +780,10 @@ class IAAIScraper:
|
||||
pass
|
||||
page = None
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number,
|
||||
reason="context_rotation",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number,
|
||||
make=make,
|
||||
@@ -794,6 +815,10 @@ class IAAIScraper:
|
||||
pass
|
||||
page = None
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number,
|
||||
reason=f"collect_failed:{type(op_exc).__name__}",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number,
|
||||
make=make,
|
||||
@@ -839,6 +864,10 @@ class IAAIScraper:
|
||||
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
|
||||
page.close()
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number,
|
||||
reason="empty_page_after_retry",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number,
|
||||
make=make,
|
||||
@@ -1069,6 +1098,10 @@ class IAAIScraper:
|
||||
pass
|
||||
page = None
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number + 1,
|
||||
reason=f"next_page_failed:{type(nav_exc).__name__}",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number + 1,
|
||||
make=make,
|
||||
|
||||
@@ -12,6 +12,8 @@ logger = logging.getLogger("iaai_scraper.worker.self_heal")
|
||||
|
||||
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||
SELF_HEAL_RESTART_LOCK_KEY = "iaai:state:self_heal_restart_in_progress"
|
||||
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
||||
SIGKILL_FALLBACK = getattr(signal, "SIGKILL", signal.SIGTERM)
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
@@ -76,6 +78,22 @@ def _read_last_progress_ts(redis_client: Redis) -> int | None:
|
||||
return max_ts or None
|
||||
|
||||
|
||||
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)
|
||||
has_lock = 1 if redis_client.get(SYNC_LISTING_LOCK_KEY) else 0
|
||||
has_task_progress = 0
|
||||
for _ in redis_client.scan_iter(match="iaai:state:task_progress:*"):
|
||||
has_task_progress = 1
|
||||
break
|
||||
flags = {
|
||||
"queue_len": queue_len,
|
||||
"has_lock": has_lock,
|
||||
"has_task_progress": has_task_progress,
|
||||
}
|
||||
return (queue_len > 0 or has_lock == 1 or has_task_progress == 1), flags
|
||||
|
||||
|
||||
def _kill_worker_process() -> None:
|
||||
pid_file = "/tmp/celery-worker.pid"
|
||||
pid: int | None = None
|
||||
@@ -101,7 +119,7 @@ def _kill_worker_process() -> None:
|
||||
# Если процесс ещё жив — принудительно убиваем.
|
||||
os.kill(pid, 0)
|
||||
logger.error("Self-heal: worker pid=%s did not stop after SIGTERM; sending SIGKILL", pid)
|
||||
os.kill(pid, signal.SIGKILL)
|
||||
os.kill(pid, SIGKILL_FALLBACK)
|
||||
except ProcessLookupError:
|
||||
pass
|
||||
except Exception:
|
||||
@@ -137,8 +155,8 @@ def main() -> None:
|
||||
redis_client = _get_redis()
|
||||
redis_client.ping()
|
||||
|
||||
queue_len = _safe_int(redis_client.llen(queue_name), 0)
|
||||
if queue_len <= 0:
|
||||
has_inflight, inflight = _has_inflight_work(redis_client, queue_name)
|
||||
if not has_inflight:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
|
||||
@@ -151,8 +169,8 @@ def main() -> None:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
logger.warning(
|
||||
"Self-heal: queue=%d but no progress timestamp found after startup grace",
|
||||
queue_len,
|
||||
"Self-heal: inflight=%s but no progress timestamp found after startup grace",
|
||||
inflight,
|
||||
)
|
||||
age = stall_seconds + 1
|
||||
|
||||
@@ -174,8 +192,8 @@ def main() -> None:
|
||||
continue
|
||||
|
||||
logger.error(
|
||||
"Self-heal: detected global stall (queue=%d, progress_age=%ss > %ss). Restarting worker process...",
|
||||
queue_len,
|
||||
"Self-heal: detected global stall (inflight=%s, progress_age=%ss > %ss). Restarting worker process...",
|
||||
inflight,
|
||||
age,
|
||||
stall_seconds,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user