gendesign/tradein-mvp/backend/app/tasks/listing_source_snapshot.py
bot-backend e4d44025e5
Some checks failed
CI Trade-In / changes (pull_request) Successful in 26s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 31s
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) Failing after 7m37s
fix(tradein/snapshot): снимок источника пишется только при изменении (#2993, PR-C)
listing_source_snapshot каждую ночь копировал listing_sources целиком: прод
14-17.09 — 290-293 тыс. строк в сутки, из них отличались от предыдущего
снимка 7.7-11.4 тыс. (3-4 %). Решение владельца 2026-08-23 — строка на
изменение; читатель (пол переобхода, PR-A #3056) и рельсы объёма (PR-B
#3066) уже смержены.

Строка пишется, если снимка ещё нет или (price_rub, is_active, last_seen_at,
payload_hash) отличается от последнего снимка источника, включая сегодняшний
(повторный прогон в те же сутки перезаписывает строку, только если значение
снова сдвинулось). Состояние на дату D = последняя строка с snapshot_date <= D
— для всех четырёх колонок то же, что дала бы суточная копия, поэтому пол
переобхода и события не меняются. Старые суточные снимки не трогаются.

Замер чтением на проде 17.09: отбор с условием — 1.0 с на 293 602 источника
(бюджет 900 с). Миграция 325 — только COMMENT на таблицу и view: прежние
«ежедневный снимок»/«на дату снимка» стали бы неправдой.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-17 14:46:36 +05:00

354 lines
24 KiB
Python
Raw 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.

"""Daily per-source snapshot writer (#570).
Берёт текущее состояние listing_sources (последний снимок на canonical listing × source)
и раз в сутки пишет в listing_source_snapshots строку на КАЖДОЕ ИЗМЕНЕНИЕ (#2993, см.
_SNAPSHOT_SQL; раньше — полная копия каждые сутки) + change-log в
listing_source_events. Так история per-source цены копится СРАЗУ, независимо от
(сейчас DORMANT) скраперов — см. шапку data/sql/079_listing_source_history.sql.
Задача синхронная (DB-only, никаких внешних HTTP-вызовов) — запускается kit-scheduler'ом
через product_handlers._job_listing_source_snapshot,
по образцу import_rosreestr_dkp (sync task в run_in_executor).
Вся работа — два set-based SQL statement'а (snapshot upsert + event-diff), никакого
row-by-row Python.
#2607 — root cause висящих прогонов (ежедневный zombie с минимум 19 июля, всегда ровно 6h
до zombie-порога): event-diff раньше писал "prior" как CTE `DISTINCT ON (listing_source_id)
... ORDER BY listing_source_id, snapshot_date DESC` по ВСЕЙ listing_source_snapshots (~2.6-2.8M
строк) и джойнил её с "today" через обычный JOIN. Планировщик оценивает `snapshot_date =
CURRENT_DATE` в 1 строку (статистика ANALYZE ещё не видела свежевставленные в этой же
транзакции строки today — CURRENT_DATE всегда за пределами гистограммы), выбирает Nested
Loop БЕЗ Materialize на внутренней стороне и на КАЖДУЮ реальную строку today (~80-140k)
заново пересчитывает DISTINCT ON по всей таблице (Unique + Index Scan ~2.7M строк) —
EXPLAIN на проде показал cost≈300k именно на этом шаге. Реально это никогда не завершалось
за 6h, оставляя backend 'active' на сутки после того как zombie-детектор помечал
scrape_runs.status='zombie' (детектор НЕ убивает backend, см. reap_zombies) — держало
backend_xmin, блокируя autovacuum на listings/listing_sources.
Fix: `prior` переписан через `JOIN LATERAL (... ORDER BY snapshot_date DESC LIMIT 1) ON true`
— форсирует per-row индексный point-lookup по idx_lss_source_date (listing_source_id,
snapshot_date DESC) вместо полного DISTINCT ON по таблице; EXPLAIN на проде: cost внутреннего
подзапроса упал с ~298 627 до ~4.4 за строку today. Плюс defense-in-depth: budget_sec →
SET LOCAL statement_timeout (см. snapshot_listing_sources) — если что-то опять разрегрессирует
план, прогон честно падает в mark_failed вместо того чтобы висеть сутками.
"""
from __future__ import annotations
import logging
from typing import Any
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.services import scrape_runs as runs_mod
logger = logging.getLogger(__name__)
# Окно свежести: источник считается активным, если last_seen_at не старше N дней.
FRESHNESS_WINDOW_DAYS = 7
# ── Wall-clock budget (#2607 п.4) ─────────────────────────────────────────────
# Задача не батчится Python-циклом (два set-based statement'а) — единственный способ
# гарантированно оборвать зависший statement это Postgres-нативный statement_timeout,
# выставленный SET LOCAL (per-transaction scope, НЕ трогает server/role-level timeout —
# это issue #2607 п.2, отдельное решение с согласованием). По образцу budget_sec из
# app/tasks/geocode_missing.py (run_geocode_missing_listings), только здесь это не Python
# loop-budget, а SQL statement_timeout.
# Default/clamp: см. data/sql/202_listing_source_snapshot_budget_sec.sql (default_params
# budget_sec=900 — 15 мин, с большим запасом над ожидаемым временем выполнения после
# LATERAL-фикса (секунды) и далеко от 6h zombie-порога).
DEFAULT_BUDGET_SEC = 900.0
_MIN_BUDGET_SEC = 30.0
_MAX_BUDGET_SEC = 3600.0 # hard ceiling — не даём budget_sec случайно воссоздать "висит вечно"
def _clamp_budget_sec(raw: Any) -> float:
"""Валидировать/зажать budget_sec из default_params — защита от 0/отрицательного/мусора."""
try:
val = float(raw)
except (TypeError, ValueError):
val = DEFAULT_BUDGET_SEC
return max(_MIN_BUDGET_SEC, min(val, _MAX_BUDGET_SEC))
# ── Snapshot upsert: строка только на изменение (#2993) ───────────────────────
# Снимок на (listing_source_id, CURRENT_DATE) пишется, только если у источника ещё нет
# ни одного снимка ИЛИ хоть одно из четырёх значений (price_rub, is_active, last_seen_at,
# payload_hash) отличается от его ПОСЛЕДНЕГО снимка. Было: полная копия listing_sources
# каждые сутки, прод 14-17.09 — 290-293 тыс. строк в сутки, из них отличались от
# предыдущего снимка 7.7-11.4 тыс. (3-4 %). Решение владельца 2026-08-23.
#
# Состояние источника на дату D = его последний снимок с snapshot_date <= D. Для всех
# четырёх колонок это то же значение, что записала бы суточная модель: сутки, в которые
# ни одна из них не менялась, и есть пропущенные строки. Поэтому last_seen_at обязан
# быть в сравнении: пол переобхода (deactivate_stale_avito._revisit_floor_from_where_sql,
# PR-A #3056) берёт last_seen_at последнего снимка не позже якоря — без него пары
# «прошлое наблюдение → текущее» считались бы от устаревшей свежести. is_active там же:
# он выводится из now(), и переход «свежий → протух» при неизменном last_seen_at — тоже
# изменение состояния на дату.
#
# p — последний снимок ВКЛЮЧАЯ сегодняшний: повторный прогон в те же сутки сравнивает
# с тем, что уже записано сегодня, и перезаписывает строку только если значение снова
# сдвинулось (ON CONFLICT → last-write-wins, как раньше). Пер-строчный LATERAL
# point-lookup по PK, как в event-diff (#2607); прод 17.09: 293 602 lookup'а — 1.0 с.
# Старые суточные снимки не трогаются — история до перехода остаётся как есть.
#
# is_active derived: last_seen_at в пределах окна свежести на момент снимка.
# payload_hash = md5(raw_payload::text) — ::text на колонке допустим (это не bind-param).
# run_id через CAST(:run_id AS bigint) — psycopg v3 (никогда :run_id::bigint).
_SNAPSHOT_SQL = text(
"""
INSERT INTO listing_source_snapshots (
listing_source_id, snapshot_date, price_rub, is_active,
last_seen_at, payload_hash, observed_at, run_id
)
SELECT
cur.id,
CURRENT_DATE,
cur.price_rub,
cur.is_active,
cur.last_seen_at,
cur.payload_hash,
now(),
CAST(:run_id AS bigint)
FROM (
SELECT
id,
price_rub,
(last_seen_at > now() - make_interval(days => :freshness_days)) AS is_active,
last_seen_at,
md5(raw_payload::text) AS payload_hash
FROM listing_sources
) cur
LEFT JOIN LATERAL (
SELECT s.snapshot_date, s.price_rub, s.is_active, s.last_seen_at, s.payload_hash
FROM listing_source_snapshots s
WHERE s.listing_source_id = cur.id
ORDER BY s.snapshot_date DESC
LIMIT 1
) p ON true
WHERE p.snapshot_date IS NULL
OR (cur.price_rub, cur.is_active, cur.last_seen_at, cur.payload_hash)
IS DISTINCT FROM (p.price_rub, p.is_active, p.last_seen_at, p.payload_hash)
ON CONFLICT (listing_source_id, snapshot_date) DO UPDATE SET
price_rub = EXCLUDED.price_rub,
is_active = EXCLUDED.is_active,
last_seen_at = EXCLUDED.last_seen_at,
payload_hash = EXCLUDED.payload_hash,
observed_at = EXCLUDED.observed_at,
run_id = EXCLUDED.run_id
"""
)
# ── Event diff: три выводимых типа событий из пяти в схеме ────────────────────
# Для каждого источника сравниваем сегодняшний снимок (snapshot_date = CURRENT_DATE) с
# самым свежим ПРЕДЫДУЩИМ (snapshot_date < CURRENT_DATE).
# today — снимок за сегодня (только что записан _SNAPSHOT_SQL, в той же транзакции).
# С #2993 здесь только изменившиеся и новые источники: у остальных все четыре
# значения равны последнему снимку, а значит ни одно событие ниже сработать
# не могло бы — набор событий тот же, что при суточной копии.
# p — последний снимок строго ДО сегодня, per-row LATERAL point-lookup (#2607).
#
# #2674: схема (079) знает пять типов событий, писатель умел один — price_change,
# 8288 строк. Дописаны два:
# edited — payload_hash изменился, а цена нет (изменение цены уже описано
# отдельным событием price_change — дублировать его как «редактирование»
# значило бы считать одно изменение дважды). Прошлый хеш обязан быть
# непустым: md5(NULL) = NULL, и «payload появился впервые» — это не
# правка, а первое наблюдение;
# first_seen — предыдущего снимка нет вовсе (LEFT JOIN LATERAL даёт p.* = NULL).
#
# delisted и relisted НЕ ПИШУТСЯ НАМЕРЕННО — они НЕ ВЫВОДИМЫ из наших данных.
# is_active в снимке — derived-признак «last_seen_at свежее FRESHNESS_WINDOW_DAYS»,
# то есть «мы видели», а не «объявление есть на площадке». При покрытии обхода 10-35%
# такой переход рождается тем, что скрейпер СНОВА ДОШЁЛ до источника, а не тем, что
# объявление вернулось/ушло. Контрольная группа в наших же данных (14-18.07):
# domklik, покрытие 99.9-100%: снятий 1/2/0/2/4 в сутки, возвратов — РОВНО 0 все дни;
# yandex, покрытие 34-43%: снятий 343-433 в сутки, возвратов до 155.
# Тот же обход, тот же день — разница только в покрытии. Отсюда же всплески:
# avito 13.07 (день остановки обхода) — 3023 «снятия» за сутки против контрольной
# ставки 1-4, точность события ≈4%; 4705 «возвратов» из 5493 за 12 дней (86%) — это
# два дня после возобновления обхода 2-3.08.
# Сузить окно свежести НЕ поможет — станет хуже (больше флапаний); окно шире
# максимального интервала повторного визита обессмысливает само событие.
# Честный ответ схеме — не писать эти два типа, а не наполнять журнал догадками.
# Единственный жёсткий сигнал снятия — 404 при поштучном обходе, он пишется в
# listings_snapshots.status='closed' (avito_detail_backfill).
#
# Оставшиеся три события утверждают факты о НАШИХ СОБСТВЕННЫХ строках («появился новый
# источник», «хеш изменился при той же цене», «цена другая»), а не о поведении площадки.
#
# JOIN → LEFT JOIN LATERAL: без LEFT источники без предыдущего снимка отбрасывались
# join'ом, поэтому first_seen был недостижим по построению. План #2607 не меняется —
# LEFT JOIN LATERAL так же форсирует per-row индексный point-lookup по
# idx_lss_source_date, просто не отбрасывает строку при отсутствии предыдущей.
#
# Ветки разворачиваются CROSS JOIN LATERAL (VALUES ...) — одна строка сравнения даёт
# до трёх строк-кандидатов, из которых WHERE e.fires оставляет сработавшие. Это
# по-прежнему ОДИН set-based statement (никакого Python-цикла), просто три предиката
# вместо одного.
#
# #2607: раньше `p` был отдельным CTE `DISTINCT ON (listing_source_id) ... FROM
# listing_source_snapshots WHERE snapshot_date < CURRENT_DATE` и джойнился обычным JOIN.
# Планировщик оценивает `today` в 1 строку (свежевставленные в этой же транзакции строки
# ANALYZE ещё не видел) → Nested Loop БЕЗ Materialize на внутренней стороне → DISTINCT ON
# по ВСЕЙ таблице (~2.6-2.8M строк, Index Scan + Unique) пересчитывался ЗАНОВО на каждую
# из ~80-140k реальных строк today — на проде EXPLAIN показал cost≈300k на этом шаге,
# запрос не укладывался ни в 6h zombie-порог, ни в сутки. LATERAL форсирует per-row
# индексный lookup через idx_lss_source_date (listing_source_id, snapshot_date DESC) —
# `ORDER BY s.snapshot_date DESC LIMIT 1` даёт тот же единственный "последний снимок до
# сегодня" на listing_source_id, что и старый DISTINCT ON (PK (listing_source_id,
# snapshot_date) исключает дубликаты snapshot_date на одном источнике — семантика
# идентична), но за O(log n) на строку вместо полного скана таблицы. EXPLAIN на проде:
# cost внутреннего подзапроса упал с ~298 627 до ~4.4 за строку today.
#
# Полностью set-based: один INSERT … SELECT по всем источникам, без Python-цикла (LATERAL
# — это внутренний план Postgres, не Python-итерация).
#
# change_time = date_trunc('day', now()), а НЕ now() (#2674): с now() уникальность
# UNIQUE(listing_source_id, change_time, event_type) работала только ВНУТРИ прогона —
# второй прогон в те же сутки перезаписывал сегодняшний снимок, предикаты срабатывали
# заново с другим временем и давали дубли (2 августа таких прогонов было два).
# Суточная гранулярность честнее для суточного же сравнения и включает заявленную
# идемпотентность: ON CONFLICT DO NOTHING теперь действительно гасит повтор за день.
#
# NULLIF(p.price_rub, 0) в diff_percent обязателен: выражения VALUES вычисляются ДО
# фильтра `WHERE e.fires`, поэтому предикат "p.price_rub <> 0" от деления на ноль уже
# не спасает — без NULLIF первый же источник с нулевой прошлой ценой уронил бы весь
# прогон. Результат при этом тот же: строка с NULL-диффом не проходит e.fires.
#
# Внешний SELECT над data-modifying CTE считает вставленное ПО ТИПАМ (RETURNING отдаёт
# только реально вставленные строки, не съеденные ON CONFLICT), сразу в виде ключей
# счётчиков `<event_type>_events` — писатель получает готовый dict без Python-агрегации.
# Ровно этот счётчик и показал бы четыре нуля из пяти, если бы существовал раньше.
_EVENT_DIFF_SQL = text(
"""
WITH today AS (
SELECT listing_source_id, price_rub, payload_hash
FROM listing_source_snapshots
WHERE snapshot_date = CURRENT_DATE
),
inserted AS (
INSERT INTO listing_source_events (
listing_source_id, change_time, event_type, price_rub, diff_percent
)
SELECT
t.listing_source_id,
date_trunc('day', now()),
e.event_type,
t.price_rub,
e.diff_percent
FROM today t
LEFT JOIN LATERAL (
SELECT s.snapshot_date, s.price_rub, s.payload_hash
FROM listing_source_snapshots s
WHERE s.listing_source_id = t.listing_source_id
AND s.snapshot_date < CURRENT_DATE
ORDER BY s.snapshot_date DESC
LIMIT 1
) p ON true
CROSS JOIN LATERAL (VALUES
(
'first_seen',
NULL::numeric,
p.snapshot_date IS NULL
),
(
'price_change',
round((t.price_rub - p.price_rub)::numeric
/ NULLIF(p.price_rub, 0) * 100, 4),
t.price_rub IS NOT NULL
AND p.price_rub IS NOT NULL
AND p.price_rub <> 0
AND t.price_rub <> p.price_rub
),
(
'edited',
NULL::numeric,
p.payload_hash IS NOT NULL
AND t.payload_hash IS DISTINCT FROM p.payload_hash
AND t.price_rub IS NOT DISTINCT FROM p.price_rub
)
) AS e(event_type, diff_percent, fires)
WHERE e.fires
ON CONFLICT (listing_source_id, change_time, event_type) DO NOTHING
RETURNING event_type
)
SELECT event_type || '_events' AS counter_key, count(*) AS n
FROM inserted
GROUP BY 1
"""
)
def snapshot_listing_sources(
db: Session, run_id: int, params: dict[str, Any] | None = None
) -> dict[str, int]:
"""Записать дневной снимок listing_sources + события изменений.
Sync (вызывается scheduler-триггером в executor, как import_rosreestr_dkp).
Два set-based statement'а в одной транзакции:
1. upsert снимка на (listing_source_id, CURRENT_DATE) только для источников,
изменившихся с последнего снимка (#2993) — last-write-wins.
2. diff сегодняшнего снимка против последнего предыдущего → три события,
выводимые из наших данных (#2674). delisted/relisted схема разрешает, но
они НЕ выводимы при покрытии обхода 10-35% — см. _EVENT_DIFF_SQL.
Params (из default_params jsonb в scrape_schedules, #2607):
budget_sec: float — SET LOCAL statement_timeout на транзакцию (default 900,
clamp [30, 3600]). Единственный способ гарантированно оборвать зависший
statement у не-батчащейся (два statement'а, не Python-цикл) задачи — если
план снова разрегрессирует, прогон честно упадёт в mark_failed вместо того
чтобы висеть часами/сутками (root cause #2607 — см. шапку файла и
_EVENT_DIFF_SQL).
Финализирует scrape_runs (mark_done / mark_failed) и пишет counters.
Returns {"snapshotted": N, "<event_type>_events": M} — по счётчику на каждый из
трёх пишущихся типов, всегда все три ключа (тип, который за прогон не сработал
ни разу, честно показывает 0, а не пропадает из counters).
"""
params = params or {}
budget_sec = _clamp_budget_sec(params.get("budget_sec", DEFAULT_BUDGET_SEC))
counters: dict[str, int] = {
"snapshotted": 0,
"price_change_events": 0,
"edited_events": 0,
"first_seen_events": 0,
}
try:
# statement_timeout НЕ принимает bind-параметр ($1/:name) — синтаксис Postgres SET
# запрещает placeholder на этом месте (проверено вживую на проде: "syntax error at
# or near \"$1\""). budget_sec провалидирован/clamp'нут в _clamp_budget_sec выше
# (источник — scrape_schedules.default_params, не user input) — f-string здесь
# безопасен (единственный практический способ выставить эту GUC динамически).
# SET LOCAL — per-transaction scope, сбрасывается на COMMIT/ROLLBACK, НЕ трогает
# server/role-level statement_timeout (issue #2607 п.2 — отдельное решение).
timeout_ms = int(budget_sec * 1000)
db.execute(text(f"SET LOCAL statement_timeout = {timeout_ms}"))
snap_result = db.execute(
_SNAPSHOT_SQL,
{"freshness_days": FRESHNESS_WINDOW_DAYS, "run_id": run_id},
)
counters["snapshotted"] = snap_result.rowcount or 0
# Statement возвращает уже готовые пары (counter_key, n) по типам событий —
# dict(...) без Python-агрегации, набор ключей задан инициализацией counters
# выше, так что не сработавшие типы остаются нулями, а не исчезают.
event_rows = db.execute(_EVENT_DIFF_SQL).fetchall()
counters.update(dict(event_rows))
db.commit()
runs_mod.mark_done(db, run_id, counters)
logger.info("snapshot_listing_sources run_id=%d done: %s", run_id, counters)
return counters
except Exception as exc:
logger.exception(
"snapshot_listing_sources run_id=%d failed (budget_sec=%.0f)", run_id, budget_sec
)
db.rollback()
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters)
raise