gendesign/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py
bot-backend ab01f7cc48
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-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 3m1s
fix(tradein): убрать невыводимые события, развести «снято» и «протухло» (#2674)
Ревью PR #2682 нашло контрольную группу в наших же данных. Перепроверено
собственными запросами к проду — сходится, местами хуже заявленного.

1. delisted/relisted УБРАНЫ из писателя событий.
   Покрытие обхода за 14-18.07: domklik 99.9-100%, yandex 34-43%, cian 21-27%,
   avito 1.6-3.4%. Переходы за те же дни: domklik — снятий 1/2/0/2/4 в сутки и
   возвратов РОВНО 0 все пять суток; yandex — снятий 343-433 в сутки. Тот же
   обход, тот же день, разница только в покрытии: событие рождается тем, что
   скрейпер снова дошёл, а не тем, что объявление вернулось. Подтверждения:
   avito 13.07 (день остановки обхода) — 3023 «снятия» за сутки против
   контрольной ставки 1-4 (точность ≈4%); 4705 возвратов из 5493 за 12 дней
   (85.7%) — это 2-3.08, два дня после возобновления обхода.
   Сужение окна свежести сделало бы хуже (больше флапаний). Журнал из догадок
   хуже пустого журнала — не пишем. is_active убран из запроса целиком.
   Гейт-тест ослаблен до трёх типов + новый гейт «невыводимые НЕ пишутся».

2. TTL-путь пишет 'stale', а не 'closed'.
   Прогон по домклику 02.08 деактивировал 6131 объявление за раз (TTL 14 суток
   против 12 суток простоя обхода) — под общим статусом это 6131 фальшивая
   «дата продажи» одной датой. 'closed' остаётся только за 404: там ответила
   площадка. CHECK на колонке нет, миграция 212 обновляет только COMMENT.

3. change_time усечён до суток (date_trunc). С now() UNIQUE(source, change_time,
   type) работал только внутри прогона: второй прогон в те же сутки (2 августа
   их было два) давал дубли. Теперь заявленная идемпотентность действительно
   работает.

4. Комнатность в разборе заголовка стала необязательной: 1991 заголовок из
   25 055 (7.9%) — «Квартира-студия, 34,2 м², 9/10 эт.», обязательная группа
   роняла match и обнуляла все четыре поля. Чинит обоих писателей сразу
   (house_suggestions + house_placement_history, там 8.8% без площади).
   Студия → rooms=0 по конвенции kit'а, а не None.

Фальсификация: вернуть delisted — 1 красный; 'closed' на TTL-пути — 6;
обязательная комнатность — 2; now() вместо date_trunc — 1.
2026-08-06 03:01:05 +05:00

240 lines
14 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""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).
- #2674: деактивация в той же транзакции пишет снимок listings_snapshots со статусом
'stale' за текущую дату -- «мы N суток не видели». Жёсткое 'closed' (площадка
ответила 404) пишет только avito_detail_backfill: смешивать факт с догадкой дорого.
Задача синхронная (DB-only, никаких внешних HTTP-вызовов) -- запускается kit-scheduler'ом
через product_handlers._job_deactivate_stale (wildcard-handler deactivate_stale_*),
по образцу 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__)
# Разрешённые колонки-таймстемпы для проверки свежести (staleness).
# Whitelist ЖЁСТКО ограничивает f-string-интерполяцию имени колонки в SQL:
# значение колонки НИКОГДА не берётся из пользовательского ввода без этой проверки.
# - last_seen_at: дефолт; двигается при каждом «увиденном» объявлении в прогоне.
# - scraped_at: честная свежесть для источников с нетрекаемым bulk-touch
# last_seen_at (domklik: #2204 — bulk UPDATE двигает last_seen_at всем строкам
# одним timestamp, поэтому TTL по last_seen_at бесполезен; scraped_at двигает
# только реальный скрейп).
_ALLOWED_STALENESS_COLUMNS = frozenset({"last_seen_at", "scraped_at"})
# ── Снимок «протухло» в дневной истории (#2674) ───────────────────────────────
# listings_snapshots.status до этого фикса был константой 'active' у всех строк
# (394 299 на момент находки) — оба места вызова upsert_listing_snapshot передавали
# литерал 'active', и это честно: там объявление ДЕЙСТВИТЕЛЬНО видели. А деактивация
# TTL-задачей не оставляла в истории вообще никакого следа. Из-за этого дата снятия
# объявления (лучший доступный сигнал «скорее всего продано») не запрашивалась из
# истории, а восстанавливалась на глаз: последний показ + предполагаемый срок жизни.
#
# ПОЧЕМУ 'stale', А НЕ 'closed'. Эта задача НЕ знает, что объявление снято, — она
# знает только, что МЫ его N суток не видели, а это разные факты, когда TTL короче
# простоя обхода. Замер: прогон по домклику 02.08 снял 6131 объявление за раз (TTL
# 14 суток против 12 суток простоя обхода) — с общим статусом это были бы 6131
# фальшивая «дата продажи» одной датой. Продукт про цены, смешивать факт с догадкой
# дорого. Поэтому:
# 'closed' — только путь 404: площадка ответила «нет» (avito_detail_backfill);
# 'stale' — этот путь: «мы N суток не смотрели».
# Дата всё равно фиксируется, но читатель отличает одно от другого. Ограничения
# CHECK на колонке нет (проверено на проде), миграция не нужна — только COMMENT.
#
# Пишем снимок в ТОЙ ЖЕ транзакции, что и UPDATE флага: деактивация без снимка (или
# наоборот) невозможна по построению — один statement, data-modifying CTE.
# 1:1 по строкам: `stale` возвращает уникальные listings.id (PK), каждая даёт ровно
# одну затронутую строку listings_snapshots (INSERT либо DO UPDATE — оба считаются
# в rowcount), поэтому rowcount statement'а по-прежнему равен числу деактивированных.
# price_rub берём из listings (NOT NULL в схеме) — это последняя известная цена.
# ON CONFLICT: если снимок за сегодня уже есть (объявление видели активным утром,
# а вечером сработал TTL) — только переводим статус в 'stale', цену не переписываем.
_STALE_SNAPSHOT_TAIL = """
INSERT INTO listings_snapshots
(listing_id, snapshot_date, run_id, price_rub, status, observed_at)
SELECT id, CURRENT_DATE, CAST(:run_id AS bigint), price_rub, 'stale', NOW()
FROM stale
ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET
status = 'stale',
observed_at = EXCLUDED.observed_at
"""
def _build_all_segments_sql(staleness_column: str) -> Any:
"""UPDATE без фильтра по сегменту: все сегменты для данного source.
staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings,
поэтому f-string-подстановка имени колонки безопасна. Значения (:listing_source,
:ttl_days, :run_id) остаются param-binding — psycopg v3 safe (никаких :param::type).
"""
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)
RETURNING id, price_rub
)
{_STALE_SNAPSHOT_TAIL}
"""
)
def _build_segments_sql(staleness_column: str) -> Any:
"""UPDATE с фильтром по сегменту (segments задан): только указанные сегменты.
staleness_column уже прошёл whitelist-проверку. = ANY(CAST(:segments AS text[]))
-- psycopg v3 адаптирует Python list -> text[].
"""
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 = ANY(CAST(:segments AS text[]))
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")
_DEACTIVATE_SQL_SEGMENTS = _build_segments_sql("last_seen_at")
# Алиас для обратной совместимости -- использовался в тестах через 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,
staleness_column: str = "last_seen_at",
) -> dict[str, int]:
"""Пометить is_active=false объявления, чья свежесть старше ttl_days дней.
Параметры:
listing_source: значение поля source в таблице listings ('avito', 'yandex',
'cian', 'domklik').
ttl_days: количество дней TTL; объявления старше этого порога деактивируются.
segments: если задан -- деактивировать только объявления с указанными
listing_segment значениями. None -> все сегменты (поведение avito по умолчанию).
staleness_column: колонка-таймстемп, по которой считается свежесть. Whitelist
{"last_seen_at", "scraped_at"} — иначе ValueError ДО любого SQL. Дефолт
last_seen_at. Для domklik (#2204) — scraped_at: нетрекаемый bulk-touch
двигает last_seen_at всем строкам одним timestamp, поэтому честная
свежесть = scraped_at (двигается только реальным скрейпом).
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
(data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed).
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
Raises:
ValueError: если staleness_column не входит в whitelist (проверка ДО SQL,
никакой интерполяции пользовательского ввода в запрос).
"""
counters: dict[str, int] = {"deactivated": 0}
try:
# Whitelist-проверка ДО построения/выполнения SQL: только после неё имя колонки
# интерполируется f-string'ом. Значения по-прежнему идут через param-binding.
# Внутри try -> невалидная колонка финализирует run как failed (mark_failed),
# а не оставляет его «running» (единый контракт с прочими сбоями).
if staleness_column not in _ALLOWED_STALENESS_COLUMNS:
raise ValueError(
f"invalid staleness_column={staleness_column!r}; "
f"allowed: {sorted(_ALLOWED_STALENESS_COLUMNS)}"
)
# 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,
"run_id": run_id,
}
result = db.execute(_build_segments_sql(staleness_column), params)
else:
params = {
"listing_source": listing_source,
"ttl_days": ttl_days,
"run_id": run_id,
}
result = db.execute(_build_all_segments_sql(staleness_column), 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, staleness_column=%s)",
listing_source,
run_id,
counters["deactivated"],
ttl_days,
segments,
staleness_column,
)
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,
)