Prepare mobile de parser release
This commit is contained in:
3
iaai_scraper/worker/__init__.py
Normal file
3
iaai_scraper/worker/__init__.py
Normal file
@@ -0,0 +1,3 @@
|
||||
from .celery_app import celery_app
|
||||
|
||||
__all__ = ["celery_app"]
|
||||
207
iaai_scraper/worker/celery_app.py
Normal file
207
iaai_scraper/worker/celery_app.py
Normal file
@@ -0,0 +1,207 @@
|
||||
# Инициализация 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
|
||||
from redis import Redis
|
||||
|
||||
from ..core.config import settings
|
||||
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"
|
||||
MOBILEDE_SYNC_QUEUE = "mobilede_sync"
|
||||
PROGRESS_KEY_PREFIX = "iaai:state:task_progress:"
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
raw = os.getenv(name, "true" if default else "false").strip().lower()
|
||||
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).
|
||||
level = settings.log_level if settings.log_level else "INFO"
|
||||
setup_logging(level, None)
|
||||
|
||||
|
||||
@worker_process_init.connect
|
||||
def _on_worker_process_init(**kwargs):
|
||||
# Повторно настраиваем логирование в каждом дочернем prefork-процессе,
|
||||
# чтобы StreamHandler(stderr) корректно работал после fork.
|
||||
level = settings.log_level if settings.log_level else "INFO"
|
||||
setup_logging(level, None)
|
||||
|
||||
|
||||
def _broker_url() -> str:
|
||||
return settings.celery.broker_url or settings.redis.url
|
||||
|
||||
|
||||
def _result_backend() -> str:
|
||||
return settings.celery.result_backend or settings.redis.url
|
||||
|
||||
|
||||
celery_app = Celery(
|
||||
"iaai_scraper",
|
||||
broker=_broker_url(),
|
||||
backend=_result_backend(),
|
||||
)
|
||||
|
||||
# Auto-clamp: если hard limit слишком далёк от soft (> soft + 120),
|
||||
# ограничиваем, чтобы зависший worker не жил вечно.
|
||||
_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.info(
|
||||
"CELERY_TASK_TIME_LIMIT=%d too far from CELERY_TASK_SOFT_TIME_LIMIT=%d; "
|
||||
"clamping hard limit to %d",
|
||||
_hard, _soft, _max_hard,
|
||||
)
|
||||
_hard = _max_hard
|
||||
|
||||
celery_app.conf.update(
|
||||
task_serializer="json",
|
||||
accept_content=["json"],
|
||||
result_serializer="json",
|
||||
timezone="UTC",
|
||||
enable_utc=True,
|
||||
task_soft_time_limit=_soft,
|
||||
task_time_limit=_hard,
|
||||
task_acks_late=True,
|
||||
task_reject_on_worker_lost=True,
|
||||
task_track_started=True,
|
||||
worker_concurrency=settings.celery.worker_concurrency,
|
||||
worker_max_tasks_per_child=settings.celery.worker_max_tasks_per_child,
|
||||
worker_pool=settings.celery.worker_pool,
|
||||
worker_prefetch_multiplier=1,
|
||||
broker_connection_retry_on_startup=True,
|
||||
broker_transport_options={
|
||||
"visibility_timeout": settings.celery.broker_visibility_timeout,
|
||||
},
|
||||
result_expires=86400,
|
||||
worker_redirect_stdouts=False,
|
||||
worker_hijack_root_logger=False,
|
||||
beat_schedule={
|
||||
"periodic-mobilede-sync-search": {
|
||||
"task": "mobilede.sync_runtime_segments",
|
||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"args": (),
|
||||
"kwargs": {
|
||||
"delay_seconds": float(os.getenv("MOBILEDE_REQUEST_DELAY_SECONDS", "0.7")),
|
||||
"use_cursor": _env_bool("MOBILEDE_CURSOR_ENABLED", True),
|
||||
"continuous": _env_bool("MOBILEDE_CONTINUOUS_SYNC_ENABLED", True),
|
||||
},
|
||||
"options": {
|
||||
"queue": MOBILEDE_SYNC_QUEUE,
|
||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
},
|
||||
}
|
||||
},
|
||||
task_routes={
|
||||
"mobilede.sync_runtime_segments": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
"mobilede.sync_search": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
"mobilede.sync_detail": {"queue": MOBILEDE_SYNC_QUEUE},
|
||||
"iaai.sync_cars_feed": {"queue": IAAI_SYNC_QUEUE},
|
||||
"iaai_scraper.worker.tasks.*": {"queue": IAAI_SYNC_QUEUE},
|
||||
},
|
||||
)
|
||||
|
||||
celery_app.autodiscover_tasks(["iaai_scraper.worker"])
|
||||
|
||||
|
||||
@worker_ready.connect
|
||||
def _on_worker_ready(**kwargs):
|
||||
"""При старте worker отправляем первый sync_listing, если очередь пуста."""
|
||||
if not _env_bool("IAAI_STARTUP_SYNC_ENABLED", True):
|
||||
logger.info("Worker ready: startup sync dispatch disabled by IAAI_STARTUP_SYNC_ENABLED")
|
||||
return
|
||||
|
||||
redis_client = None
|
||||
try:
|
||||
redis_client = Redis.from_url(
|
||||
settings.redis.url,
|
||||
decode_responses=True,
|
||||
socket_connect_timeout=settings.redis.socket_connect_timeout_seconds,
|
||||
socket_timeout=settings.redis.socket_timeout_seconds,
|
||||
health_check_interval=settings.redis.health_check_interval_seconds,
|
||||
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 and not has_fresh_progress:
|
||||
redis_client.delete(stale_key)
|
||||
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 inspect stale lock %s on startup", stale_key, exc_info=True)
|
||||
|
||||
try:
|
||||
queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) or 0)
|
||||
except Exception:
|
||||
queue_len = 0
|
||||
if queue_len > 0:
|
||||
logger.info("Worker ready: iaai_sync queue already has %d task(s); skip startup dispatch", queue_len)
|
||||
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
|
||||
finally:
|
||||
if redis_client is not None:
|
||||
try:
|
||||
redis_client.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if not should_dispatch:
|
||||
logger.info("Worker ready immediate sync already dispatched recently; skipping duplicate enqueue")
|
||||
return
|
||||
|
||||
logger.info("Worker ready — dispatching initial mobile.de sync_search task")
|
||||
celery_app.send_task(
|
||||
"mobilede.sync_runtime_segments",
|
||||
kwargs={
|
||||
"delay_seconds": float(os.getenv("MOBILEDE_REQUEST_DELAY_SECONDS", "0.7")),
|
||||
"use_cursor": _env_bool("MOBILEDE_CURSOR_ENABLED", True),
|
||||
"continuous": _env_bool("MOBILEDE_CONTINUOUS_SYNC_ENABLED", True),
|
||||
},
|
||||
queue=MOBILEDE_SYNC_QUEUE,
|
||||
expires=settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
)
|
||||
276
iaai_scraper/worker/self_heal.py
Normal file
276
iaai_scraper/worker/self_heal.py
Normal file
@@ -0,0 +1,276 @@
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import random
|
||||
import signal
|
||||
import time
|
||||
|
||||
from redis import Redis
|
||||
|
||||
|
||||
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)
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
value = os.getenv(name)
|
||||
if value is None:
|
||||
return default
|
||||
return value.strip().lower() in {"1", "true", "yes", "on"}
|
||||
|
||||
|
||||
def _env_int(name: str, default: int) -> int:
|
||||
value = os.getenv(name)
|
||||
if value is None:
|
||||
return default
|
||||
try:
|
||||
return int(value.strip())
|
||||
except Exception:
|
||||
return default
|
||||
|
||||
|
||||
def _get_redis() -> Redis:
|
||||
url = os.getenv("IAAI_REDIS_URL", "redis://redis:6379/0")
|
||||
return Redis.from_url(
|
||||
url,
|
||||
decode_responses=True,
|
||||
socket_connect_timeout=5.0,
|
||||
socket_timeout=10.0,
|
||||
health_check_interval=30,
|
||||
retry_on_timeout=True,
|
||||
)
|
||||
|
||||
|
||||
def _safe_int(value: str | None, default: int = 0) -> int:
|
||||
if value is None:
|
||||
return default
|
||||
try:
|
||||
return int(str(value).strip())
|
||||
except Exception:
|
||||
return default
|
||||
|
||||
|
||||
def _read_last_progress_ts(redis_client: Redis) -> int | None:
|
||||
raw = redis_client.get(GLOBAL_PROGRESS_TS_KEY)
|
||||
if raw:
|
||||
ts = _safe_int(raw)
|
||||
if ts > 0:
|
||||
return ts
|
||||
|
||||
# Fallback: если глобальный ключ не найден, берём max(ts) из task_progress:*.
|
||||
# Это дороже, но выполняется только при отсутствии основного маркера.
|
||||
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)
|
||||
ts = _safe_int(data.get("ts"), 0)
|
||||
if ts > max_ts:
|
||||
max_ts = ts
|
||||
except Exception:
|
||||
continue
|
||||
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)
|
||||
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
|
||||
try:
|
||||
with open(pid_file, "r", encoding="utf-8") as f:
|
||||
pid = int(f.read().strip())
|
||||
except Exception:
|
||||
pid = None
|
||||
|
||||
if not pid:
|
||||
logger.error("Self-heal: failed to read worker pid from %s", pid_file)
|
||||
return
|
||||
|
||||
logger.error("Self-heal: terminating stuck worker process pid=%s", pid)
|
||||
try:
|
||||
os.kill(pid, signal.SIGTERM)
|
||||
except Exception:
|
||||
logger.exception("Self-heal: failed to send SIGTERM to pid=%s", pid)
|
||||
return
|
||||
|
||||
time.sleep(20)
|
||||
try:
|
||||
# Если процесс ещё жив — принудительно убиваем.
|
||||
os.kill(pid, 0)
|
||||
logger.error("Self-heal: worker pid=%s did not stop after SIGTERM; sending SIGKILL", pid)
|
||||
os.kill(pid, SIGKILL_FALLBACK)
|
||||
except ProcessLookupError:
|
||||
pass
|
||||
except Exception:
|
||||
logger.exception("Self-heal: failed to send SIGKILL to pid=%s", pid)
|
||||
|
||||
|
||||
def main() -> None:
|
||||
if not _env_bool("IAAI_SELF_HEAL_ENABLED", True):
|
||||
logger.info("Self-heal watchdog disabled via IAAI_SELF_HEAL_ENABLED")
|
||||
return
|
||||
|
||||
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.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,
|
||||
)
|
||||
|
||||
started_at = time.time()
|
||||
redis_client: Redis | None = None
|
||||
|
||||
while True:
|
||||
try:
|
||||
if redis_client is None:
|
||||
redis_client = _get_redis()
|
||||
redis_client.ping()
|
||||
|
||||
has_inflight, inflight = _has_inflight_work(redis_client, queue_name)
|
||||
if not has_inflight:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
|
||||
last_progress_ts = _read_last_progress_ts(redis_client)
|
||||
now_ts = int(time.time())
|
||||
age = None if last_progress_ts is None else max(0, now_ts - int(last_progress_ts))
|
||||
|
||||
if age is None:
|
||||
if now_ts - int(started_at) < startup_grace:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
logger.warning(
|
||||
"Self-heal: inflight=%s but no progress timestamp found after startup grace",
|
||||
inflight,
|
||||
)
|
||||
age = stall_seconds + 1
|
||||
|
||||
restart_reason = f"progress_age={age}s > {stall_seconds}s"
|
||||
db_idle_restart = False
|
||||
if age <= stall_seconds:
|
||||
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(
|
||||
redis_client.set(
|
||||
SELF_HEAL_RESTART_LOCK_KEY,
|
||||
str(now_ts),
|
||||
nx=True,
|
||||
ex=restart_cooldown,
|
||||
)
|
||||
)
|
||||
if not acquired:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
|
||||
logger.error(
|
||||
"Self-heal: detected stall (inflight=%s, %s). Restarting worker process...",
|
||||
inflight,
|
||||
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))
|
||||
_kill_worker_process()
|
||||
# После kill pid1 контейнер будет перезапущен Docker restart-policy.
|
||||
# На случай неуспеха не молотим цикл.
|
||||
time.sleep(check_interval)
|
||||
|
||||
except Exception:
|
||||
logger.exception("Self-heal watchdog iteration failed")
|
||||
redis_client = None
|
||||
time.sleep(check_interval)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
logging.basicConfig(
|
||||
level=os.getenv("IAAI_LOG_LEVEL", "INFO"),
|
||||
format="%(asctime)s | %(levelname)s | %(name)s | %(message)s",
|
||||
)
|
||||
main()
|
||||
3358
iaai_scraper/worker/tasks.py
Normal file
3358
iaai_scraper/worker/tasks.py
Normal file
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user