All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 12s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m58s
compute_next_run_at держит суточную гранулярность (interval_days, минимум 1),
подчасовой такт делается хуком reschedule_after_minutes как post_claim. У
avito_detail_backfill он есть (180 мин), у yandex_detail_backfill и
cian_detail_backfill не было — отсюда один прогон в сутки.
Цена простоя по замеру прода 31.08:
Яндекс: 375-450 карточек за прогон, блоков НОЛЬ за неделю, очередь 11110
→ 25 суток при нынешнем такте
Циан: блоков ноль, очередь 20501
Такты разные, и это не произвол:
yandex — 180 мин (8 прогонов/сутки). Ходит через resolve_proxy_url: берёт
URL узла, но НЕ лизует его, поэтому чужие прогоны не блокирует.
cian — 360 мин (4 прогона/сутки). Ходит через BrowserFetcher и ДЕРЖИТ
lease весь прогон, то есть отнимает узел у Авито и Домклика. Пул
дефицитен (#2638), поэтому осторожнее.
Асимметрия зафиксирована комментарием у обоих хендлеров и в докстринге
миграции — иначе следующий читатель выровняет интервалы и сожжёт пул. Тест
test_cian_default_interval_is_360_minutes_not_180 ассертит именно неравенство,
чтобы выравнивание без замера покраснело.
283_scrape_schedules_cadence_yandex_cian_detail.sql — идемпотентный
UPDATE ... SET default_params = default_params || jsonb, остальные ключи
параметров не трогает.
Значения 180/360 — консервативная отправная точка по аналогии с Авито, а не
найденный оптимум: двигать вниз только по замеру нескольких суток, глядя и на
свипы тоже (тот же довод, что в комментарии у avito_detail_backfill).
Тесты (5): наличие post_claim у обоих, дефолтные интервалы, переопределение
через params. Прогон: 397 passed, 1 skipped, ruff чист.
823 lines
45 KiB
Python
823 lines
45 KiB
Python
"""Продуктовые Handler'ы для kit-scheduler (#2192, strangler-миграция).
|
||
|
||
`build_product_handlers(ctx)` собирает реестр НЕ-sweep source'ов, тело которых
|
||
осталось в `app` (rosreestr_dkp / sber_index / deactivate_stale_* / *_backfill /
|
||
cadastral_geo_match / house_* / proxy_healthcheck / …). Kit-native sweep-оркестраторы
|
||
(avito/yandex/cian/domclick full-load) регистрируются самим `build_registry` через
|
||
`_default_kit_handlers` — здесь их НЕТ. Исключение: `domclick_city_sweep` (#3264) —
|
||
kit-native тело осталось прежним (`run_domclick_city_sweep`), но здесь переопределён
|
||
Handler, который ДО вызова читает БД за куки сессии (kit не имеет права на app.*/БД,
|
||
см. докстринг `_job_domclick_city_sweep` ниже); `build_registry` явно допускает такое
|
||
переопределение продуктовым ключом.
|
||
|
||
Дизайн-инвариант: НЕ дублируем логику. Каждый `job` переиспользует существующую боевую
|
||
джоб-функцию (`app.services.scheduler.import_rosreestr_dkp`, `app.tasks.*`,
|
||
`app.services.*`) — то же тело, что крутит боевой `trigger_*_run`. Kit `_dispatch` уже
|
||
делает claim + fresh-session + spawn + close, поэтому `job` содержит ТОЛЬКО работу
|
||
(без claim/spawn-boilerplate). Run-lifecycle (mark_done/mark_failed/heartbeat) для
|
||
джоб, что не владеют им сами, делается через инжектированный `ctx.runs`.
|
||
|
||
Ship-dark: модуль импортируется лишь когда `settings.use_kit_scheduler` True
|
||
(scheduler_main._run_kit_scheduler). Боевой `app.services.scheduler` не тронут.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
from datetime import UTC, datetime, timedelta
|
||
from typing import TYPE_CHECKING, Any
|
||
|
||
from scraper_kit.orchestration import runs as kit_runs
|
||
from scraper_kit.orchestration.scheduler import (
|
||
Handler,
|
||
reschedule_after_minutes,
|
||
)
|
||
from scraper_kit.orchestration.scheduler import (
|
||
_defer_next_run_at as kit_defer_next_run_at,
|
||
)
|
||
from scraper_kit.orchestration.scheduler import (
|
||
_pick_resume as kit_pick_resume,
|
||
)
|
||
|
||
if TYPE_CHECKING:
|
||
from scraper_kit.orchestration.scheduler import SchedulerContext
|
||
from sqlalchemy.orm import Session
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
# ── cian_history_backfill — cookie-gated backfill ────────────────────────────
|
||
# Машиночитаемые причины пропуска (#2658) — пишутся в scrape_runs.error строки
|
||
# со status='skipped'. Отделены от kit-причин (already_running и т.п.): по слагу
|
||
# видно, встал ли сбор из-за кук или из-за конкурентного прогона.
|
||
SKIP_CIAN_COOKIES_MISSING = "cian_cookies_missing"
|
||
SKIP_CIAN_COOKIES_EXPIRED = "cian_cookies_expired"
|
||
SKIP_CIAN_COOKIES_INVALID = "cian_cookies_invalid"
|
||
|
||
|
||
def _alert_cian_cookies(source: str, detail: str) -> None:
|
||
"""Громкий алерт «сбор встал из-за кук» — logger.error, НЕ capture_message(warning).
|
||
|
||
В scraper-контейнере GlitchTip поднят с LoggingIntegration(event_level=ERROR)
|
||
(scheduler_main.py) — ERROR-запись сама становится событием, а прежний
|
||
`capture_message(..., level="warning")` до этого уровня не дотягивал (и стоял в
|
||
недостижимой ветке, см. докстринг _cian_pre_claim). Заодно причина остаётся в
|
||
docker-логах и в строке scrape_runs, которая переживает редеплой.
|
||
"""
|
||
logger.error(
|
||
"scheduler: %s пропущен — %s. Перезалейте куки Циана через админку "
|
||
"(до этого backfill истории стоит)",
|
||
source,
|
||
detail,
|
||
)
|
||
|
||
|
||
async def _cian_pre_claim(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) -> bool:
|
||
"""Pre-claim gate: проверить наличие/валидность cian-cookies ДО claim (#1522).
|
||
|
||
Cookies отсутствуют/протухли → пишем строку прогона status='skipped' с причиной,
|
||
двигаем next_run_at на следующее окно и skip (иначе get_due_schedules переотбирает
|
||
schedule каждые 60с и verify_session долбит Cian круглосуточно).
|
||
|
||
#2658 — что было не так. Первая ветка (load_session вернул None) молчала: warning в
|
||
docker-лог, сдвиг next_run_at, `return False`. Ни строки в scrape_runs, ни изменения
|
||
last_run_at — снаружи 37 дней простоя выглядели как «всё по расписанию». Sentry-алерт
|
||
стоял во ВТОРОЙ ветке (verify_session вернул None), до которой на протухших куках
|
||
исполнение не доходит НИКОГДА: load_session сам фильтрует expires_at_estimate > NOW()
|
||
и отдаёт None ещё в первой. Теперь громко в обеих + предупреждение ЗАРАНЕЕ, пока куки
|
||
ещё валидны (COOKIE_EXPIRY_WARN_DAYS) — обновление кук ручное, ему нужен запас.
|
||
"""
|
||
from app.services.cian_session import (
|
||
COOKIE_EXPIRY_WARN_DAYS,
|
||
load_session,
|
||
session_expires_at,
|
||
verify_session,
|
||
)
|
||
|
||
source: str = schedule_row["source"]
|
||
now = datetime.now(tz=UTC)
|
||
|
||
cookies = load_session(db)
|
||
if cookies is None:
|
||
expires_at = session_expires_at(db)
|
||
if expires_at is None:
|
||
reason, detail = SKIP_CIAN_COOKIES_MISSING, "кук Циана нет в БД"
|
||
elif expires_at <= now:
|
||
reason = SKIP_CIAN_COOKIES_EXPIRED
|
||
detail = (
|
||
f"куки Циана протухли {expires_at:%Y-%m-%d} ({(now - expires_at).days} дн. назад)"
|
||
)
|
||
else:
|
||
reason = SKIP_CIAN_COOKIES_INVALID
|
||
detail = "куки Циана помечены невалидными (last_invalid_at)"
|
||
_alert_cian_cookies(source, detail)
|
||
kit_runs.mark_skipped(db, source=source, reason=reason, details=detail)
|
||
kit_defer_next_run_at(db, schedule_row)
|
||
return False
|
||
|
||
state = await verify_session(cookies)
|
||
if state is None:
|
||
# verify вернул именно None (401 / isAuthenticated=false) — куки числятся
|
||
# валидными по сроку, но Циан их не принимает. Sentinel-ответы (бан / источник
|
||
# недоступен / сменилась вёрстка) сюда НЕ попадают, они truthy — см. cian_session.
|
||
detail = "Циан не принимает куки (разлогин)"
|
||
_alert_cian_cookies(source, detail)
|
||
kit_runs.mark_skipped(db, source=source, reason=SKIP_CIAN_COOKIES_INVALID, details=detail)
|
||
kit_defer_next_run_at(db, schedule_row)
|
||
return False
|
||
|
||
# Куки рабочие — предупреждаем, пока есть время их обновить без простоя сбора.
|
||
# valid_only=True: срок ИМЕННО той записи, которую взял load_session (при нескольких
|
||
# аккаунтах свежайшая-любая может быть чужой протухшей строкой).
|
||
expires_at = session_expires_at(db, valid_only=True)
|
||
if expires_at is not None and expires_at - now <= timedelta(days=COOKIE_EXPIRY_WARN_DAYS):
|
||
logger.error(
|
||
"scheduler: куки Циана протухнут %s (осталось %.1f дн.) — обновите заранее, "
|
||
"иначе %s встанет молча",
|
||
expires_at.date().isoformat(),
|
||
(expires_at - now).total_seconds() / 86400,
|
||
source,
|
||
)
|
||
return True
|
||
|
||
|
||
async def _job_cian_history_backfill(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.services.scheduler import _execute_cian_backfill
|
||
|
||
await _execute_cian_backfill(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── rosreestr_dkp_import — sync FDW-импорт в executor ─────────────────────────
|
||
async def _job_rosreestr_dkp(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.services.scheduler import import_rosreestr_dkp
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(None, import_rosreestr_dkp, db, run_id, params)
|
||
|
||
|
||
# ── listing_source_snapshot — sync DB-snapshot в executor ────────────────────
|
||
# params прокинуты (#2607) — snapshot_listing_sources теперь читает budget_sec из
|
||
# default_params (SET LOCAL statement_timeout, см. app/tasks/listing_source_snapshot.py).
|
||
async def _job_listing_source_snapshot(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.listing_source_snapshot import snapshot_listing_sources
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(None, snapshot_listing_sources, db, run_id, params)
|
||
|
||
|
||
# ── asking_to_sold_ratio_refresh — sync re-derive в executor ─────────────────
|
||
async def _job_asking_to_sold_ratio(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.asking_to_sold_ratio import recompute_asking_to_sold_ratios
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(None, recompute_asking_to_sold_ratios, db, run_id)
|
||
|
||
|
||
# ── deal_city_price_bands_refresh — sync tier-aware re-derive в executor ──────
|
||
async def _job_deal_city_price_bands_refresh(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.deal_city_price_bands_refresh import refresh_deal_city_price_bands
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(None, refresh_deal_city_price_bands, db, run_id)
|
||
|
||
|
||
# ── refresh_search_matview — REFRESH MATVIEW CONCURRENTLY (own connection) ────
|
||
async def _job_refresh_search_matview(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.refresh_search_matview import refresh_search_matview
|
||
|
||
loop = asyncio.get_event_loop()
|
||
try:
|
||
# refresh_search_matview() держит собственное autocommit-соединение
|
||
# (REFRESH CONCURRENTLY нельзя внутри транзакции).
|
||
await loop.run_in_executor(None, refresh_search_matview)
|
||
ctx.runs.mark_done(db, run_id, {})
|
||
except Exception:
|
||
logger.exception("scheduler: refresh_search_matview crashed run_id=%d", run_id)
|
||
ctx.runs.mark_failed(db, run_id, "refresh_search_matview failed", {})
|
||
|
||
|
||
# ── yandex_address_backfill — async, owns lifecycle ──────────────────────────
|
||
async def _job_yandex_address_backfill(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.yandex_address_backfill import run_yandex_address_backfill
|
||
|
||
await run_yandex_address_backfill(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── deactivate_stale_* — sync UPDATE в executor (wildcard-семейство) ──────────
|
||
async def _job_deactivate_stale(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.core.config import settings as _settings
|
||
from app.tasks.deactivate_stale_avito import (
|
||
CAP_MULT,
|
||
DEFAULT_FLOOR_DROP_RATIO,
|
||
DEFAULT_MAX_DEACTIVATED,
|
||
DEFAULT_MIN_CONFIRMATIONS,
|
||
DEFAULT_MIN_FLOOR_PAIRS,
|
||
DEFAULT_REVISIT_FLOOR_QUANTILE,
|
||
deactivate_stale_listings,
|
||
)
|
||
|
||
listing_source: str = params.get("listing_source", "avito")
|
||
ttl_days: int = params.get("ttl_days", _settings.avito_stale_ttl_days)
|
||
segments: list[str] | None = params.get("segments")
|
||
staleness_column: str = params.get("staleness_column", "last_seen_at")
|
||
# Гейт по здоровью сбора (#2659) включён по умолчанию: незасеянное расписание
|
||
# получает страховочный порог, а не «деактивируй вслепую». Посчитанные по
|
||
# источнику пороги приходят из default_params (миграция 219).
|
||
min_confirmations: int = params.get("min_confirmations", DEFAULT_MIN_CONFIRMATIONS)
|
||
# Пол TTL по измеренному циклу переобхода (#2659) — тоже включён по умолчанию:
|
||
# незасеянное расписание не должно снимать объявления по порогу ниже собственного
|
||
# хвоста обхода. Снять ручку вручную: revisit_floor_quantile = 0.
|
||
revisit_floor_quantile: float = params.get(
|
||
"revisit_floor_quantile", DEFAULT_REVISIT_FLOOR_QUANTILE
|
||
)
|
||
# Пустой (NULL) listing_segment -- легаси-строки до миграции 011 + жертвы
|
||
# отсутствующего COALESCE в ON CONFLICT (base.py upsert никогда не перезаписывает
|
||
# listing_segment на повторном скрейпе). Отдельный явный предикат IS NULL, а не
|
||
# элемент :segments (ANY(...) никогда не матчит NULL) -- см. deactivate_stale_avito.py.
|
||
null_segment_only: bool = params.get("null_segment_only", False)
|
||
# Потолок эффективного TTL (см. CAP_MULT в deactivate_stale_avito.py) — множитель,
|
||
# а не голая константа: источник с непропорционально длинным хвостом переобхода
|
||
# относительно своего ttl_days переопределяет его через default_params (ключ
|
||
# "cap_mult"), не трогая дефолт для остальных источников.
|
||
cap_mult: float = params.get("cap_mult", CAP_MULT)
|
||
# Гейт деградации пола (PR-B, #2659 продолжение) -- включён по умолчанию, тот
|
||
# же принцип, что у min_confirmations/revisit_floor_quantile выше: незасеянное
|
||
# расписание получает страховку, а не «деактивируй вслепую».
|
||
min_floor_pairs: int = params.get("min_floor_pairs", DEFAULT_MIN_FLOOR_PAIRS)
|
||
floor_drop_ratio: float = params.get("floor_drop_ratio", DEFAULT_FLOOR_DROP_RATIO)
|
||
# Аварийный (не рабочий) потолок объёма снятия за один прогон -- см.
|
||
# DEFAULT_MAX_DEACTIVATED в deactivate_stale_avito.py.
|
||
max_deactivated: int = params.get("max_deactivated", DEFAULT_MAX_DEACTIVATED)
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(
|
||
None,
|
||
lambda: deactivate_stale_listings(
|
||
db,
|
||
run_id,
|
||
listing_source=listing_source,
|
||
ttl_days=ttl_days,
|
||
segments=segments,
|
||
staleness_column=staleness_column,
|
||
min_confirmations=min_confirmations,
|
||
revisit_floor_quantile=revisit_floor_quantile,
|
||
null_segment_only=null_segment_only,
|
||
cap_mult=cap_mult,
|
||
min_floor_pairs=min_floor_pairs,
|
||
floor_drop_ratio=floor_drop_ratio,
|
||
max_deactivated=max_deactivated,
|
||
),
|
||
)
|
||
|
||
|
||
# ── sber_index_pull — async, owns lifecycle ──────────────────────────────────
|
||
async def _job_sber_index_pull(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.sber_index_pull import run_sber_index_pull
|
||
|
||
await run_sber_index_pull(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── rosreestr_quarter_poll — async, owns lifecycle ───────────────────────────
|
||
async def _job_rosreestr_quarter_poll(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.rosreestr_quarter_poll import run_rosreestr_quarter_poll
|
||
|
||
await run_rosreestr_quarter_poll(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── deals_freshness_monitor — sync DB-only freshness check в executor ─────────
|
||
async def _job_deals_freshness_monitor(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.deals_freshness_monitor import check_deals_freshness
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(None, check_deals_freshness, db, run_id, params)
|
||
|
||
|
||
# ── landing_stats_refresh — sync DB-only пересчёт витрины в executor ─────────
|
||
async def _job_landing_stats(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.landing_stats import refresh_landing_stats
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(None, refresh_landing_stats, db, run_id, params)
|
||
|
||
|
||
# ── sber_freshness_monitor — sync DB-only freshness check в executor ──────────
|
||
async def _job_sber_freshness_monitor(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.sber_freshness_monitor import check_sber_freshness
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(None, check_sber_freshness, db, run_id, params)
|
||
|
||
|
||
# ── newbuilding_enrich — async, owns lifecycle ───────────────────────────────
|
||
async def _job_newbuilding_enrich(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.newbuilding_enrich_backfill import run_newbuilding_enrich
|
||
|
||
await run_newbuilding_enrich(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── yandex_newbuilding_sweep — async, lifecycle в job ────────────────────────
|
||
async def _job_yandex_newbuilding_sweep(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.yandex_newbuilding_sweep import enrich_yandex_newbuilding_sweep
|
||
|
||
try:
|
||
result = await enrich_yandex_newbuilding_sweep(
|
||
db,
|
||
limit=int(params.get("limit", 5)),
|
||
request_delay_sec=float(params.get("request_delay_sec", 8.0)),
|
||
city=str(params.get("city", "ekaterinburg")),
|
||
)
|
||
counters = result.to_dict()
|
||
# #2860: прогон, который обработал дома и не разрешил НИ ОДНОГО, успешным
|
||
# называть нельзя. Четырнадцать таких прогонов подряд (16.07-17.08.2026)
|
||
# стояли `done` с пустым error_text — механизм исполнялся, счётчик был
|
||
# честный, вывод из него не делал никто, и очередь тихо росла 351 → 397.
|
||
#
|
||
# НЕ подгоняем счётчик: если дома действительно не разрешаются, честный
|
||
# исход — назвать прогон неуспешным, а не дотянуть succeeded до ненуля.
|
||
processed = int(counters.get("processed") or 0)
|
||
succeeded = int(counters.get("succeeded") or 0)
|
||
if processed > 0 and succeeded == 0:
|
||
logger.warning(
|
||
"yandex_newbuilding_sweep run_id=%d: обработано %d, разрешено 0 — "
|
||
"помечаю прогон неуспешным",
|
||
run_id,
|
||
processed,
|
||
)
|
||
ctx.runs.mark_failed(
|
||
db,
|
||
run_id,
|
||
f"обработано {processed} домов, разрешено 0 — полный отказ разрешения slug",
|
||
counters,
|
||
)
|
||
else:
|
||
ctx.runs.mark_done(db, run_id, counters)
|
||
except Exception:
|
||
logger.exception("scheduler: enrich_yandex_newbuilding_sweep crashed run_id=%d", run_id)
|
||
try:
|
||
ctx.runs.mark_failed(db, run_id, "crashed", {})
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
# ── geocode_missing_listings — async, owns lifecycle ─────────────────────────
|
||
async def _job_geocode_missing_listings(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.geocode_missing import run_geocode_missing_listings
|
||
|
||
await run_geocode_missing_listings(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── avito_detail_backfill — async, owns lifecycle ────────────────────────────
|
||
async def _job_avito_detail_backfill(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.avito_detail_backfill import run_avito_detail_backfill
|
||
|
||
await run_avito_detail_backfill(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── yandex_detail_backfill — async, owns lifecycle ───────────────────────────
|
||
async def _job_yandex_detail_backfill(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.yandex_detail_backfill import run_yandex_detail_backfill
|
||
|
||
await run_yandex_detail_backfill(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── domclick_detail_backfill — async, owns lifecycle (issue #2000 Layer B) ───
|
||
async def _job_domclick_detail_backfill(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.domclick_detail_backfill import run_domclick_detail_backfill
|
||
|
||
await run_domclick_detail_backfill(db, run_id=run_id, params=params)
|
||
|
||
|
||
# ── domclick_city_sweep — override kit-native с инъекцией куки сессии (#3264) ─
|
||
# Kit-native `_job_domclick_city_sweep` (scraper_kit.orchestration.scheduler) ходит на
|
||
# bff-search-web.domclick.ru БЕЗ кук: kit не имеет права импортировать app.* / читать БД
|
||
# (strangler-инвариант #2133). Каждый запрос свипа поэтому был обязан решать QRATOR PoW
|
||
# с нуля — прод 30.08 (прогон 5330): 9/9 запросов зависли на challenge, 0 листингов.
|
||
# `domclick_detail_backfill` (добор карточек) уже решает эту задачу инъекцией того же
|
||
# снимка (domclick_session.load_session) — здесь тот же приём для свипа.
|
||
#
|
||
# Регистрируем ЗДЕСЬ (app-side) под тем же ключом "domclick_city_sweep", а не правим
|
||
# kit: build_registry явно допускает переопределение kit-native продуктовым Handler'ом
|
||
# (см. докстринг build_registry — "последнее слово за продуктом"), а строгий запрет на
|
||
# app.* внутри kit при этом не нарушается — БД читает только этот модуль.
|
||
async def _job_domclick_city_sweep(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from scraper_kit.orchestration.pipeline import run_domclick_city_sweep
|
||
|
||
from app.services import domclick_session as domclick_session_svc
|
||
|
||
# None — сессии в БД нет либо протухла (load_session сам фильтрует
|
||
# expires_at_estimate > NOW()). Это НЕ авария: свип продолжает работать без
|
||
# инъекции, ровно как до #3264, только логируем факт один раз для видимости.
|
||
cookies = domclick_session_svc.load_session(db)
|
||
if cookies is None:
|
||
logger.info(
|
||
"domclick_city_sweep run_id=%d: куки Sber ID сессии отсутствуют/протухли — "
|
||
"свип идёт без инъекции (см. #3264)",
|
||
run_id,
|
||
)
|
||
|
||
await run_domclick_city_sweep(
|
||
db,
|
||
run_id=run_id,
|
||
config=ctx.config,
|
||
matcher=ctx.matcher,
|
||
shutdown_requested=ctx.shutdown_requested,
|
||
proxy_provider=ctx.proxy_provider,
|
||
city_id=int(params.get("city_id", 4)),
|
||
rooms=params.get("rooms"),
|
||
pages=int(params.get("pages_per_anchor", 5)),
|
||
request_delay_sec=float(params.get("request_delay_sec", 6.0)),
|
||
resume_run_id=kit_pick_resume(db, run_id),
|
||
cookies=cookies,
|
||
)
|
||
|
||
|
||
# ── house_coords_from_listings — sync set-based UPDATE в executor (#2771) ─────
|
||
async def _job_house_coords_from_listings(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.backfill_house_coords_from_listings import run_house_coords_from_listings
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(
|
||
None,
|
||
lambda: run_house_coords_from_listings(db, run_id=run_id, params=params),
|
||
)
|
||
|
||
|
||
# ── geoportal_coords_backfill — sync local exact match в executor (#1967) ─────
|
||
async def _job_geoportal_coords_backfill(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.backfill_listings_coords_geoportal import run_geoportal_coords_backfill
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(
|
||
None,
|
||
lambda: run_geoportal_coords_backfill(db, run_id=run_id, params=params),
|
||
)
|
||
|
||
|
||
# ── cadastral_geo_match — sync FDW-refresh+match в executor ───────────────────
|
||
async def _job_cadastral_geo_match(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.cadastral_geo_match import run_cadastral_geo_match
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(
|
||
None,
|
||
lambda: run_cadastral_geo_match(db, run_id=run_id, params=params),
|
||
)
|
||
|
||
|
||
# ── osm_poi_ekb_refresh — sync FDW-refresh в executor ─────────────────────────
|
||
async def _job_osm_poi_ekb_refresh(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.osm_poi_ekb_refresh import run_osm_poi_ekb_refresh
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(
|
||
None,
|
||
lambda: run_osm_poi_ekb_refresh(db, run_id=run_id, params=params),
|
||
)
|
||
|
||
|
||
# ── house_imv_backfill — async Avito-IMV с heartbeat, lifecycle в job ─────────
|
||
async def _job_house_imv_backfill(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
batch_size = int(params.get("batch_size", 50))
|
||
request_delay_sec = float(params.get("request_delay_sec", 5.0))
|
||
only_status = str(params.get("only_status", "pending"))
|
||
|
||
def _heartbeat() -> None:
|
||
# heartbeat каждые N домов — иначе reap_zombies снимет живой долгий IMV-run (#1363).
|
||
ctx.runs.update_heartbeat(db, run_id, {})
|
||
|
||
try:
|
||
result = await ctx.enrichment.house_imv_backfill(
|
||
db,
|
||
batch_size=batch_size,
|
||
request_delay_sec=request_delay_sec,
|
||
only_status=only_status,
|
||
heartbeat=_heartbeat,
|
||
)
|
||
counters = {
|
||
"checked": result.checked,
|
||
"saved": result.saved,
|
||
"skipped": result.skipped,
|
||
"errors": result.errors,
|
||
"duration_sec": int(result.duration_sec),
|
||
# #2674: _column_counts (scrape_runs.py) берёт выделенные колонки из
|
||
# ключей total_seen|lots_fetched и new_count|lots_inserted — ни одного
|
||
# из них тут не было, поэтому все 39 прогонов этого source лежат в БД
|
||
# с total_seen=0. А mark_done по этой же колонке шлёт алерт «3 подряд
|
||
# done с нулевым результатом» (#2625) — то есть даже идеальный прогон
|
||
# с 50 сохранёнными считался бы нулевым и через три дня выстрелил бы
|
||
# ложной тревогой про капчу.
|
||
# Трейд-офф: на исчерпанной очереди checked=0 три дня подряд тоже даст
|
||
# алерт — но пустая очередь при ежедневном расписании это и правда сигнал.
|
||
"total_seen": result.checked,
|
||
"new_count": result.saved,
|
||
# #2674: из скольких слотов пакета взяты дома на ПОВТОР (transient_error)
|
||
# и сколько домов ушло в no_params до пакета одним запросом. Без этих
|
||
# двух счётчиков в scrape_runs.counters проверить, что застрявшие
|
||
# действительно возвращаются в очередь, можно только по houses.
|
||
"retried": result.retried,
|
||
"premarked": result.premarked,
|
||
}
|
||
# Честный статус (#2674, тот же класс, что #2670/#2657): успех — это
|
||
# «сделали то, что собирались», а не «не поймали известное исключение».
|
||
# На проде так ушли в done 31 прогон подряд: saved=0 при errors≈35 из 50.
|
||
# Ноль сохранённых БЕЗ ошибок (всё отфильтровано в skipped) — честная
|
||
# пустота, она по-прежнему done.
|
||
if result.saved == 0 and result.errors > 0:
|
||
ctx.runs.mark_failed(
|
||
db,
|
||
run_id,
|
||
f"saved=0 при errors={result.errors} (checked={result.checked})",
|
||
counters,
|
||
)
|
||
else:
|
||
ctx.runs.mark_done(db, run_id, counters)
|
||
except Exception as exc:
|
||
logger.exception("scheduler: house_imv_backfill crashed run_id=%d", run_id)
|
||
try:
|
||
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
|
||
except Exception:
|
||
logger.exception("scheduler: mark_failed crashed run_id=%d", run_id)
|
||
|
||
|
||
# ── domrf_kapremont_load — sync загрузка open data ДОМ.РФ в executor ─────────
|
||
# #2674: loader (services/domrf_kapremont_loader.py) и CLI (tasks/domrf_kapremont_load.py)
|
||
# написаны и покрыты тестами с #2013, но Handler'а и строки расписания не было — источник
|
||
# запускали руками ровно один раз, 12.07.2026 (29 978 строк, один и тот же loaded_at у всех).
|
||
# Это не мёртвый код, а оборванная проводка: нечему было его вызвать.
|
||
async def _job_domrf_kapremont_load(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
"""Скачать КР1.1+КР1.2 ДОМ.РФ → staging → backfill houses → propagate listings.
|
||
|
||
Тело переиспользует те же три функции, что и CLI (дизайн-инвариант модуля: не
|
||
дублируем логику). Lifecycle не свой — mark_done/mark_failed здесь, как у
|
||
_job_yandex_newbuilding_sweep.
|
||
|
||
Счётчики кладём в total_seen/new_count: `scrape_runs._column_counts` берёт выделенные
|
||
колонки именно из этих ключей, и по ним же mark_done ловит «три подряд нулевых
|
||
прогона» (#2625) — без них идеальный прогон лежал бы в БД как нулевой (тот же
|
||
промах, что чинили у house_imv_backfill).
|
||
"""
|
||
from app.services.domrf_kapremont_loader import (
|
||
backfill_houses_from_domrf,
|
||
load_domrf_kapremont,
|
||
propagate_listings_year_from_houses,
|
||
)
|
||
|
||
def _run() -> dict[str, int]:
|
||
load_counts = load_domrf_kapremont(db)
|
||
db.commit()
|
||
houses_counts = backfill_houses_from_domrf(db)
|
||
listings_counts = propagate_listings_year_from_houses(db)
|
||
db.commit()
|
||
return {
|
||
"kr11_rows": load_counts["kr11_rows"],
|
||
"upserted": load_counts["upserted"],
|
||
"houses_updated": houses_counts["houses_updated"],
|
||
"listings_updated": listings_counts["listings_updated"],
|
||
# см. докстринг: выделенные колонки прогона + гейт «нулевой прогон».
|
||
"total_seen": load_counts["kr11_rows"],
|
||
"new_count": houses_counts["houses_updated"] + listings_counts["listings_updated"],
|
||
}
|
||
|
||
loop = asyncio.get_event_loop()
|
||
try:
|
||
counters = await loop.run_in_executor(None, _run)
|
||
ctx.runs.mark_done(db, run_id, counters)
|
||
except Exception as exc:
|
||
logger.exception("scheduler: domrf_kapremont_load crashed run_id=%d", run_id)
|
||
db.rollback()
|
||
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
|
||
|
||
|
||
# ── purge_expired_trade_in_data — ЭТАП 4 B2C retention (152-ФЗ) ───────────────
|
||
async def _job_purge_expired_trade_in_data(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.tasks.purge_expired_trade_in_data import purge_expired_trade_in_data
|
||
|
||
batch_size = params.get("batch_size")
|
||
max_batches = params.get("max_batches")
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(
|
||
None,
|
||
lambda: purge_expired_trade_in_data(
|
||
db, run_id, batch_size=batch_size, max_batches=max_batches
|
||
),
|
||
)
|
||
|
||
|
||
# ── house_dedup_merge — sync destructive merge в executor, owns lifecycle ─────
|
||
async def _job_house_dedup_merge(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.services.house_dedup_merge import run_house_dedup_merge
|
||
|
||
loop = asyncio.get_event_loop()
|
||
await loop.run_in_executor(
|
||
None,
|
||
lambda: run_house_dedup_merge(db, run_id=run_id, params=params),
|
||
)
|
||
|
||
|
||
# ── proxy_healthcheck — async, sub-hourly reschedule + lifecycle в job ────────
|
||
async def _job_proxy_healthcheck(
|
||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||
) -> None:
|
||
from app.services.proxy_pool import run_proxy_healthcheck
|
||
|
||
try:
|
||
counters = await run_proxy_healthcheck(db)
|
||
ctx.runs.mark_done(db, run_id, counters)
|
||
except Exception as exc:
|
||
logger.exception("scheduler: run_proxy_healthcheck crashed run_id=%d", run_id)
|
||
try:
|
||
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
|
||
except Exception:
|
||
logger.exception("scheduler: mark_failed crashed run_id=%d", run_id)
|
||
|
||
|
||
def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]:
|
||
"""Реестр НЕ-sweep продуктовых source→Handler для kit build_registry.
|
||
|
||
Kit-native sweeps (avito/yandex/cian/domclick city/full-load/newbuilding) НЕ здесь —
|
||
их даёт build_registry(_default_kit_handlers). Здесь — именованные + 1 wildcard
|
||
(deactivate_stale_*), покрывающие каждый НЕ-sweep source боевого scheduler-dispatch.
|
||
(Число намеренно не названо: прежнее «19» разошлось с реальностью на пять записей.)
|
||
|
||
`ctx` — принят для симметрии контракта; сами Handler-job'ы получают ctx во время
|
||
dispatch (см. kit `_dispatch`), поэтому здесь он не замыкается.
|
||
"""
|
||
return {
|
||
"cian_history_backfill": Handler(
|
||
_job_cian_history_backfill,
|
||
"cian_history_backfill",
|
||
pre_claim=_cian_pre_claim,
|
||
),
|
||
# #3284: то же тело, другая очередь. Разные source нужны именно как РАЗНЫЕ
|
||
# строки расписания — у них свои окна, свой next_run_at и свой счётчик
|
||
# прогонов; гонять оба режима под одним source нельзя, планировщик держит
|
||
# на source ровно один активный прогон.
|
||
# Каденс — в минутах (замер прода 31.08.2026): та же суточная дыра
|
||
# compute_next_run_at, что у avito/yandex_detail_backfill — Циан держит
|
||
# темп ~23 с/карточку при очереди 20 501 и НОЛЕ блоков за неделю.
|
||
#
|
||
# 360 минут (4 прогона/сутки), НЕ 180 как у Яндекса — АСИММЕТРИЯ ниже
|
||
# ОБЯЗАТЕЛЬНА к прочтению перед тем, как выравнивать эти два интервала:
|
||
# Циан ходит через BrowserFetcher и ДЕРЖИТ lease узла на ВЕСЬ прогон, а
|
||
# не только берёт URL прокси (как Яндекс через resolve_proxy_url/
|
||
# proxy_egress). Пул сейчас 3 живых узла — дефицитный ресурс, который
|
||
# Циан на время прогона отбирает у avito_detail_backfill и
|
||
# domclick_detail_backfill. Учащение Циана до 180 мин без замера того,
|
||
# как это скажется на соседях по пулу, — прямой риск сжечь мощность,
|
||
# которой и так не хватает на троих.
|
||
"cian_detail_backfill": Handler(
|
||
_job_cian_history_backfill,
|
||
"cian_detail_backfill",
|
||
pre_claim=_cian_pre_claim,
|
||
post_claim=reschedule_after_minutes(param="interval_minutes", default=360),
|
||
),
|
||
"rosreestr_dkp_import": Handler(_job_rosreestr_dkp, "rosreestr_dkp_import"),
|
||
"listing_source_snapshot": Handler(_job_listing_source_snapshot, "listing_source_snapshot"),
|
||
"asking_to_sold_ratio_refresh": Handler(
|
||
_job_asking_to_sold_ratio, "asking_to_sold_ratio_refresh"
|
||
),
|
||
"deal_city_price_bands_refresh": Handler(
|
||
_job_deal_city_price_bands_refresh, "deal_city_price_bands_refresh"
|
||
),
|
||
"refresh_search_matview": Handler(_job_refresh_search_matview, "refresh_search_matview"),
|
||
"yandex_address_backfill": Handler(_job_yandex_address_backfill, "yandex_address_backfill"),
|
||
"sber_index_pull": Handler(_job_sber_index_pull, "sber_index_pull"),
|
||
"rosreestr_quarter_poll": Handler(_job_rosreestr_quarter_poll, "rosreestr_quarter_poll"),
|
||
"deals_freshness_monitor": Handler(_job_deals_freshness_monitor, "deals_freshness_monitor"),
|
||
"sber_freshness_monitor": Handler(_job_sber_freshness_monitor, "sber_freshness_monitor"),
|
||
"landing_stats_refresh": Handler(_job_landing_stats, "landing_stats_refresh"),
|
||
"newbuilding_enrich": Handler(_job_newbuilding_enrich, "newbuilding_enrich"),
|
||
"yandex_newbuilding_sweep": Handler(
|
||
_job_yandex_newbuilding_sweep, "yandex_newbuilding_sweep"
|
||
),
|
||
"geoportal_coords_backfill": Handler(
|
||
_job_geoportal_coords_backfill, "geoportal_coords_backfill"
|
||
),
|
||
"house_coords_from_listings": Handler(
|
||
_job_house_coords_from_listings, "house_coords_from_listings"
|
||
),
|
||
"geocode_missing_listings": Handler(
|
||
_job_geocode_missing_listings, "geocode_missing_listings"
|
||
),
|
||
# Каденс — в минутах, а не в сутках (замер 2026-08-22). Дефолтная
|
||
# гранулярность compute_next_run_at — сутки, и бэкфилл получал ровно один
|
||
# прогон в день. При этом прогон умирает по бану через 17-83 минуты, то
|
||
# есть 23 часа из 24 задание простаивало: 238 обогащённых карточек за
|
||
# сутки при 9 951 активном объявлении — очередь разбиралась бы месяцами.
|
||
#
|
||
# Хук сам себя throttl'ит: пока прогон идёт, has_running_run в _claim_run
|
||
# возвращает None и next_run_at не сбрасывается.
|
||
#
|
||
# 180 минут — ОСОЗНАННО консервативная отправная точка, а не найденный
|
||
# оптимум. Данных для подбора нет: два прогона дали противоречивую
|
||
# картину (4562 — 83 мин и 175 карточек; 4586 через 4.3 часа, когда пул
|
||
# прокси был давно чист, — 17 мин и 42 карточки). Значит память Авито
|
||
# длиннее часов, и учащение может ухудшить выход, а не улучшить.
|
||
#
|
||
# Риск, который надо держать в голове при подборе: те же 4 прокси
|
||
# обслуживают SERP-свипы — первичный сбор. Сжечь их на обогащении хуже,
|
||
# чем медленно обогащать. Двигать интервал вниз только по замеру
|
||
# нескольких суток, глядя и на свипы тоже.
|
||
"avito_detail_backfill": Handler(
|
||
_job_avito_detail_backfill,
|
||
"avito_detail_backfill",
|
||
post_claim=reschedule_after_minutes(param="interval_minutes", default=180),
|
||
),
|
||
# Каденс — в минутах (замер прода 31.08.2026), тот же дефект, что у
|
||
# avito_detail_backfill выше: без post_claim compute_next_run_at держит
|
||
# суточную гранулярность, а прогон Яндекса разбирает 375-450 карточек в
|
||
# час (budget_sec=3600) при очереди 11 110 и НУЛЕ блоков за неделю — узкое
|
||
# место чисто в частоте запуска, не в самом фетчере.
|
||
#
|
||
# 180 минут (8 прогонов/сутки) безопасны именно из-за асимметрии с Циан:
|
||
# Яндекс ходит через resolve_proxy_url (proxy_egress) — берёт URL прокси,
|
||
# но НЕ держит lease узла на весь прогон, поэтому учащение не отнимает
|
||
# узел у других source.
|
||
"yandex_detail_backfill": Handler(
|
||
_job_yandex_detail_backfill,
|
||
"yandex_detail_backfill",
|
||
post_claim=reschedule_after_minutes(param="interval_minutes", default=180),
|
||
),
|
||
"domclick_detail_backfill": Handler(
|
||
_job_domclick_detail_backfill, "domclick_detail_backfill"
|
||
),
|
||
# Override kit-native (#3264) — инъекция куки сессии, см. докстринг job'а выше.
|
||
"domclick_city_sweep": Handler(_job_domclick_city_sweep, "domclick_city_sweep"),
|
||
"cadastral_geo_match": Handler(_job_cadastral_geo_match, "cadastral_geo_match"),
|
||
"osm_poi_ekb_refresh": Handler(_job_osm_poi_ekb_refresh, "osm_poi_ekb_refresh"),
|
||
"house_imv_backfill": Handler(_job_house_imv_backfill, "house_imv_backfill"),
|
||
"house_dedup_merge": Handler(_job_house_dedup_merge, "house_dedup_merge"),
|
||
"domrf_kapremont_load": Handler(_job_domrf_kapremont_load, "domrf_kapremont_load"),
|
||
"purge_expired_trade_in_data": Handler(
|
||
_job_purge_expired_trade_in_data, "purge_expired_trade_in_data"
|
||
),
|
||
"proxy_healthcheck": Handler(
|
||
_job_proxy_healthcheck,
|
||
"proxy_healthcheck",
|
||
post_claim=reschedule_after_minutes(param="interval_minutes", default=30),
|
||
),
|
||
# Wildcard: deactivate_stale_avito / _yandex / _cian — один Handler на семейство
|
||
# (resolve_handler матчит по префиксу key[:-1]="deactivate_stale_").
|
||
"deactivate_stale_*": Handler(_job_deactivate_stale, "deactivate_stale"),
|
||
}
|
||
|
||
|
||
__all__ = ["build_product_handlers"]
|