All checks were successful
CI / changes (push) Successful in 8s
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI / backend-tests (push) Has been skipped
CI / frontend-tests (push) Has been skipped
CI / openapi-codegen-check (push) Has been skipped
CI / changes (pull_request) Successful in 6s
CI / backend-tests (pull_request) Has been skipped
Generalize the avito-only stale-listing deactivation into a per-source,
segment-aware task and enable it for yandex/cian vtorichka.
Why segment-aware (not blanket TTL like avito): yandex/cian city sweeps do
not maintain full inventory coverage (cian_city_sweep disabled, yandex narrow
5-anchor nightly sweep refreshes ~48/day of 3969 active). A blanket TTL would
deactivate live inventory — including 9659 active first-party novostroyki —
because 'not seen in N days' means 'not re-covered', not 'sold'. So yandex/cian
deactivate ONLY listing_segment='vtorichka' at TTL=30; novostroyki and
NULL-segment are protected. avito unchanged (all segments, TTL=10).
- deactivate_stale_listings(db, run_id, *, listing_source, ttl_days, segments)
generic fn; deactivate_stale_avito_listings kept as back-compat wrapper.
- scheduler dispatch: source.startswith('deactivate_stale_') -> generic trigger
reading listing_source/ttl_days/segments from default_params jsonb.
- 114 seed: deactivate_stale_yandex + deactivate_stale_cian (vtorichka, TTL=30).
- segments is None => all segments; [] => none (ANY(ARRAY[]) matches nothing).
Refs #759
137 lines
6.2 KiB
Python
137 lines
6.2 KiB
Python
"""Deactivate stale listings that have not been seen in >TTL days.
|
||
|
||
Первоначально создано для avito (#759). Теперь поддерживает любой источник
|
||
(listing_source) с опциональной фильтрацией по listing_segment.
|
||
|
||
Ключевые решения:
|
||
- Cian/Yandex не поддерживают full-coverage sweep -> паушальный TTL сломает живой
|
||
инвентарь. DECISION: для yandex/cian деактивировать ТОЛЬКО listing_segment='vtorichka',
|
||
TTL=30. novostroyki (9659 активных первичных строк) и NULL-сегмент не трогаем.
|
||
- avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений.
|
||
- Строки НЕ удаляются -- история нужна для бэктеста (#667).
|
||
|
||
Задача синхронная (DB-only, никаких внешних HTTP-вызовов) -- запускается in-app
|
||
scheduler'ом через trigger_deactivate_stale_run() (scheduler.py),
|
||
по образцу snapshot_listing_sources / recompute_asking_to_sold_ratios
|
||
(sync task в run_in_executor).
|
||
|
||
TTL для avito берётся из settings.avito_stale_ttl_days (env AVITO_STALE_TTL_DAYS, default 10).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from typing import Any
|
||
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.config import settings
|
||
from app.services import scrape_runs as runs_mod
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Базовый UPDATE без фильтра по сегменту: все сегменты для данного source.
|
||
# CAST(:ttl_days || ' days' AS interval) -- psycopg v3 safe (никаких :param::type).
|
||
_DEACTIVATE_SQL_ALL_SEGMENTS = text(
|
||
"""
|
||
UPDATE listings
|
||
SET is_active = false
|
||
WHERE source = :listing_source
|
||
AND is_active = true
|
||
AND last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||
"""
|
||
)
|
||
|
||
# UPDATE с фильтром по сегменту (segments задан): затрагивает только указанные сегменты.
|
||
# = ANY(CAST(:segments AS text[])) -- psycopg v3 адаптирует Python list -> text[].
|
||
_DEACTIVATE_SQL_SEGMENTS = text(
|
||
"""
|
||
UPDATE listings
|
||
SET is_active = false
|
||
WHERE source = :listing_source
|
||
AND is_active = true
|
||
AND last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||
AND listing_segment = ANY(CAST(:segments AS text[]))
|
||
"""
|
||
)
|
||
|
||
# Алиас для обратной совместимости -- использовался в тестах через task_mod._DEACTIVATE_SQL
|
||
_DEACTIVATE_SQL = _DEACTIVATE_SQL_ALL_SEGMENTS
|
||
|
||
|
||
def deactivate_stale_listings(
|
||
db: Session,
|
||
run_id: int,
|
||
*,
|
||
listing_source: str,
|
||
ttl_days: int,
|
||
segments: list[str] | None = None,
|
||
) -> dict[str, int]:
|
||
"""Пометить is_active=false объявления с last_seen_at > ttl_days дней назад.
|
||
|
||
Параметры:
|
||
listing_source: значение поля source в таблице listings ('avito', 'yandex', 'cian').
|
||
ttl_days: количество дней TTL; объявления старше этого порога деактивируются.
|
||
segments: если задан -- деактивировать только объявления с указанными
|
||
listing_segment значениями. None -> все сегменты (поведение avito по умолчанию).
|
||
|
||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||
Один UPDATE в транзакции. Финализирует scrape_runs (mark_done / mark_failed).
|
||
|
||
Returns {"deactivated": N} -- количество обновлённых строк.
|
||
"""
|
||
counters: dict[str, int] = {"deactivated": 0}
|
||
try:
|
||
# segments is None -> все сегменты (поведение avito). segments=[...] -> только
|
||
# перечисленные сегменты. Используем `is not None` (НЕ truthy): пустой список []
|
||
# означает "ни один сегмент" (= ANY(ARRAY[]) ничего не матчит, деактивирует 0),
|
||
# а НЕ "все сегменты" — иначе случайный [] стёр бы весь источник.
|
||
if segments is not None:
|
||
params: dict[str, Any] = {
|
||
"listing_source": listing_source,
|
||
"ttl_days": ttl_days,
|
||
"segments": segments,
|
||
}
|
||
result = db.execute(_DEACTIVATE_SQL_SEGMENTS, params)
|
||
else:
|
||
params = {"listing_source": listing_source, "ttl_days": ttl_days}
|
||
result = db.execute(_DEACTIVATE_SQL_ALL_SEGMENTS, params)
|
||
|
||
counters["deactivated"] = result.rowcount or 0
|
||
|
||
db.commit()
|
||
runs_mod.mark_done(db, run_id, counters)
|
||
logger.info(
|
||
"deactivate_stale source=%s run_id=%d done: deactivated=%d "
|
||
"(ttl_days=%d, segments=%r)",
|
||
listing_source,
|
||
run_id,
|
||
counters["deactivated"],
|
||
ttl_days,
|
||
segments,
|
||
)
|
||
return counters
|
||
except Exception as exc:
|
||
logger.exception("deactivate_stale source=%s run_id=%d failed", listing_source, run_id)
|
||
db.rollback()
|
||
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters)
|
||
raise
|
||
|
||
|
||
def deactivate_stale_avito_listings(db: Session, run_id: int) -> dict[str, int]:
|
||
"""Обёртка для обратной совместимости -- deactivate_stale_avito (все сегменты, TTL из settings).
|
||
|
||
Пометить is_active=false все avito-объявления с last_seen_at > TTL дней назад.
|
||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||
Один UPDATE в транзакции. Финализирует scrape_runs (mark_done / mark_failed).
|
||
|
||
Returns {"deactivated": N} -- количество обновлённых строк.
|
||
"""
|
||
return deactivate_stale_listings(
|
||
db,
|
||
run_id,
|
||
listing_source="avito",
|
||
ttl_days=settings.avito_stale_ttl_days,
|
||
segments=None,
|
||
)
|