fix(tradein/scraper): deactivate stale yandex/cian listings with NULL segment
544 yandex + 224 cian active rows carry listing_segment=NULL (legacy rows predating migration 011, plus a small trickle that can never self-heal since the upsert ON CONFLICT never rewrites listing_segment on re-scrape). 94-97% of them are frozen at ~86 days old, yet the estimator's Tier A (same-building) and Tier C (micro-radius) anchor queries filter only is_active=true -- no freshness column -- so these stale asking prices anchor live valuations. deactivate_stale_yandex/_cian (migration 115) already run at TTL=30 but scope segments=['vtorichka'] only: `= ANY(CAST(:segments AS text[]))` never matches NULL, so the NULL bucket was invisible to both existing jobs and to avito/domklik/n1 (which are source-blanket or vtorichka-only respectively). Adds null_segment_only kwarg to deactivate_stale_listings() building an explicit `listing_segment IS NULL` predicate (confirmations/revisit-floor builders extended in parallel for correctness, though both gates are kept off for this slice -- population too small for thresholds calibrated on a full vtorichka sweep, would permanently skip as unhealthy). Two new scrape_schedules rows (migration 264) run it per source, untouched novostroyki/vtorichka jobs unaffected. TTL=60d (vs 30d for vtorichka): revisit-floor is not computable here (rows out of dedicated sweep scope have no revisit-gap history), so the margin is folded into the TTL directly -- 2.26x/1.4x over the measured p99 revisit gaps already on file for cian/yandex vtorichka (26.6d/43.0d). Bimodal age distribution means this costs almost nothing in coverage (30d vs 60d: 744 vs 734 deactivated). First run: 734 of 768 NULL rows deactivated (211 cian, 523 yandex), 34 remain (fresher than 60d, still incidentally re-touched). Comparable pool (is_active AND segment IN (NULL, vtorichka)) after: cian -2.7% (7739->7528), yandex -10.2% (5133->4610) -- yandex crosses the 10% flag threshold. All 734 removed rows already had scraped_at frozen >60d, i.e. already excluded from Tier S/H (which do filter freshness, 14-60d window) -- the drop is real for is_active headcount but zero-impact there; it only prunes Tier A/C, where it removes stale prices rather than live comps. Separate finding (not fixed here): base.py's ON CONFLICT DO UPDATE omits listing_segment from SET entirely, so a legacy NULL row can never heal even though cian/yandex SERP always compute segment deterministically on re-scrape.
This commit is contained in:
parent
7def4973bd
commit
cfb4c159ab
5 changed files with 410 additions and 15 deletions
|
|
@ -236,6 +236,11 @@ async def _job_deactivate_stale(
|
|||
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)
|
||||
|
||||
loop = asyncio.get_event_loop()
|
||||
await loop.run_in_executor(
|
||||
|
|
@ -249,6 +254,7 @@ async def _job_deactivate_stale(
|
|||
staleness_column=staleness_column,
|
||||
min_confirmations=min_confirmations,
|
||||
revisit_floor_quantile=revisit_floor_quantile,
|
||||
null_segment_only=null_segment_only,
|
||||
),
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -6,8 +6,17 @@
|
|||
Ключевые решения:
|
||||
- Cian/Yandex не поддерживают full-coverage sweep -> паушальный TTL сломает живой
|
||||
инвентарь. DECISION: для yandex/cian деактивировать ТОЛЬКО listing_segment='vtorichka',
|
||||
TTL=30. novostroyki (9659 активных первичных строк) и NULL-сегмент не трогаем.
|
||||
TTL=30. novostroyki (9659 активных первичных строк) не трогаем.
|
||||
- avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений.
|
||||
- NULL-сегмент (легаси-строки до миграции 011 + жертвы бага в ON CONFLICT -- upsert
|
||||
никогда не пишет listing_segment повторно, поэтому раз рождённая NULL-строка сама
|
||||
себя не чинит даже при живой ежедневной досдаче) деактивируется ОТДЕЛЬНОЙ джобой per
|
||||
source (null_segment_only=True, миграция 264): явный `listing_segment IS NULL`
|
||||
предикат, а не ANY(:segments) -- этот оператор NULL никогда не матчит. Гейт
|
||||
здоровья/пол переобхода для этой джобы выключены (min_confirmations=0,
|
||||
revisit_floor_quantile=0) -- население нерепрезентативно мало (единицы подтверждений
|
||||
в сутки против сотен-тысяч у обычного vtorichka-среза), калиброванный под vtorichka
|
||||
порог держал бы джобу вечно skipped_unhealthy.
|
||||
- Строки НЕ удаляются -- история нужна для бэктеста (#667).
|
||||
- #2674: деактивация в той же транзакции пишет снимок listings_snapshots со статусом
|
||||
'stale' за текущую дату -- «мы N суток не видели». Жёсткое 'closed' (площадка
|
||||
|
|
@ -127,6 +136,11 @@ DEFAULT_MIN_CONFIRMATIONS = 500
|
|||
|
||||
_CONFIRMATIONS_SEGMENT_FILTER = "\n AND listing_segment = ANY(CAST(:segments AS text[]))"
|
||||
|
||||
# NULL-сегмент: `= ANY(...)` НИКОГДА не матчит NULL (SQL, не баг), поэтому
|
||||
# для null_segment_only-режима нужен отдельный явный предикат IS NULL, а не элемент
|
||||
# в :segments. См. _build_null_segment_sql ниже -- тот же принцип для самого UPDATE.
|
||||
_CONFIRMATIONS_NULL_SEGMENT_FILTER = "\n AND listing_segment IS NULL"
|
||||
|
||||
|
||||
# ── Пол TTL по измеренному циклу переобхода (#2659) ───────────────────────────
|
||||
# Гейт выше отвечает на вопрос «источник вообще собирается?». Он НЕ отвечает на
|
||||
|
|
@ -188,9 +202,12 @@ _CONFIRMATIONS_SEGMENT_FILTER = "\n AND listing_segment = ANY(CAST(:seg
|
|||
DEFAULT_REVISIT_FLOOR_QUANTILE = 0.99
|
||||
|
||||
_REVISIT_FLOOR_SEGMENT_FILTER = "\n AND l.listing_segment = ANY(CAST(:segments AS text[]))"
|
||||
_REVISIT_FLOOR_NULL_SEGMENT_FILTER = "\n AND l.listing_segment IS NULL"
|
||||
|
||||
|
||||
def _build_revisit_floor_sql(staleness_column: str, *, with_segments: bool) -> Any:
|
||||
def _build_revisit_floor_sql(
|
||||
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
||||
) -> Any:
|
||||
"""Квантиль возраста, при котором свип за окно ДОКАЗАЛ, что строка жива.
|
||||
|
||||
Пара «предыдущее наблюдение (снимок) → текущее наблюдение (listings)» даёт
|
||||
|
|
@ -198,10 +215,20 @@ def _build_revisit_floor_sql(staleness_column: str, *, with_segments: bool) -> A
|
|||
Только строки, у которых свежесть реально сдвинулась, — то есть выжившие,
|
||||
а не «мы к ним не приходили».
|
||||
|
||||
null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments)
|
||||
(ANY никогда не матчит NULL). На практике для null_segment_only-джобы этот запрос
|
||||
не строится вовсе (revisit_floor_quantile=0 -- см. модульный докстринг), но вариант
|
||||
нужен для корректности, если порог когда-нибудь включат.
|
||||
|
||||
staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings.
|
||||
Значения — param-binding, psycopg v3 safe (CAST(... AS ...), никаких :param::type).
|
||||
"""
|
||||
segment_filter = _REVISIT_FLOOR_SEGMENT_FILTER if with_segments else ""
|
||||
if null_segment_only:
|
||||
segment_filter = _REVISIT_FLOOR_NULL_SEGMENT_FILTER
|
||||
elif with_segments:
|
||||
segment_filter = _REVISIT_FLOOR_SEGMENT_FILTER
|
||||
else:
|
||||
segment_filter = ""
|
||||
return text(
|
||||
f"""
|
||||
SELECT percentile_disc(CAST(:revisit_quantile AS double precision))
|
||||
|
|
@ -229,14 +256,25 @@ def _build_revisit_floor_sql(staleness_column: str, *, with_segments: bool) -> A
|
|||
)
|
||||
|
||||
|
||||
def _build_confirmations_sql(staleness_column: str, *, with_segments: bool) -> Any:
|
||||
def _build_confirmations_sql(
|
||||
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
||||
) -> Any:
|
||||
"""SELECT count(*) подтверждённых за окно строк — тот же срез, что и у UPDATE.
|
||||
|
||||
null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments).
|
||||
Для null_segment_only-джобы min_confirmations=0 по умолчанию (см. модульный
|
||||
докстринг), так что на практике этот путь не строится -- оставлен для корректности.
|
||||
|
||||
staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings.
|
||||
Значения (:listing_source, :health_window_days, :segments) — param-binding,
|
||||
psycopg v3 safe (CAST(... AS ...), никаких :param::type).
|
||||
"""
|
||||
segment_filter = _CONFIRMATIONS_SEGMENT_FILTER if with_segments else ""
|
||||
if null_segment_only:
|
||||
segment_filter = _CONFIRMATIONS_NULL_SEGMENT_FILTER
|
||||
elif with_segments:
|
||||
segment_filter = _CONFIRMATIONS_SEGMENT_FILTER
|
||||
else:
|
||||
segment_filter = ""
|
||||
return text(
|
||||
f"""
|
||||
SELECT count(*)
|
||||
|
|
@ -292,6 +330,33 @@ def _build_segments_sql(staleness_column: str) -> Any:
|
|||
)
|
||||
|
||||
|
||||
def _build_null_segment_sql(staleness_column: str) -> Any:
|
||||
"""UPDATE строго по listing_segment IS NULL (null_segment_only=True).
|
||||
|
||||
НЕ переиспользует _build_segments_sql: `= ANY(CAST(:segments AS text[]))` никогда
|
||||
не матчит NULL (SQL-семантика, не баг -- та же ловушка задокументирована выше у
|
||||
novostroyki-гарда), поэтому NULL-сегмент не выразить через список segments и нужен
|
||||
отдельный явный предикат. Целенаправленно НЕ трогает 'vtorichka'/'novostroyki' --
|
||||
их деактивация идёт через _build_segments_sql в отдельных, уже существующих джобах.
|
||||
|
||||
staleness_column уже прошёл whitelist-проверку. Без :segments-параметра вовсе.
|
||||
"""
|
||||
return text(
|
||||
f"""
|
||||
WITH stale AS (
|
||||
UPDATE listings
|
||||
SET is_active = false
|
||||
WHERE source = :listing_source
|
||||
AND is_active = true
|
||||
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||||
AND listing_segment IS NULL
|
||||
RETURNING id, price_rub
|
||||
)
|
||||
{_STALE_SNAPSHOT_TAIL}
|
||||
"""
|
||||
)
|
||||
|
||||
|
||||
# Дефолтные (last_seen_at) варианты SQL -- сохранены как модульные константы для
|
||||
# обратной совместимости (тесты читают .text, product_handlers/scheduler не менялись).
|
||||
_DEACTIVATE_SQL_ALL_SEGMENTS = _build_all_segments_sql("last_seen_at")
|
||||
|
|
@ -312,6 +377,7 @@ def deactivate_stale_listings(
|
|||
min_confirmations: int = 0,
|
||||
health_window_days: int = _HEALTH_WINDOW_DAYS,
|
||||
revisit_floor_quantile: float = 0.0,
|
||||
null_segment_only: bool = False,
|
||||
) -> dict[str, int]:
|
||||
"""Пометить is_active=false объявления, чья свежесть старше ttl_days дней.
|
||||
|
||||
|
|
@ -321,6 +387,7 @@ def deactivate_stale_listings(
|
|||
ttl_days: количество дней TTL; объявления старше этого порога деактивируются.
|
||||
segments: если задан -- деактивировать только объявления с указанными
|
||||
listing_segment значениями. None -> все сегменты (поведение avito по умолчанию).
|
||||
Несовместимо с null_segment_only=True (см. ниже).
|
||||
staleness_column: колонка-таймстемп, по которой считается свежесть. Whitelist
|
||||
{"last_seen_at", "scraped_at"} — иначе ValueError ДО любого SQL. Дефолт
|
||||
last_seen_at. Для domklik (#2204) — scraped_at: нетрекаемый bulk-touch
|
||||
|
|
@ -337,6 +404,12 @@ def deactivate_stale_listings(
|
|||
эффективный TTL = max(ttl_days, этот пол). 0 -> пол выключен (так
|
||||
вызывают старые тесты и совместимая обёртка), рабочее значение —
|
||||
DEFAULT_REVISIT_FLOOR_QUANTILE, см. комментарий выше.
|
||||
null_segment_only: True -> WHERE фильтрует `listing_segment IS NULL` вместо
|
||||
ANY(:segments). Требует segments=None (иначе ValueError -- смешивать
|
||||
бессмысленно, это два непересекающихся среза). Для этого среза гейт/пол
|
||||
обычно держат выключенными (min_confirmations=0, revisit_floor_quantile=0,
|
||||
см. миграцию 264 и модульный докстринг) -- население слишком мало для
|
||||
откалиброванных под полноценный vtorichka-свип порогов.
|
||||
|
||||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||||
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
|
||||
|
|
@ -349,7 +422,8 @@ def deactivate_stale_listings(
|
|||
|
||||
Raises:
|
||||
ValueError: если staleness_column не входит в whitelist (проверка ДО SQL,
|
||||
никакой интерполяции пользовательского ввода в запрос).
|
||||
никакой интерполяции пользовательского ввода в запрос), либо если заданы
|
||||
одновременно null_segment_only=True и segments (взаимоисключающие срезы).
|
||||
"""
|
||||
counters: dict[str, int] = {"deactivated": 0}
|
||||
try:
|
||||
|
|
@ -362,6 +436,11 @@ def deactivate_stale_listings(
|
|||
f"invalid staleness_column={staleness_column!r}; "
|
||||
f"allowed: {sorted(_ALLOWED_STALENESS_COLUMNS)}"
|
||||
)
|
||||
# null_segment_only + segments одновременно -- неоднозначный запрос:
|
||||
# IS NULL и ANY(:segments) -- разные, непересекающиеся предикаты, а не
|
||||
# композиция. Явный ValueError лучше молчаливого выбора одного из двух.
|
||||
if null_segment_only and segments is not None:
|
||||
raise ValueError("null_segment_only=True несовместимо с заданным segments")
|
||||
|
||||
# Гейт по здоровью сбора (#2659) — ДО любого UPDATE. Деактивация необратима
|
||||
# на практике (вернуть «живость» может только повторный сбор), поэтому
|
||||
|
|
@ -375,7 +454,11 @@ def deactivate_stale_listings(
|
|||
health_params["segments"] = segments
|
||||
confirmations = (
|
||||
db.execute(
|
||||
_build_confirmations_sql(staleness_column, with_segments=segments is not None),
|
||||
_build_confirmations_sql(
|
||||
staleness_column,
|
||||
with_segments=segments is not None,
|
||||
null_segment_only=null_segment_only,
|
||||
),
|
||||
health_params,
|
||||
).scalar()
|
||||
or 0
|
||||
|
|
@ -390,13 +473,15 @@ def deactivate_stale_listings(
|
|||
logger.warning(
|
||||
"deactivate_stale source=%s run_id=%d SKIPPED: сбор нездоров — "
|
||||
"подтверждений за %d сут %d < порога %d "
|
||||
"(segments=%r, staleness_column=%s); ни одна строка не тронута",
|
||||
"(segments=%r, null_segment_only=%s, staleness_column=%s); "
|
||||
"ни одна строка не тронута",
|
||||
listing_source,
|
||||
run_id,
|
||||
health_window_days,
|
||||
confirmations,
|
||||
min_confirmations,
|
||||
segments,
|
||||
null_segment_only,
|
||||
staleness_column,
|
||||
)
|
||||
return counters
|
||||
|
|
@ -413,7 +498,11 @@ def deactivate_stale_listings(
|
|||
if segments is not None:
|
||||
floor_params["segments"] = segments
|
||||
floor_days = db.execute(
|
||||
_build_revisit_floor_sql(staleness_column, with_segments=segments is not None),
|
||||
_build_revisit_floor_sql(
|
||||
staleness_column,
|
||||
with_segments=segments is not None,
|
||||
null_segment_only=null_segment_only,
|
||||
),
|
||||
floor_params,
|
||||
).scalar()
|
||||
# NULL = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
||||
|
|
@ -439,12 +528,21 @@ def deactivate_stale_listings(
|
|||
ttl_days,
|
||||
)
|
||||
|
||||
# segments is None -> все сегменты (поведение avito). segments=[...] -> только
|
||||
# перечисленные сегменты. Используем `is not None` (НЕ truthy): пустой список []
|
||||
# означает "ни один сегмент" (= ANY(ARRAY[]) ничего не матчит, деактивирует 0),
|
||||
# а НЕ "все сегменты" — иначе случайный [] стёр бы весь источник.
|
||||
if segments is not None:
|
||||
# null_segment_only -> IS NULL, отдельный явный предикат (ANY(:segments)
|
||||
# никогда не матчит NULL). segments is None -> все сегменты (поведение avito).
|
||||
# segments=[...] -> только перечисленные сегменты. Используем `is not None`
|
||||
# (НЕ truthy): пустой список [] означает "ни один сегмент" (= ANY(ARRAY[])
|
||||
# ничего не матчит, деактивирует 0), а НЕ "все сегменты" — иначе случайный []
|
||||
# стёр бы весь источник.
|
||||
if null_segment_only:
|
||||
params: dict[str, Any] = {
|
||||
"listing_source": listing_source,
|
||||
"ttl_days": effective_ttl_days,
|
||||
"run_id": run_id,
|
||||
}
|
||||
result = db.execute(_build_null_segment_sql(staleness_column), params)
|
||||
elif segments is not None:
|
||||
params = {
|
||||
"listing_source": listing_source,
|
||||
"ttl_days": effective_ttl_days,
|
||||
"segments": segments,
|
||||
|
|
@ -465,13 +563,15 @@ def deactivate_stale_listings(
|
|||
runs_mod.mark_done(db, run_id, counters)
|
||||
logger.info(
|
||||
"deactivate_stale source=%s run_id=%d done: deactivated=%d "
|
||||
"(ttl_days=%d эффективный, задан %d, segments=%r, staleness_column=%s)",
|
||||
"(ttl_days=%d эффективный, задан %d, segments=%r, null_segment_only=%s, "
|
||||
"staleness_column=%s)",
|
||||
listing_source,
|
||||
run_id,
|
||||
counters["deactivated"],
|
||||
effective_ttl_days,
|
||||
ttl_days,
|
||||
segments,
|
||||
null_segment_only,
|
||||
staleness_column,
|
||||
)
|
||||
return counters
|
||||
|
|
|
|||
|
|
@ -0,0 +1,109 @@
|
|||
-- 264_seed_deactivate_stale_null_segment_yandex_cian.sql
|
||||
-- Деактивация протухших yandex/cian объявлений с ПУСТЫМ listing_segment.
|
||||
--
|
||||
-- Замер на проде 2026-08-15 (is_active=true, listing_segment IS NULL):
|
||||
-- source | активных | старше 30 сут | макс возраст
|
||||
-- yandex | 544 | 533 | 86.3 сут
|
||||
-- cian | 224 | 211 | 86.3 сут
|
||||
-- 97% / 94% этих строк протухли, вплоть до 86 суток. При этом estimator их
|
||||
-- ИСПОЛЬЗУЕТ как comps без freshness-фильтра (Tier A "тот же дом" / Tier C
|
||||
-- micro-radius в app/services/estimator.py фильтруют только is_active=true,
|
||||
-- без scraped_at-фильтра свежести — в отличие от Tier S/H, у которых он есть).
|
||||
--
|
||||
-- ПОЧЕМУ NULL, А НЕ ANY(:segments). deactivate_stale_yandex / deactivate_stale_cian
|
||||
-- (миграция 115) уже деактивируют segments=['vtorichka'] — пустой сегмент они НЕ видят:
|
||||
-- `listing_segment = ANY(CAST(:segments AS text[]))` в SQL никогда не матчит NULL
|
||||
-- (задокументировано в 115 у novostroyki-гарда). Нужен отдельный явный предикат
|
||||
-- IS NULL — app/tasks/deactivate_stale_avito.py получил kwarg null_segment_only=True,
|
||||
-- строящий `... AND listing_segment IS NULL` вместо ANY(:segments).
|
||||
--
|
||||
-- ПОЧЕМУ ОТДЕЛЬНАЯ ДЖОБА, А НЕ РАСШИРЕНИЕ deactivate_stale_yandex/_cian. Гейт
|
||||
-- здоровья сбора (#2659, migration 219) и пол переобхода (#2659) откалиброваны под
|
||||
-- полноценный vtorichka-свип (сотни-тысячи подтверждений в сутки, см. 219). У
|
||||
-- NULL-сегмента подтверждений на 2-3 порядка меньше (замер того же дня: 9 cian +
|
||||
-- 4 yandex строк с last_seen_at < 7 суток) — с общим min_confirmations джоба
|
||||
-- вечно давала бы skipped_unhealthy и никогда не деактивировала бы ни строки.
|
||||
-- Отдельная джоба с собственными (выключенными) порогами не трогает работающие
|
||||
-- deactivate_stale_yandex/_cian и их пол/гейт.
|
||||
--
|
||||
-- НЕ ЗАТРАГИВАЕТ novostroyki: null_segment_only-предикат — строго `IS NULL`, ни
|
||||
-- 'novostroyki', ни 'vtorichka' в него не попадают ни при каких условиях (в отличие
|
||||
-- от паушального TTL по всему source, который снёс бы все ~22,5к первичных строк).
|
||||
--
|
||||
-- TTL=60 суток — консервативный, обоснование числом:
|
||||
-- Строки этого среза по определению не переобходятся систематически (иначе у них
|
||||
-- был бы сегмент — свежий обход cian/yandex SERP всегда вычисляет listing_segment
|
||||
-- детерминированно, см. providers/cian/serp.py:955-958, providers/yandex/serp.py:177).
|
||||
-- Значит «пол переобхода» (revisit_floor, #2659) здесь измерять нечем: он квантиль
|
||||
-- разрывов НАБЛЮДАЕМОГО повторного обхода, а для строки вне скоупа обхода такого
|
||||
-- ряда нет — вычислять его было бы фикцией. Поэтому revisit_floor_quantile=0 явно
|
||||
-- (выключен), а весь запас закладываем в сам TTL:
|
||||
-- deactivate_stale_avito.py документирует измеренные p99 разрывов переобхода
|
||||
-- vtorichka (тот же тип строк, тот же source, разница только в сегменте):
|
||||
-- cian/vtorichka p99 = 26.6 сут
|
||||
-- yandex/vtorichka p99 = 43.0 сут
|
||||
-- TTL=60 даёт запас 2.26x над cian p99 и 1.4x над yandex p99 — комфортный отступ
|
||||
-- без специального замера под null-сегмент (население слишком мало для устойчивого
|
||||
-- перцентиля). При этом бимодальность выборки (замер 2026-08-15: gt30d/gt45d/gt60d
|
||||
-- почти не меняются — 211/211/211 cian, 533/525/523 yandex) означает, что более
|
||||
-- консервативный TTL стоит ПОЧТИ НИЧЕГО в охвате: первый прогон снимет 734 из 768
|
||||
-- строк (95.6%) вместо 744 при TTL=30 — разница 10 строк, зато вдвое больший
|
||||
-- защитный запас над измеренным хвостом обхода.
|
||||
--
|
||||
-- min_confirmations=0, revisit_floor_quantile=0 — оба гейта ВЫКЛЮЧЕНЫ явно (не через
|
||||
-- умолчание product_handlers.py, которое иначе подставило бы DEFAULT_MIN_CONFIRMATIONS
|
||||
-- = 500 и DEFAULT_REVISIT_FLOOR_QUANTILE = 0.99 — оба откалиброваны под другую шкалу
|
||||
-- популяции и держали бы эту джобу в вечном skipped_unhealthy, см. выше).
|
||||
--
|
||||
-- Schedule window 07:00-08:00 UTC — тот же слот, что и deactivate_stale_yandex/_cian
|
||||
-- (migration 115) и deactivate_stale_domklik/_n1 (migration 160): после ночных sweep'ов
|
||||
-- (02:00-05:00 UTC), так что реально переобойдённые строки не деактивируются.
|
||||
--
|
||||
-- next_run_at bootstrapped на завтра 07:00 UTC — тот же паттерн, что 090/115/160,
|
||||
-- чтобы не сработать сразу на деплое.
|
||||
--
|
||||
-- Идемпотентно: ON CONFLICT (source) DO NOTHING — безопасно при повторном применении.
|
||||
--
|
||||
-- Dependencies:
|
||||
-- 052_scrape_schedules.sql (таблица + UNIQUE(source)).
|
||||
-- listings.listing_segment (011_listings_alter.sql).
|
||||
-- 115_scrape_schedules_seed_deactivate_stale_yandex_cian.sql (соседние джобы, тот же слот).
|
||||
-- app/tasks/deactivate_stale_avito.py — null_segment_only kwarg.
|
||||
-- app/services/product_handlers.py — _job_deactivate_stale читает null_segment_only
|
||||
-- из default_params и пробрасывает в deactivate_stale_listings.
|
||||
--
|
||||
-- Deploy order: применять ПОСЛЕ деплоя backend-кода (null_segment_only kwarg), иначе
|
||||
-- первый прогон свалится с TypeError на неизвестный параметр default_params.
|
||||
|
||||
BEGIN;
|
||||
|
||||
INSERT INTO scrape_schedules (
|
||||
source,
|
||||
enabled,
|
||||
window_start_hour,
|
||||
window_end_hour,
|
||||
next_run_at,
|
||||
default_params
|
||||
)
|
||||
VALUES
|
||||
(
|
||||
'deactivate_stale_yandex_null_segment',
|
||||
true, -- SAFE: pure internal DB UPDATE, no ext calls
|
||||
7,
|
||||
8,
|
||||
((CURRENT_DATE + INTERVAL '1 day') + make_interval(hours => 7)) AT TIME ZONE 'UTC',
|
||||
'{"listing_source":"yandex","ttl_days":60,"null_segment_only":true,'
|
||||
'"min_confirmations":0,"revisit_floor_quantile":0}'::jsonb
|
||||
),
|
||||
(
|
||||
'deactivate_stale_cian_null_segment',
|
||||
true, -- SAFE: pure internal DB UPDATE, no ext calls
|
||||
7,
|
||||
8,
|
||||
((CURRENT_DATE + INTERVAL '1 day') + make_interval(hours => 7)) AT TIME ZONE 'UTC',
|
||||
'{"listing_source":"cian","ttl_days":60,"null_segment_only":true,'
|
||||
'"min_confirmations":0,"revisit_floor_quantile":0}'::jsonb
|
||||
)
|
||||
ON CONFLICT (source) DO NOTHING;
|
||||
|
||||
COMMIT;
|
||||
|
|
@ -252,3 +252,4 @@
|
|||
261_listings_search_mv_drop_placeholder_columns.sql
|
||||
262_scrape_schedules_seed_oblast_city_sweeps_wave2.sql
|
||||
263_scrape_schedules_wave2_cian_newbuilding_only_false.sql
|
||||
264_seed_deactivate_stale_null_segment_yandex_cian.sql
|
||||
|
|
|
|||
|
|
@ -424,3 +424,182 @@ def test_migration_160_is_transactional() -> None:
|
|||
def test_migration_160_no_psycopg_trap() -> None:
|
||||
sql = _MIGRATION_160.read_text("utf-8")
|
||||
assert not re.search(r":\w+::", sql)
|
||||
|
||||
|
||||
# ── null_segment_only (пустой listing_segment yandex/cian, никогда не переобходится) ──
|
||||
|
||||
|
||||
def test_null_segment_sql_uses_is_null_not_any() -> None:
|
||||
sql = str(task_mod._build_null_segment_sql("last_seen_at").text)
|
||||
assert "listing_segment IS NULL" in sql
|
||||
assert "ANY(CAST(:segments AS text[]))" not in sql
|
||||
assert ":segments" not in sql
|
||||
|
||||
|
||||
def test_null_segment_sql_filters_is_active_and_source() -> None:
|
||||
sql = str(task_mod._build_null_segment_sql("last_seen_at").text)
|
||||
assert "is_active = true" in sql
|
||||
assert ":listing_source" in sql
|
||||
assert "SET is_active = false" in sql
|
||||
assert "DELETE" not in sql.upper()
|
||||
|
||||
|
||||
def test_null_segment_sql_no_psycopg_trap() -> None:
|
||||
sql = str(task_mod._build_null_segment_sql("last_seen_at").text)
|
||||
assert not re.search(r":\w+::", sql)
|
||||
assert "CAST(:ttl_days || ' days' AS interval)" in sql
|
||||
|
||||
|
||||
def test_null_segment_only_and_segments_raises(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""null_segment_only=True + segments заданы -- неоднозначный запрос, ValueError."""
|
||||
failed: dict[str, Any] = {}
|
||||
monkeypatch.setattr(task_mod.runs_mod, "mark_done", lambda *a, **k: None)
|
||||
monkeypatch.setattr(
|
||||
task_mod.runs_mod,
|
||||
"mark_failed",
|
||||
lambda _db, run_id, err, counters: failed.update(run_id=run_id, err=err),
|
||||
)
|
||||
db = _FakeDB(rowcount=0)
|
||||
with pytest.raises(ValueError, match="null_segment_only"):
|
||||
task_mod.deactivate_stale_listings(
|
||||
db,
|
||||
run_id=20,
|
||||
listing_source="cian",
|
||||
ttl_days=60,
|
||||
segments=["vtorichka"],
|
||||
null_segment_only=True,
|
||||
) # type: ignore[arg-type]
|
||||
assert db.executed == []
|
||||
assert failed["run_id"] == 20
|
||||
|
||||
|
||||
def test_null_segment_only_deactivates_via_is_null(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
marked: dict[str, Any] = {}
|
||||
monkeypatch.setattr(
|
||||
task_mod.runs_mod,
|
||||
"mark_done",
|
||||
lambda _db, run_id, counters: marked.update(run_id=run_id, counters=dict(counters)),
|
||||
)
|
||||
monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||
db = _FakeDB(rowcount=211)
|
||||
out = task_mod.deactivate_stale_listings(
|
||||
db,
|
||||
run_id=21,
|
||||
listing_source="cian",
|
||||
ttl_days=60,
|
||||
null_segment_only=True,
|
||||
) # type: ignore[arg-type]
|
||||
assert out == {"deactivated": 211}
|
||||
assert db.committed is True
|
||||
stmt, params = db.executed[0]
|
||||
sql = _sql_text(stmt)
|
||||
assert "listing_segment IS NULL" in sql
|
||||
assert params is not None
|
||||
assert "segments" not in params
|
||||
assert params["listing_source"] == "cian"
|
||||
assert params["ttl_days"] == 60
|
||||
assert marked["counters"] == {"deactivated": 211}
|
||||
|
||||
|
||||
def test_null_segment_only_confirmations_sql_uses_is_null() -> None:
|
||||
sql = str(
|
||||
task_mod._build_confirmations_sql(
|
||||
"last_seen_at", with_segments=False, null_segment_only=True
|
||||
).text
|
||||
)
|
||||
assert "listing_segment IS NULL" in sql
|
||||
assert ":segments" not in sql
|
||||
|
||||
|
||||
def test_null_segment_only_revisit_floor_sql_uses_is_null() -> None:
|
||||
sql = str(
|
||||
task_mod._build_revisit_floor_sql(
|
||||
"last_seen_at", with_segments=False, null_segment_only=True
|
||||
).text
|
||||
)
|
||||
assert "l.listing_segment IS NULL" in sql
|
||||
assert ":segments" not in sql
|
||||
|
||||
|
||||
def test_null_segment_only_default_is_false(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""Обратная совместимость: старые вызовы без null_segment_only ведут себя как раньше."""
|
||||
monkeypatch.setattr(task_mod.runs_mod, "mark_done", lambda *a, **k: None)
|
||||
monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||
db = _FakeDB(rowcount=3)
|
||||
task_mod.deactivate_stale_listings(
|
||||
db, run_id=22, listing_source="avito", ttl_days=10, segments=None
|
||||
) # type: ignore[arg-type]
|
||||
stmt, _params = db.executed[0]
|
||||
sql = _sql_text(stmt)
|
||||
assert "listing_segment IS NULL" not in sql
|
||||
|
||||
|
||||
# ── Migration 264 (deactivate_stale_yandex_null_segment / _cian_null_segment) ───────
|
||||
|
||||
_MIGRATION_264 = _SQL_DIR / "264_seed_deactivate_stale_null_segment_yandex_cian.sql"
|
||||
|
||||
|
||||
def test_migration_264_exists() -> None:
|
||||
assert _MIGRATION_264.is_file(), f"missing migration: {_MIGRATION_264}"
|
||||
|
||||
|
||||
def test_migration_264_seeds_yandex_and_cian_null_segment() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert "'deactivate_stale_yandex_null_segment'" in sql
|
||||
assert "'deactivate_stale_cian_null_segment'" in sql
|
||||
|
||||
|
||||
def test_migration_264_null_segment_only_true() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert '"null_segment_only":true' in sql
|
||||
|
||||
|
||||
def test_migration_264_ttl_60_days() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert '"ttl_days":60' in sql
|
||||
|
||||
|
||||
def test_migration_264_gates_disabled() -> None:
|
||||
"""min_confirmations/revisit_floor_quantile выключены явно -- население слишком
|
||||
мало для порогов, откалиброванных под полноценный vtorichka-свип (см. файл)."""
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert '"min_confirmations":0' in sql
|
||||
assert '"revisit_floor_quantile":0' in sql
|
||||
|
||||
|
||||
def test_migration_264_is_idempotent() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert "ON CONFLICT (source) DO NOTHING" in sql
|
||||
|
||||
|
||||
def test_migration_264_is_transactional() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert "BEGIN;" in sql
|
||||
assert "COMMIT;" in sql
|
||||
|
||||
|
||||
def test_migration_264_enabled_true() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert "true" in sql
|
||||
|
||||
|
||||
def test_migration_264_window_7_to_8_utc() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert re.search(r"\b7\b", sql), "window_start_hour 7 missing"
|
||||
assert re.search(r"\b8\b", sql), "window_end_hour 8 missing"
|
||||
|
||||
|
||||
def test_migration_264_no_psycopg_trap() -> None:
|
||||
sql = _MIGRATION_264.read_text("utf-8")
|
||||
assert not re.search(r":\w+::", sql)
|
||||
|
||||
|
||||
def test_handler_wires_null_segment_only_from_schedule_params() -> None:
|
||||
"""Читаем исходник файлом (как test_handler_wires_revisit_floor_from_schedule_params):
|
||||
product_handlers тянет scraper_kit, которого в юнит-окружении может не быть."""
|
||||
handlers = Path(__file__).resolve().parents[1] / "app" / "services" / "product_handlers.py"
|
||||
src = handlers.read_text("utf-8")
|
||||
job = src.split("async def _job_deactivate_stale")[1].split("\nasync def ")[0]
|
||||
flat = " ".join(job.split())
|
||||
assert 'params.get("null_segment_only", False)' in flat
|
||||
assert "null_segment_only=null_segment_only" in job
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue