fix(tradein/snapshot): починить зависающий запрос снапшотов и добавить бюджет времени (#2607)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 8s
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 2m37s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 8s
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 2m37s
Root cause: event-diff CTE джойнил "today" (снимок за CURRENT_DATE) с "prior" (DISTINCT ON по всей listing_source_snapshots, ~2.6-2.8M строк) обычным JOIN. Планировщик оценивал today в 1 строку (свежевставленные в той же транзакции строки ANALYZE ещё не видел) → Nested Loop без Materialize пересчитывал DISTINCT ON по всей таблице заново на каждую из ~80-140k реальных строк today (EXPLAIN на проде: cost≈300k на этом шаге) — прогон не укладывался ни в 6h zombie-порог, ни в сутки, каждую ночь минимум с 19 июля. Переписано на JOIN LATERAL (per-row indexed point-lookup через idx_lss_source_date, cost упал до ~4.4/строку). Плюс budget_sec → SET LOCAL statement_timeout как defense-in-depth (по образцу geocode_missing_listings) — задача теперь честно падает в mark_failed вместо того чтобы висеть сутками, если план когда-нибудь разрегрессирует снова. Зомби-детектор (reap_zombies) не тронут — он только помечает scrape_runs.status, не убивает backend (нет pid/application_name в схеме run'а); pg_terminate_backend для этого — отдельный follow-up, не в этом PR.
This commit is contained in:
parent
16237573ed
commit
b586b5ff68
4 changed files with 214 additions and 28 deletions
|
|
@ -94,13 +94,15 @@ async def _job_rosreestr_dkp(
|
||||||
|
|
||||||
|
|
||||||
# ── listing_source_snapshot — sync DB-snapshot в executor ────────────────────
|
# ── listing_source_snapshot — sync DB-snapshot в executor ────────────────────
|
||||||
|
# params прокинуты (#2607) — snapshot_listing_sources теперь читает budget_sec из
|
||||||
|
# default_params (SET LOCAL statement_timeout, см. app/tasks/listing_source_snapshot.py).
|
||||||
async def _job_listing_source_snapshot(
|
async def _job_listing_source_snapshot(
|
||||||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||||||
) -> None:
|
) -> None:
|
||||||
from app.tasks.listing_source_snapshot import snapshot_listing_sources
|
from app.tasks.listing_source_snapshot import snapshot_listing_sources
|
||||||
|
|
||||||
loop = asyncio.get_event_loop()
|
loop = asyncio.get_event_loop()
|
||||||
await loop.run_in_executor(None, snapshot_listing_sources, db, run_id)
|
await loop.run_in_executor(None, snapshot_listing_sources, db, run_id, params)
|
||||||
|
|
||||||
|
|
||||||
# ── asking_to_sold_ratio_refresh — sync re-derive в executor ─────────────────
|
# ── asking_to_sold_ratio_refresh — sync re-derive в executor ─────────────────
|
||||||
|
|
|
||||||
|
|
@ -9,13 +9,34 @@ listing_source_events. Так история per-source цены копится
|
||||||
через product_handlers._job_listing_source_snapshot,
|
через product_handlers._job_listing_source_snapshot,
|
||||||
по образцу import_rosreestr_dkp (sync task в run_in_executor).
|
по образцу import_rosreestr_dkp (sync task в run_in_executor).
|
||||||
|
|
||||||
Вся работа — два set-based SQL statement'а (snapshot upsert + event-diff CTE),
|
Вся работа — два set-based SQL statement'а (snapshot upsert + event-diff), никакого
|
||||||
никакого row-by-row Python: 18 355 строк обслуживаются одним INSERT … SELECT каждый.
|
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
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
@ -27,6 +48,30 @@ logger = logging.getLogger(__name__)
|
||||||
# Окно свежести: источник считается активным, если last_seen_at не старше N дней.
|
# Окно свежести: источник считается активным, если last_seen_at не старше N дней.
|
||||||
FRESHNESS_WINDOW_DAYS = 7
|
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))
|
||||||
|
|
||||||
|
|
||||||
# ── Daily snapshot upsert ─────────────────────────────────────────────────────
|
# ── Daily snapshot upsert ─────────────────────────────────────────────────────
|
||||||
# Снимок на (listing_source_id, CURRENT_DATE). ON CONFLICT → last-write-wins за день
|
# Снимок на (listing_source_id, CURRENT_DATE). ON CONFLICT → last-write-wins за день
|
||||||
# (повторный прогон в те же сутки перезаписывает снимок свежими значениями).
|
# (повторный прогон в те же сутки перезаписывает снимок свежими значениями).
|
||||||
|
|
@ -63,9 +108,25 @@ _SNAPSHOT_SQL = text(
|
||||||
# Для каждого источника сравниваем сегодняшнюю цену (snapshot_date = CURRENT_DATE) с
|
# Для каждого источника сравниваем сегодняшнюю цену (snapshot_date = CURRENT_DATE) с
|
||||||
# самым свежим ПРЕДЫДУЩИМ снимком (snapshot_date < CURRENT_DATE). Если цена изменилась
|
# самым свежим ПРЕДЫДУЩИМ снимком (snapshot_date < CURRENT_DATE). Если цена изменилась
|
||||||
# (обе NOT NULL, old <> 0) — пишем price_change.
|
# (обе NOT NULL, old <> 0) — пишем price_change.
|
||||||
# today — снимок за сегодня (только что записан _SNAPSHOT_SQL).
|
# today — снимок за сегодня (только что записан _SNAPSHOT_SQL, в той же транзакции).
|
||||||
# prior — последний снимок строго ДО сегодня (DISTINCT ON … ORDER BY date DESC).
|
# p — последний снимок строго ДО сегодня, per-row LATERAL point-lookup (#2607).
|
||||||
# Полностью set-based: один INSERT … SELECT по всем источникам, без 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 = now() детерминирует UNIQUE(listing_source_id, change_time, event_type)
|
# change_time = now() детерминирует UNIQUE(listing_source_id, change_time, event_type)
|
||||||
# в пределах прогона → ON CONFLICT DO NOTHING делает писатель идемпотентным.
|
# в пределах прогона → ON CONFLICT DO NOTHING делает писатель идемпотентным.
|
||||||
_EVENT_DIFF_SQL = text(
|
_EVENT_DIFF_SQL = text(
|
||||||
|
|
@ -74,13 +135,6 @@ _EVENT_DIFF_SQL = text(
|
||||||
SELECT listing_source_id, price_rub
|
SELECT listing_source_id, price_rub
|
||||||
FROM listing_source_snapshots
|
FROM listing_source_snapshots
|
||||||
WHERE snapshot_date = CURRENT_DATE
|
WHERE snapshot_date = CURRENT_DATE
|
||||||
),
|
|
||||||
prior AS (
|
|
||||||
SELECT DISTINCT ON (listing_source_id)
|
|
||||||
listing_source_id, price_rub
|
|
||||||
FROM listing_source_snapshots
|
|
||||||
WHERE snapshot_date < CURRENT_DATE
|
|
||||||
ORDER BY listing_source_id, snapshot_date DESC
|
|
||||||
)
|
)
|
||||||
INSERT INTO listing_source_events (
|
INSERT INTO listing_source_events (
|
||||||
listing_source_id, change_time, event_type, price_rub, diff_percent
|
listing_source_id, change_time, event_type, price_rub, diff_percent
|
||||||
|
|
@ -92,7 +146,14 @@ _EVENT_DIFF_SQL = text(
|
||||||
t.price_rub,
|
t.price_rub,
|
||||||
round((t.price_rub - p.price_rub)::numeric / p.price_rub * 100, 4)
|
round((t.price_rub - p.price_rub)::numeric / p.price_rub * 100, 4)
|
||||||
FROM today t
|
FROM today t
|
||||||
JOIN prior p ON p.listing_source_id = t.listing_source_id
|
JOIN LATERAL (
|
||||||
|
SELECT s.price_rub
|
||||||
|
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
|
||||||
WHERE t.price_rub IS NOT NULL
|
WHERE t.price_rub IS NOT NULL
|
||||||
AND p.price_rub IS NOT NULL
|
AND p.price_rub IS NOT NULL
|
||||||
AND p.price_rub <> 0
|
AND p.price_rub <> 0
|
||||||
|
|
@ -102,7 +163,9 @@ _EVENT_DIFF_SQL = text(
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def snapshot_listing_sources(db: Session, run_id: int) -> dict[str, int]:
|
def snapshot_listing_sources(
|
||||||
|
db: Session, run_id: int, params: dict[str, Any] | None = None
|
||||||
|
) -> dict[str, int]:
|
||||||
"""Записать дневной снимок listing_sources + price_change-события.
|
"""Записать дневной снимок listing_sources + price_change-события.
|
||||||
|
|
||||||
Sync (вызывается scheduler-триггером в executor, как import_rosreestr_dkp).
|
Sync (вызывается scheduler-триггером в executor, как import_rosreestr_dkp).
|
||||||
|
|
@ -110,12 +173,32 @@ def snapshot_listing_sources(db: Session, run_id: int) -> dict[str, int]:
|
||||||
1. upsert снимка на (listing_source_id, CURRENT_DATE) — last-write-wins.
|
1. upsert снимка на (listing_source_id, CURRENT_DATE) — last-write-wins.
|
||||||
2. diff сегодняшней цены против последнего предыдущего снимка → price_change-события.
|
2. diff сегодняшней цены против последнего предыдущего снимка → price_change-события.
|
||||||
|
|
||||||
|
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.
|
Финализирует scrape_runs (mark_done / mark_failed) и пишет counters.
|
||||||
|
|
||||||
Returns {"snapshotted": N, "price_change_events": M}.
|
Returns {"snapshotted": N, "price_change_events": M}.
|
||||||
"""
|
"""
|
||||||
|
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}
|
counters: dict[str, int] = {"snapshotted": 0, "price_change_events": 0}
|
||||||
try:
|
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(
|
snap_result = db.execute(
|
||||||
_SNAPSHOT_SQL,
|
_SNAPSHOT_SQL,
|
||||||
{"freshness_days": FRESHNESS_WINDOW_DAYS, "run_id": run_id},
|
{"freshness_days": FRESHNESS_WINDOW_DAYS, "run_id": run_id},
|
||||||
|
|
@ -135,7 +218,9 @@ def snapshot_listing_sources(db: Session, run_id: int) -> dict[str, int]:
|
||||||
)
|
)
|
||||||
return counters
|
return counters
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.exception("snapshot_listing_sources run_id=%d failed", run_id)
|
logger.exception(
|
||||||
|
"snapshot_listing_sources run_id=%d failed (budget_sec=%.0f)", run_id, budget_sec
|
||||||
|
)
|
||||||
db.rollback()
|
db.rollback()
|
||||||
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters)
|
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters)
|
||||||
raise
|
raise
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,38 @@
|
||||||
|
-- 202_listing_source_snapshot_budget_sec.sql
|
||||||
|
-- #2607 — listing_source_snapshot зависал каждую ночь (минимум с 19 июля): scrape_runs
|
||||||
|
-- всегда добирал до 'zombie' ровно за 6h (порог zombie-детектора), но backend в Postgres
|
||||||
|
-- продолжал жечь CPU СУТКАМИ после этого (zombie-детектор в scraper_kit.orchestration.
|
||||||
|
-- scheduler.reap_zombies только помечает строку scrape_runs — не убивает backend), держа
|
||||||
|
-- backend_xmin и блокируя autovacuum на listings/listing_sources.
|
||||||
|
--
|
||||||
|
-- ROOT CAUSE (тот же PR, app/tasks/listing_source_snapshot.py): event-diff CTE джойнил
|
||||||
|
-- "today" (снимок за CURRENT_DATE) с "prior" — DISTINCT ON по ВСЕЙ listing_source_snapshots
|
||||||
|
-- (~2.6-2.8M строк) обычным JOIN. Планировщик оценивал "today" в 1 строку (свежевставленные
|
||||||
|
-- в той же транзакции строки ANALYZE ещё не видел) → выбирал Nested Loop БЕЗ Materialize на
|
||||||
|
-- внутренней стороне → DISTINCT ON пересчитывался заново на КАЖДУЮ из ~80-140k реальных
|
||||||
|
-- строк today. EXPLAIN на проде: cost внутреннего подзапроса ~298 627. Запрос переписан на
|
||||||
|
-- JOIN LATERAL (per-row indexed point-lookup, cost ~4.4/строку) — устраняет корневую причину.
|
||||||
|
--
|
||||||
|
-- ЭТА миграция — ДОПОЛНИТЕЛЬНЫЙ предохранитель (issue #2607 п.4): budget_sec в default_params
|
||||||
|
-- теперь читается snapshot_listing_sources() и выставляется как SET LOCAL statement_timeout
|
||||||
|
-- (per-transaction, НЕ server/role-level — тот отдельный вопрос issue #2607 п.2, требует
|
||||||
|
-- согласования, здесь намеренно не трогается). Если план когда-нибудь снова разрегрессирует,
|
||||||
|
-- прогон честно упадёт в mark_failed вместо того чтобы висеть сутками.
|
||||||
|
--
|
||||||
|
-- 900 сек (15 мин) — по образцу migration 110 (geocode_missing_listings budget_sec=1800),
|
||||||
|
-- с большим запасом над ожидаемым временем выполнения после LATERAL-фикса (секунды) и
|
||||||
|
-- далеко от 6h zombie-порога и от окна 01:00-02:00 UTC (052/079).
|
||||||
|
--
|
||||||
|
-- ЗАВИСИМОСТИ: 079_listing_source_history.sql (создаёт scrape_schedules row, source=
|
||||||
|
-- 'listing_source_snapshot', default_params='{}'::jsonb).
|
||||||
|
-- Idempotent: UPDATE ... || jsonb-merge — безопасно перезапускать (всегда приводит
|
||||||
|
-- default_params.budget_sec к 900 независимо от предыдущего состояния).
|
||||||
|
-- Apply after: 201_purge_dead_mobileproxy_proxies.sql
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
UPDATE scrape_schedules
|
||||||
|
SET default_params = COALESCE(default_params, '{}'::jsonb) || '{"budget_sec": 900}'::jsonb
|
||||||
|
WHERE source = 'listing_source_snapshot';
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
|
|
@ -75,17 +75,29 @@ def test_snapshot_derives_is_active_and_payload_hash() -> None:
|
||||||
# ── Event-diff CTE SQL ────────────────────────────────────────────────────────
|
# ── Event-diff CTE SQL ────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
def test_event_diff_is_set_based_cte_not_python_loop() -> None:
|
def test_event_diff_is_set_based_lateral_not_python_loop() -> None:
|
||||||
"""Event diff is one set-based INSERT … SELECT over a CTE — never a row-by-row loop."""
|
"""Event diff is one set-based INSERT … SELECT with a LATERAL join — no Python loop.
|
||||||
|
|
||||||
|
#2607: prior used to be a `DISTINCT ON (listing_source_id) ... FROM
|
||||||
|
listing_source_snapshots` CTE joined via plain JOIN — the planner's Nested Loop
|
||||||
|
(no Materialize, misestimated `today` row count) re-executed the DISTINCT ON over
|
||||||
|
the whole table once per today-row, hanging for days. Rewritten as `JOIN LATERAL
|
||||||
|
(... ORDER BY snapshot_date DESC LIMIT 1) ON true` — forces a per-row indexed
|
||||||
|
point-lookup via idx_lss_source_date instead of a full-table DISTINCT ON.
|
||||||
|
"""
|
||||||
assert "WITH today AS" in _EVENT_DIFF_SQL
|
assert "WITH today AS" in _EVENT_DIFF_SQL
|
||||||
assert "prior AS" in _EVENT_DIFF_SQL
|
assert "JOIN LATERAL" in _EVENT_DIFF_SQL
|
||||||
assert "DISTINCT ON (listing_source_id)" in _EVENT_DIFF_SQL
|
assert "prior AS" not in _EVENT_DIFF_SQL, "prior CTE removed — replaced by LATERAL (#2607)"
|
||||||
|
assert "DISTINCT ON" not in _EVENT_DIFF_SQL, "DISTINCT ON over full table removed (#2607)"
|
||||||
assert "INSERT INTO listing_source_events" in _EVENT_DIFF_SQL
|
assert "INSERT INTO listing_source_events" in _EVENT_DIFF_SQL
|
||||||
# Prior = most-recent snapshot strictly before today.
|
# LATERAL subquery: most-recent snapshot strictly before today, per listing_source_id.
|
||||||
assert "snapshot_date < CURRENT_DATE" in _EVENT_DIFF_SQL
|
assert "s.listing_source_id = t.listing_source_id" in _EVENT_DIFF_SQL
|
||||||
|
assert "s.snapshot_date < CURRENT_DATE" in _EVENT_DIFF_SQL
|
||||||
assert "snapshot_date = CURRENT_DATE" in _EVENT_DIFF_SQL
|
assert "snapshot_date = CURRENT_DATE" in _EVENT_DIFF_SQL
|
||||||
assert "ORDER BY listing_source_id, snapshot_date DESC" in _EVENT_DIFF_SQL
|
assert "ORDER BY s.snapshot_date DESC" in _EVENT_DIFF_SQL
|
||||||
# No Python iteration over rows in the writer body (set-based only).
|
assert "LIMIT 1" in _EVENT_DIFF_SQL
|
||||||
|
# No Python iteration over rows in the writer body (set-based only — LATERAL is a
|
||||||
|
# Postgres execution-plan construct, not a Python loop).
|
||||||
body = _WRITER_SRC.split('"""', 2)[-1]
|
body = _WRITER_SRC.split('"""', 2)[-1]
|
||||||
assert "for " not in body, "writer must be set-based — no Python row loop"
|
assert "for " not in body, "writer must be set-based — no Python row loop"
|
||||||
|
|
||||||
|
|
@ -222,20 +234,69 @@ def test_counter_logic_with_fake_db(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
)
|
)
|
||||||
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||||
|
|
||||||
db = _FakeDB(rowcounts=[18355, 42]) # snapshot rowcount, then event rowcount
|
# rowcounts: SET LOCAL statement_timeout (ignored), snapshot upsert, event-diff insert.
|
||||||
|
db = _FakeDB(rowcounts=[0, 18355, 42])
|
||||||
out = snap_mod.snapshot_listing_sources(db, run_id=99) # type: ignore[arg-type]
|
out = snap_mod.snapshot_listing_sources(db, run_id=99) # type: ignore[arg-type]
|
||||||
|
|
||||||
assert out == {"snapshotted": 18355, "price_change_events": 42}
|
assert out == {"snapshotted": 18355, "price_change_events": 42}
|
||||||
assert db.committed is True
|
assert db.committed is True
|
||||||
assert len(db.executed) == 2
|
assert len(db.executed) == 3
|
||||||
# run_id threaded into the snapshot statement's bind params.
|
# First statement sets the per-transaction wall-clock budget (#2607).
|
||||||
_stmt, params = db.executed[0]
|
stmt0, _params0 = db.executed[0]
|
||||||
|
assert "SET LOCAL statement_timeout" in str(stmt0)
|
||||||
|
# run_id threaded into the snapshot statement's bind params (now executed[1]).
|
||||||
|
_stmt, params = db.executed[1]
|
||||||
assert params is not None and params["run_id"] == 99
|
assert params is not None and params["run_id"] == 99
|
||||||
# Run finalised via mark_done with the same counters.
|
# Run finalised via mark_done with the same counters.
|
||||||
assert marked["run_id"] == 99
|
assert marked["run_id"] == 99
|
||||||
assert marked["counters"] == {"snapshotted": 18355, "price_change_events": 42}
|
assert marked["counters"] == {"snapshotted": 18355, "price_change_events": 42}
|
||||||
|
|
||||||
|
|
||||||
|
# ── budget_sec / statement_timeout (#2607) ─────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_clamp_budget_sec_defaults_and_bounds() -> None:
|
||||||
|
assert snap_mod._clamp_budget_sec(snap_mod.DEFAULT_BUDGET_SEC) == snap_mod.DEFAULT_BUDGET_SEC
|
||||||
|
# Below floor / garbage / zero (the historical bug: 0 == "no timeout") clamp to the floor.
|
||||||
|
assert snap_mod._clamp_budget_sec(0) == snap_mod._MIN_BUDGET_SEC
|
||||||
|
assert snap_mod._clamp_budget_sec(-5) == snap_mod._MIN_BUDGET_SEC
|
||||||
|
assert snap_mod._clamp_budget_sec(None) == snap_mod.DEFAULT_BUDGET_SEC
|
||||||
|
assert snap_mod._clamp_budget_sec("garbage") == snap_mod.DEFAULT_BUDGET_SEC
|
||||||
|
# Above ceiling clamps down — never lets a fat-fingered value re-create "hangs forever".
|
||||||
|
assert snap_mod._clamp_budget_sec(999_999) == snap_mod._MAX_BUDGET_SEC
|
||||||
|
# Sane custom value passes through unclamped.
|
||||||
|
assert snap_mod._clamp_budget_sec(120) == 120.0
|
||||||
|
|
||||||
|
|
||||||
|
def test_snapshot_listing_sources_sets_statement_timeout_from_params(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""budget_sec from default_params is applied via SET LOCAL statement_timeout (ms)."""
|
||||||
|
monkeypatch.setattr(snap_mod.runs_mod, "mark_done", lambda *a, **k: None)
|
||||||
|
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||||
|
|
||||||
|
db = _FakeDB(rowcounts=[0, 10, 1])
|
||||||
|
snap_mod.snapshot_listing_sources(db, run_id=1, params={"budget_sec": 120}) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
stmt0, _params0 = db.executed[0]
|
||||||
|
assert "SET LOCAL statement_timeout = 120000" in str(stmt0)
|
||||||
|
|
||||||
|
|
||||||
|
def test_snapshot_listing_sources_defaults_budget_sec_when_params_missing(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""No params / no budget_sec key → DEFAULT_BUDGET_SEC applied (never unlimited/0)."""
|
||||||
|
monkeypatch.setattr(snap_mod.runs_mod, "mark_done", lambda *a, **k: None)
|
||||||
|
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||||
|
|
||||||
|
db = _FakeDB(rowcounts=[0, 10, 1])
|
||||||
|
snap_mod.snapshot_listing_sources(db, run_id=1) # type: ignore[arg-type]
|
||||||
|
|
||||||
|
stmt0, _params0 = db.executed[0]
|
||||||
|
expected_ms = int(snap_mod.DEFAULT_BUDGET_SEC * 1000)
|
||||||
|
assert f"SET LOCAL statement_timeout = {expected_ms}" in str(stmt0)
|
||||||
|
|
||||||
|
|
||||||
def test_counter_logic_failure_path_marks_failed(monkeypatch: pytest.MonkeyPatch) -> None:
|
def test_counter_logic_failure_path_marks_failed(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
"""On execute error: rollback + mark_failed + re-raise (no silent swallow)."""
|
"""On execute error: rollback + mark_failed + re-raise (no silent swallow)."""
|
||||||
failed: dict[str, Any] = {}
|
failed: dict[str, Any] = {}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue