refactor mobile.de parser, fix country mapping, update README
This commit is contained in:
105
mobilede_scraper/worker/progress.py
Normal file
105
mobilede_scraper/worker/progress.py
Normal file
@@ -0,0 +1,105 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import time
|
||||
|
||||
from redis import Redis
|
||||
|
||||
from .constants import (
|
||||
DB_PROGRESS_STAGES,
|
||||
GLOBAL_DB_PROGRESS_TS_KEY,
|
||||
GLOBAL_PROGRESS_TS_KEY,
|
||||
MOBILEDE_SEGMENT_FOLLOWUP_PENDING_KEY_FMT,
|
||||
MOBILEDE_SEGMENT_LOCK_KEY_FMT,
|
||||
STALL_WATCHDOG_DETAIL_GRACE_SECONDS,
|
||||
STALL_WATCHDOG_LONG_RUNNING_STAGES,
|
||||
STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS,
|
||||
STALL_WATCHDOG_NAVIGATION_STAGES,
|
||||
TASK_PROGRESS_KEY_FMT,
|
||||
)
|
||||
|
||||
logger = logging.getLogger("mobilede_scraper.worker.progress")
|
||||
|
||||
|
||||
def _safe_int(value) -> int | None:
|
||||
try:
|
||||
if value is None:
|
||||
return None
|
||||
return int(value)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _task_progress_key(task_id: str) -> str:
|
||||
return TASK_PROGRESS_KEY_FMT.format(task_id=task_id)
|
||||
|
||||
|
||||
def _mobilede_segment_lock_key(segment_key: str) -> str:
|
||||
return MOBILEDE_SEGMENT_LOCK_KEY_FMT.format(segment_key=segment_key)
|
||||
|
||||
|
||||
def _mobilede_followup_pending_key(segment_key: str) -> str:
|
||||
return MOBILEDE_SEGMENT_FOLLOWUP_PENDING_KEY_FMT.format(segment_key=segment_key)
|
||||
|
||||
|
||||
def _update_task_progress(
|
||||
redis_client: Redis,
|
||||
*,
|
||||
task_id: str,
|
||||
stage: str,
|
||||
ttl_seconds: int,
|
||||
**payload,
|
||||
) -> 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 not in DB_PROGRESS_STAGES 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,
|
||||
"ts": now_ts,
|
||||
**payload,
|
||||
}
|
||||
ttl = max(60, int(ttl_seconds))
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.set(
|
||||
_task_progress_key(task_id),
|
||||
json.dumps(data, ensure_ascii=False),
|
||||
ex=ttl,
|
||||
)
|
||||
# Глобальный маркер активности для внешнего guard-процесса.
|
||||
# Нужен, чтобы контейнер мог самовосстанавливаться при полном зависании воркера
|
||||
# (когда PID жив, но прогресс по задачам не двигается).
|
||||
pipe.set(GLOBAL_PROGRESS_TS_KEY, str(now_ts), ex=max(ttl, 7 * 24 * 60 * 60))
|
||||
if stage in DB_PROGRESS_STAGES:
|
||||
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)
|
||||
|
||||
|
||||
def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
|
||||
try:
|
||||
redis_client.delete(_task_progress_key(task_id))
|
||||
except Exception:
|
||||
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
|
||||
|
||||
|
||||
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)
|
||||
Reference in New Issue
Block a user