fix(tradein/deactivate): протухшие объявления с пустым сегментом больше не якорят оценку #2908
5 changed files with 426 additions and 25 deletions
|
|
@ -237,6 +237,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)
|
||||
# Потолок эффективного TTL (см. CAP_MULT в deactivate_stale_avito.py) — множитель,
|
||||
# а не голая константа: источник с непропорционально длинным хвостом переобхода
|
||||
# относительно своего ttl_days переопределяет его через default_params (ключ
|
||||
|
|
@ -255,6 +260,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,
|
||||
cap_mult=cap_mult,
|
||||
),
|
||||
)
|
||||
|
|
|
|||
|
|
@ -6,18 +6,28 @@
|
|||
Ключевые решения:
|
||||
- Cian/Yandex не поддерживают full-coverage sweep -> паушальный TTL сломает живой
|
||||
инвентарь. DECISION: для yandex/cian деактивировать ТОЛЬКО listing_segment='vtorichka',
|
||||
TTL=30. novostroyki (9659 активных первичных строк) и NULL-сегмент не трогаем.
|
||||
TTL=30. novostroyki (9659 активных первичных строк) не трогаем.
|
||||
ИЗВЕСТНЫЙ ПРОБЕЛ (ревью TTL-CAP круг 2, 2026-08-15): этот скоуп уже, чем множество
|
||||
реально протухших строк -- живой замер на проде даёт cian/novostroyki 9 483,
|
||||
cian/NULL-сегмент 211, yandex/NULL-сегмент 523 активных строки старше 60 суток, ни
|
||||
одна из них не деактивируется НИ ОДНОЙ джобой (внутри скоупа cian/vtorichka и
|
||||
yandex/vtorichka таких строк 0). Потолок cap_mult (см. CAP_MULT ниже) этот пробел
|
||||
не закрывает и закрыть не может -- он сжимает пул ВНУТРИ скоупа джобы, а не
|
||||
расширяет сам скоуп. Расширение скоупа -- отдельная задача (нужно сперва выяснить,
|
||||
поддерживают ли cian/yandex full-coverage sweep для novostroyki/NULL-сегмента
|
||||
СЕЙЧАС, иначе паушальный TTL повторит инцидент, ради которого этот DECISION и
|
||||
принят) и намеренно НЕ входит в TTL-CAP.
|
||||
реально протухших строк -- живой замер на проде даёт cian/novostroyki 9 483 активных
|
||||
строки старше 60 суток, ни одна из них не деактивируется НИ ОДНОЙ джобой (внутри
|
||||
скоупа cian/vtorichka и yandex/vtorichka таких строк 0). Потолок cap_mult (см.
|
||||
CAP_MULT ниже) этот пробел не закрывает и закрыть не может -- он сжимает пул ВНУТРИ
|
||||
скоупа джобы, а не расширяет сам скоуп. NULL-сегмент (тот же замер круга 2 давал
|
||||
cian/NULL 211, yandex/NULL 523 строки старше 60 суток) закрыт отдельно ниже
|
||||
(null_segment_only, миграция 266) -- novostroyki-часть пробела остаётся: расширение
|
||||
скоупа туда отдельная задача (нужно сперва выяснить, поддерживает ли cian/yandex
|
||||
full-coverage sweep для novostroyki СЕЙЧАС, иначе паушальный TTL повторит инцидент,
|
||||
ради которого этот DECISION и принят) и намеренно НЕ входит в TTL-CAP.
|
||||
- avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений.
|
||||
- NULL-сегмент (легаси-строки до миграции 011 + жертвы бага в ON CONFLICT -- upsert
|
||||
никогда не пишет listing_segment повторно, поэтому раз рождённая NULL-строка сама
|
||||
себя не чинит даже при живой ежедневной досдаче) деактивируется ОТДЕЛЬНОЙ джобой per
|
||||
source (null_segment_only=True, миграция 266): явный `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' (площадка
|
||||
|
|
@ -137,6 +147,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) ───────────────────────────
|
||||
# Гейт выше отвечает на вопрос «источник вообще собирается?». Он НЕ отвечает на
|
||||
|
|
@ -198,6 +213,7 @@ _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"
|
||||
|
||||
|
||||
# ── Потолок эффективного TTL (положительная обратная связь пола, найдено 2026-08-15) ──
|
||||
|
|
@ -260,7 +276,9 @@ _REVISIT_FLOOR_SEGMENT_FILTER = "\n AND l.listing_segment = ANY(CAST(:s
|
|||
CAP_MULT = 2
|
||||
|
||||
|
||||
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)» даёт
|
||||
|
|
@ -268,10 +286,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))
|
||||
|
|
@ -299,14 +327,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(*)
|
||||
|
|
@ -362,6 +401,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")
|
||||
|
|
@ -382,6 +448,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,
|
||||
cap_mult: float = CAP_MULT,
|
||||
) -> dict[str, int]:
|
||||
"""Пометить is_active=false объявления, чья свежесть старше ttl_days дней.
|
||||
|
|
@ -392,6 +459,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
|
||||
|
|
@ -409,13 +477,21 @@ def deactivate_stale_listings(
|
|||
пол поднимает TTL, но не выше потолка. 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,
|
||||
см. миграцию 266 и модульный докстринг) -- население слишком мало для
|
||||
откалиброванных под полноценный vtorichka-свип порогов.
|
||||
cap_mult: множитель потолка эффективного TTL (см. комментарий у модульной
|
||||
константы CAP_MULT). Дефолт -- сама CAP_MULT=2, но параметр, а НЕ голая
|
||||
константа: источник с непропорционально длинным хвостом переобхода
|
||||
относительно своего ttl_days (avito: p99=42.1 при ttl=10 -> дефолтный
|
||||
потолок 20 режет ниже хвоста) может переопределить его через
|
||||
default_params расписания (ключ "cap_mult"), не трогая остальные
|
||||
источники. Итоговый потолок = ttl_days * cap_mult.
|
||||
источники. Итоговый потолок = ttl_days * cap_mult. Применяется и к
|
||||
null_segment_only-джобе, но там гейт/пол выключены (см. выше), так что
|
||||
на практике не участвует.
|
||||
|
||||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||||
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
|
||||
|
|
@ -448,7 +524,9 @@ def deactivate_stale_listings(
|
|||
то есть потолок = сам ttl_days и пол молча отключается, никакого ValueError.
|
||||
jsonb `true`/`false` вместо числа -- ровно та опечатка в расписании, ради
|
||||
которой оба guard'а вообще написаны, поэтому bool отклоняется явной
|
||||
type-проверкой ДО числового сравнения для обоих параметров).
|
||||
type-проверкой ДО числового сравнения для обоих параметров), ЛИБО если
|
||||
заданы одновременно null_segment_only=True и segments (взаимоисключающие
|
||||
срезы -- IS NULL и ANY(:segments) не композируются).
|
||||
"""
|
||||
counters: dict[str, int] = {"deactivated": 0}
|
||||
try:
|
||||
|
|
@ -483,6 +561,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. Деактивация необратима
|
||||
# на практике (вернуть «живость» может только повторный сбор), поэтому
|
||||
|
|
@ -496,7 +579,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
|
||||
|
|
@ -511,13 +598,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
|
||||
|
|
@ -536,7 +625,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 = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
||||
|
|
@ -592,12 +685,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,
|
||||
|
|
@ -618,13 +720,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 @@
|
|||
-- 266_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;
|
||||
|
|
@ -254,3 +254,4 @@
|
|||
263_scrape_schedules_wave2_cian_newbuilding_only_false.sql
|
||||
264_deactivate_stale_avito_cap_mult.sql
|
||||
265_deactivate_stale_yandex_cap_mult.sql
|
||||
266_seed_deactivate_stale_null_segment_yandex_cian.sql
|
||||
|
|
|
|||
|
|
@ -424,3 +424,184 @@ 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 266 (deactivate_stale_yandex_null_segment / _cian_null_segment) ───────
|
||||
# Renumbered 264 -> 266 (collision with forgejo/main's 264_deactivate_stale_avito_cap_mult.sql
|
||||
# / 265_deactivate_stale_yandex_cap_mult.sql, merged после того как эта ветка забрала 264).
|
||||
|
||||
_MIGRATION_266 = _SQL_DIR / "266_seed_deactivate_stale_null_segment_yandex_cian.sql"
|
||||
|
||||
|
||||
def test_migration_266_exists() -> None:
|
||||
assert _MIGRATION_266.is_file(), f"missing migration: {_MIGRATION_266}"
|
||||
|
||||
|
||||
def test_migration_266_seeds_yandex_and_cian_null_segment() -> None:
|
||||
sql = _MIGRATION_266.read_text("utf-8")
|
||||
assert "'deactivate_stale_yandex_null_segment'" in sql
|
||||
assert "'deactivate_stale_cian_null_segment'" in sql
|
||||
|
||||
|
||||
def test_migration_266_null_segment_only_true() -> None:
|
||||
sql = _MIGRATION_266.read_text("utf-8")
|
||||
assert '"null_segment_only":true' in sql
|
||||
|
||||
|
||||
def test_migration_266_ttl_60_days() -> None:
|
||||
sql = _MIGRATION_266.read_text("utf-8")
|
||||
assert '"ttl_days":60' in sql
|
||||
|
||||
|
||||
def test_migration_266_gates_disabled() -> None:
|
||||
"""min_confirmations/revisit_floor_quantile выключены явно -- население слишком
|
||||
мало для порогов, откалиброванных под полноценный vtorichka-свип (см. файл)."""
|
||||
sql = _MIGRATION_266.read_text("utf-8")
|
||||
assert '"min_confirmations":0' in sql
|
||||
assert '"revisit_floor_quantile":0' in sql
|
||||
|
||||
|
||||
def test_migration_266_is_idempotent() -> None:
|
||||
sql = _MIGRATION_266.read_text("utf-8")
|
||||
assert "ON CONFLICT (source) DO NOTHING" in sql
|
||||
|
||||
|
||||
def test_migration_266_is_transactional() -> None:
|
||||
sql = _MIGRATION_266.read_text("utf-8")
|
||||
assert "BEGIN;" in sql
|
||||
assert "COMMIT;" in sql
|
||||
|
||||
|
||||
def test_migration_266_enabled_true() -> None:
|
||||
sql = _MIGRATION_266.read_text("utf-8")
|
||||
assert "true" in sql
|
||||
|
||||
|
||||
def test_migration_266_window_7_to_8_utc() -> None:
|
||||
sql = _MIGRATION_266.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_266_no_psycopg_trap() -> None:
|
||||
sql = _MIGRATION_266.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