fix(tradein/deactivate-stale): пол переобхода — LATERAL вместо equality-join по дате
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
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 4m41s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
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 4m41s
Правка меняет две вещи разом, обе намеренно. 1. Равенство по дате -> «последняя строка не позже якоря». listing_source_snapshots переходит на модель «строка на изменение». В ней equality-join теряет 95.6% пар (на 2026-08-20 изменениями являются 4458 строк из 101795): выборка для percentile_disc схлопывается до n=3-4, квантиль вырождается в максимум из трёх чисел, TTL обваливается — domklik 28->14 (под снятие сразу 128 активных строк), yandex 44->30, avito 13->10. Дыры в суточной истории ломали equality-join и до перехода: на проде 03-14.06 (12 суток подряд), 03-04.07, 12.07, 26.07, 30-31.07, 01.08 — там n_pairs=0, floor_days=NULL и TTL молча оставался как задан. 2. Глобальный max(snapshot_date) -> максимум внутри источника. Старый подзапрос брал максимум по ВСЕЙ таблице, не скоупленный по listing_source_id: источник, чья история короче общей, выпадал из выборки целиком. LATERAL ищет предшественника по строке. Замер на живом проде 2026-08-23 (health_window_days=3): обе формы дают побитово одинаковый результат на всех четырёх источниках — avito n=765 пол=49.7, cian n=1226 пол=81.6, domklik n=565 пол=35.7, yandex n=3856 пол=87.3. Сегодня суточная джоба пишет строку для каждого источника каждый день, поэтому глобальный максимум совпадает с максимумом каждого источника, и пункт 2 — no-op на текущих данных. Расхождение проявится только на дырах и после перехода на change-only. Плюс floor_n_pairs в counters — наблюдательность для будущего гейта деградации пола, сейчас ничего не блокирует. FROM..WHERE вынесен в _revisit_floor_from_where_sql, чтобы count(*) и percentile_disc гарантированно шли по одному срезу.
This commit is contained in:
parent
f5e2f5f5f3
commit
fb38d657ad
3 changed files with 548 additions and 17 deletions
|
|
@ -276,6 +276,43 @@ _REVISIT_FLOOR_NULL_SEGMENT_FILTER = "\n AND l.listing_segment IS NULL"
|
||||||
CAP_MULT = 2
|
CAP_MULT = 2
|
||||||
|
|
||||||
|
|
||||||
|
def _revisit_floor_from_where_sql(staleness_column: str, segment_filter: str) -> str:
|
||||||
|
"""FROM..WHERE, общий для _build_revisit_floor_sql и _build_revisit_floor_pairs_count_sql.
|
||||||
|
|
||||||
|
Вынесено в отдельную функцию, а не продублировано, ЦЕЛЕНАПРАВЛЕННО: пункт 3 PR-A
|
||||||
|
(#2659 продолжение) требует, чтобы n_pairs в counters считался «по тому же срезу,
|
||||||
|
что даёт выборку для percentile_disc» -- дословно, а не приблизительно. Общий текст
|
||||||
|
гарантирует это по построению; раздельные копии SELECT count(*) и SELECT
|
||||||
|
percentile_disc(...) рано или поздно разошлись бы независимой правкой одной из
|
||||||
|
двух.
|
||||||
|
|
||||||
|
JOIN LATERAL — предшественник ищется ПО СТРОКЕ (per listing_source_id), не по
|
||||||
|
глобальной дате: `ORDER BY s.snapshot_date DESC LIMIT 1` берёт последний снимок
|
||||||
|
ЭТОГО listing_source_id не позже якоря (CURRENT_DATE - health_window_days),
|
||||||
|
независимо от того, писалась ли строка именно на дату якоря. Подробно — почему
|
||||||
|
именно так и что было раньше — см. докстринг _build_revisit_floor_sql.
|
||||||
|
"""
|
||||||
|
return f"""
|
||||||
|
FROM listings l
|
||||||
|
JOIN listing_sources ls
|
||||||
|
ON ls.listing_id = l.id
|
||||||
|
AND ls.ext_source = l.source
|
||||||
|
JOIN LATERAL (
|
||||||
|
SELECT s.last_seen_at
|
||||||
|
FROM listing_source_snapshots s
|
||||||
|
WHERE s.listing_source_id = ls.id
|
||||||
|
AND s.snapshot_date
|
||||||
|
<= CURRENT_DATE - CAST(:health_window_days AS integer)
|
||||||
|
ORDER BY s.snapshot_date DESC
|
||||||
|
LIMIT 1
|
||||||
|
) prev ON true
|
||||||
|
WHERE l.source = :listing_source
|
||||||
|
AND l.{staleness_column}
|
||||||
|
> NOW() - CAST(:health_window_days || ' days' AS interval)
|
||||||
|
AND l.{staleness_column} > prev.last_seen_at{segment_filter}
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
def _build_revisit_floor_sql(
|
def _build_revisit_floor_sql(
|
||||||
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
||||||
) -> Any:
|
) -> Any:
|
||||||
|
|
@ -286,6 +323,32 @@ def _build_revisit_floor_sql(
|
||||||
Только строки, у которых свежесть реально сдвинулась, — то есть выжившие,
|
Только строки, у которых свежесть реально сдвинулась, — то есть выжившие,
|
||||||
а не «мы к ним не приходили».
|
а не «мы к ним не приходили».
|
||||||
|
|
||||||
|
«ПРЕДЫДУЩЕЕ НАБЛЮДЕНИЕ» -- ЧТО ЭТО ТОЧНО (PR-A, #2659 продолжение). Это ПОСЛЕДНЯЯ
|
||||||
|
строка listing_source_snapshots для данного listing_source_id, чей snapshot_date
|
||||||
|
не позже якоря (CURRENT_DATE - health_window_days) -- `ORDER BY snapshot_date DESC
|
||||||
|
LIMIT 1` в LATERAL-подзапросе _revisit_floor_from_where_sql. РАВЕНСТВО ПО ДАТЕ
|
||||||
|
(`prev.snapshot_date = <якорь>`) ЗДЕСЬ ЗАПРЕЩЕНО, и вот почему:
|
||||||
|
|
||||||
|
1) listing_source_snapshots переходит на модель «строка на изменение» (пишется
|
||||||
|
только когда значение отличается от предыдущего снимка, не ежедневно). В этой
|
||||||
|
модели equality-join по дате теряет подавляющее большинство пар: на 2026-08-20
|
||||||
|
«изменениями» являются 4 458 строк из 101 795 -- 95.6% пар для percentile_disc
|
||||||
|
пропадают, выборка схлопывается до n=3-4, а percentile_disc(0.99) на такой
|
||||||
|
выборке вырождается в максимум из трёх чисел, а не в реальный квантиль хвоста.
|
||||||
|
2) Дыры в ЕЖЕДНЕВНОЙ истории уже ломали equality-join и ДО перехода на change-only
|
||||||
|
модель -- на проде отмечены разрывы 03-14.06 (12 суток подряд), 03-04.07,
|
||||||
|
12.07, 26.07, 30-31.07, 01.08. В эти дни equality-join давал n_pairs=0,
|
||||||
|
floor_days=NULL, и пол молча не считался -- TTL оставался как задан, хотя
|
||||||
|
история для «ближайшего более раннего» снимка (см. комментарий про потолок
|
||||||
|
выше в модульном докстринге) в базе была, просто не РОВНО на эту дату.
|
||||||
|
|
||||||
|
Семантика при этом не меняется: между двумя изменениями last_seen_at по
|
||||||
|
определению постоянен (иначе строка была бы новым изменением), поэтому
|
||||||
|
«последний снимок не позже якоря» и «снимок ровно на дату якоря» при СПЛОШНОЙ
|
||||||
|
ежедневной истории дают одно и то же число -- разница проявляется только там,
|
||||||
|
где equality-join был неверен и раньше (гэпы) либо станет неверен при переходе
|
||||||
|
на change-only (почти повсеместно).
|
||||||
|
|
||||||
null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments)
|
null_segment_only=True переопределяет with_segments -- IS NULL вместо ANY(:segments)
|
||||||
(ANY никогда не матчит NULL). На практике для null_segment_only-джобы этот запрос
|
(ANY никогда не матчит NULL). На практике для null_segment_only-джобы этот запрос
|
||||||
не строится вовсе (revisit_floor_quantile=0 -- см. модульный докстринг), но вариант
|
не строится вовсе (revisit_floor_quantile=0 -- см. модульный докстринг), но вариант
|
||||||
|
|
@ -307,22 +370,35 @@ def _build_revisit_floor_sql(
|
||||||
ORDER BY EXTRACT(epoch FROM (l.{staleness_column} - prev.last_seen_at))
|
ORDER BY EXTRACT(epoch FROM (l.{staleness_column} - prev.last_seen_at))
|
||||||
/ 86400.0
|
/ 86400.0
|
||||||
)
|
)
|
||||||
FROM listings l
|
{_revisit_floor_from_where_sql(staleness_column, segment_filter)}
|
||||||
JOIN listing_sources ls
|
"""
|
||||||
ON ls.listing_id = l.id
|
)
|
||||||
AND ls.ext_source = l.source
|
|
||||||
JOIN listing_source_snapshots prev
|
|
||||||
ON prev.listing_source_id = ls.id
|
def _build_revisit_floor_pairs_count_sql(
|
||||||
AND prev.snapshot_date = (
|
staleness_column: str, *, with_segments: bool, null_segment_only: bool = False
|
||||||
SELECT max(snapshot_date)
|
) -> Any:
|
||||||
FROM listing_source_snapshots
|
"""count(*) пар, из которых percentile_disc в _build_revisit_floor_sql берёт квантиль.
|
||||||
WHERE snapshot_date
|
|
||||||
<= CURRENT_DATE - CAST(:health_window_days AS integer)
|
Наблюдательность (PR-A, #2659 продолжение) -- НИЧЕГО не блокирует сейчас (гейт по
|
||||||
)
|
деградации выборки, если он когда-нибудь понадобится, -- отдельная задача, PR-B).
|
||||||
WHERE l.source = :listing_source
|
Пишется в counters["floor_n_pairs"] ДО того, как схлопнувшаяся выборка станет
|
||||||
AND l.{staleness_column}
|
видна только по повторению прод-инцидента, ради которого весь пол заведён (см.
|
||||||
> NOW() - CAST(:health_window_days || ' days' AS interval)
|
ЗАМЕР НА ПРОДЕ 2026-08-09 в модульном докстринге).
|
||||||
AND l.{staleness_column} > prev.last_seen_at{segment_filter}
|
|
||||||
|
Тот же срез, что и percentile_disc -- ОБЩАЯ функция _revisit_floor_from_where_sql,
|
||||||
|
не копия WHERE, см. её докстринг про то, почему это важно.
|
||||||
|
"""
|
||||||
|
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 count(*)
|
||||||
|
{_revisit_floor_from_where_sql(staleness_column, segment_filter)}
|
||||||
"""
|
"""
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -499,7 +575,10 @@ def deactivate_stale_listings(
|
||||||
|
|
||||||
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
||||||
Если гейт не пропустил прогон: {"deactivated": 0, "confirmations": N,
|
Если гейт не пропустил прогон: {"deactivated": 0, "confirmations": N,
|
||||||
"skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода поднял
|
"skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Если пол переобхода включён
|
||||||
|
(revisit_floor_quantile > 0): дополнительно {"floor_n_pairs": N} -- размер
|
||||||
|
выборки, из которой percentile_disc посчитал квантиль (наблюдательность, PR-A
|
||||||
|
#2659 продолжение; сейчас ничего не гейтит). Если пол при этом реально поднял
|
||||||
TTL: дополнительно {"revisit_floor_days": N, "ttl_days_effective": N}. Если пол
|
TTL: дополнительно {"revisit_floor_days": N, "ttl_days_effective": N}. Если пол
|
||||||
упёрся в потолок cap_mult: дополнительно {"ttl_floor_capped": 1,
|
упёрся в потолок cap_mult: дополнительно {"ttl_floor_capped": 1,
|
||||||
"ttl_days_floor_raw": N} -- N это то, во что пол поднял бы TTL БЕЗ потолка.
|
"ttl_days_floor_raw": N} -- N это то, во что пол поднял бы TTL БЕЗ потолка.
|
||||||
|
|
@ -632,6 +711,21 @@ def deactivate_stale_listings(
|
||||||
),
|
),
|
||||||
floor_params,
|
floor_params,
|
||||||
).scalar()
|
).scalar()
|
||||||
|
# Наблюдательность (PR-A, #2659 продолжение) -- та же выборка, что и
|
||||||
|
# percentile_disc выше (общая _revisit_floor_from_where_sql). Пока ничего
|
||||||
|
# не гейтит (это PR-B), но n_pairs обязан попасть в counters ДО того, как
|
||||||
|
# схлопнувшуюся выборку станет видно только по повторению прод-инцидента.
|
||||||
|
# Считается ВСЕГДА при включённом полу, в том числе когда floor_days
|
||||||
|
# получится NULL -- тогда n_pairs=0 и объясняет, почему пол не посчитался.
|
||||||
|
floor_n_pairs = db.execute(
|
||||||
|
_build_revisit_floor_pairs_count_sql(
|
||||||
|
staleness_column,
|
||||||
|
with_segments=segments is not None,
|
||||||
|
null_segment_only=null_segment_only,
|
||||||
|
),
|
||||||
|
floor_params,
|
||||||
|
).scalar()
|
||||||
|
counters["floor_n_pairs"] = int(floor_n_pairs or 0)
|
||||||
# NULL = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
# NULL = истории снимков за окно нет вовсе (свежая БД, дыра в снимках).
|
||||||
# Тогда пола нет и TTL остаётся как задан: выдумывать пол не из чего.
|
# Тогда пола нет и TTL остаётся как задан: выдумывать пол не из чего.
|
||||||
if floor_days is not None:
|
if floor_days is not None:
|
||||||
|
|
|
||||||
|
|
@ -109,3 +109,12 @@ tests/test_2992_upsert_unchanged_gate.py::test_listing_sources_unchanged_rescrap
|
||||||
# настоящую БД (в CI она есть, #2745), локально без TEST_DATABASE_URL пропускается. Строки
|
# настоящую БД (в CI она есть, #2745), локально без TEST_DATABASE_URL пропускается. Строки
|
||||||
# t3036-* тест удаляет в finally.
|
# t3036-* тест удаляет в finally.
|
||||||
tests/test_3036_detail_house_params_to_houses.py::test_live_fill_only_then_keep_then_unlinked_untouched
|
tests/test_3036_detail_house_params_to_houses.py::test_live_fill_only_then_keep_then_unlinked_untouched
|
||||||
|
|
||||||
|
# PR-A LATERAL-поиск предшественника в поле переобхода (#2659 продолжение) — тот же
|
||||||
|
# `_live_session()`. Три live-теста вставляют синтетическую историю снимков (в т.ч.
|
||||||
|
# с дырами) в ОДНОЙ транзакции и проверяют результат JOIN LATERAL/equality-join
|
||||||
|
# прямым SELECT'ом в той же транзакции — cleanup через rollback (без commit), явный
|
||||||
|
# DELETE не нужен. Без БД — skip; в CI Trade-In идут на postgres-сервисе (#2745).
|
||||||
|
tests/test_revisit_floor_lateral_lookup.py::test_lateral_ignores_gaps_and_matches_dense_daily_history
|
||||||
|
tests/test_revisit_floor_lateral_lookup.py::test_equality_join_loses_gapped_pairs_and_gives_a_smaller_floor
|
||||||
|
tests/test_revisit_floor_lateral_lookup.py::test_missing_history_before_anchor_still_yields_no_floor_via_lateral
|
||||||
|
|
|
||||||
428
tradein-mvp/backend/tests/test_revisit_floor_lateral_lookup.py
Normal file
428
tradein-mvp/backend/tests/test_revisit_floor_lateral_lookup.py
Normal file
|
|
@ -0,0 +1,428 @@
|
||||||
|
"""LATERAL-поиск предшественника в поле переобхода (#2659 продолжение, PR-A).
|
||||||
|
|
||||||
|
`_build_revisit_floor_sql` джойнил историю переобхода ПО РАВЕНСТВУ ДАТЫ:
|
||||||
|
`prev.snapshot_date = (SELECT max(snapshot_date) FROM listing_source_snapshots
|
||||||
|
WHERE snapshot_date <= CURRENT_DATE - health_window_days)`. Это ломается ДВАЖДЫ:
|
||||||
|
|
||||||
|
1. listing_source_snapshots переходит на модель «строка на изменение» (пишется
|
||||||
|
только когда значение отличается от предыдущего снимка) -- в ней у подавляющего
|
||||||
|
большинства listing_source_id на конкретную календарную дату строки просто нет.
|
||||||
|
На 2026-08-20 «изменениями» являются 4 458 строк из 101 795 (95.6% пар теряются).
|
||||||
|
2. Дыры в истории ломали equality-join и раньше, до перехода на change-only (на
|
||||||
|
проде: 03-14.06 -- 12 суток подряд, 03-04.07, 12.07, 26.07, 30-31.07, 01.08).
|
||||||
|
|
||||||
|
Правка меняет equality-join на `JOIN LATERAL (... ORDER BY snapshot_date DESC
|
||||||
|
LIMIT 1) prev ON true` -- предшественник ищется ПО СТРОКЕ (per listing_source_id),
|
||||||
|
не по единой глобальной дате, и гэпы в истории для него прозрачны. Семантика не
|
||||||
|
меняется: между изменениями last_seen_at по определению постоянен, поэтому
|
||||||
|
«последний снимок не позже якоря» и «снимок ровно на дату якоря» при СПЛОШНОЙ
|
||||||
|
ежедневной истории дают одно и то же число -- разница проявляется только там, где
|
||||||
|
equality-join терял пары.
|
||||||
|
|
||||||
|
Тесты ниже:
|
||||||
|
- test_lateral_ignores_gaps_and_matches_dense_daily_history -- LATERAL на дырявой
|
||||||
|
(change-only) истории даёт ТОТ ЖЕ floor_days, что дала бы сплошная суточная
|
||||||
|
история для того же слушателя.
|
||||||
|
- test_equality_join_loses_gapped_pairs_and_gives_a_smaller_floor -- на ОДНИХ И ТЕХ
|
||||||
|
ЖЕ данных equality-join (воспроизведён буквально -- см. _OLD_EQUALITY_JOIN_FLOOR_SQL
|
||||||
|
ниже, это ровно то условие, что было в _build_revisit_floor_sql до этого PR)
|
||||||
|
теряет строки без снимка ровно на дату якоря и даёт МЕНЬШИЙ пол.
|
||||||
|
- Две pure-SQL проверки без БД (форма LATERAL-запроса, "count(*)" из того же среза).
|
||||||
|
|
||||||
|
Живая БД: см. `_live_session()` -- self-skip без реального Postgres, как соседние
|
||||||
|
live-тесты в этом каталоге (test_2992_upsert_unchanged_gate.py и т.д.). В CI Trade-In
|
||||||
|
есть postgres-сервис (ci-tradein.yml), локально без БД эти проверки skip'аются, а
|
||||||
|
чисто-SQL тесты внизу файла бегут всегда.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
import re
|
||||||
|
import uuid
|
||||||
|
from datetime import UTC, date, datetime, timedelta
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from sqlalchemy import text
|
||||||
|
|
||||||
|
from app.tasks import deactivate_stale_avito as task_mod
|
||||||
|
|
||||||
|
_HEALTH_WINDOW_DAYS = 5
|
||||||
|
|
||||||
|
|
||||||
|
def _live_session() -> Any | None:
|
||||||
|
"""Тот же контракт, что у соседних live-тестов (test_2992_upsert_unchanged_gate.py)."""
|
||||||
|
try:
|
||||||
|
from sqlalchemy import create_engine
|
||||||
|
from sqlalchemy.orm import sessionmaker
|
||||||
|
|
||||||
|
dsn = os.environ.get("TEST_DATABASE_URL") or os.environ.get("DATABASE_URL", "")
|
||||||
|
if not dsn or "localhost:5432/test" in dsn:
|
||||||
|
return None
|
||||||
|
engine = create_engine(dsn, future=True)
|
||||||
|
conn = engine.connect()
|
||||||
|
conn.execute(text("SELECT 1"))
|
||||||
|
conn.close()
|
||||||
|
return sessionmaker(bind=engine, future=True)()
|
||||||
|
except Exception:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
# ── Воспроизведение СТАРОГО equality-join (буквально то, что было в
|
||||||
|
# _build_revisit_floor_sql до этой правки) -- код удалён из app/, поэтому для
|
||||||
|
# сравнения "до/после" запрос встроен сюда как литерал. Единственное упрощение:
|
||||||
|
# вместо подзапроса `(SELECT max(snapshot_date) FROM listing_source_snapshots
|
||||||
|
# WHERE snapshot_date <= CURRENT_DATE - :N)` (глобальный по ВСЕЙ таблице, не
|
||||||
|
# скоупленный per-listing_source_id -- и в этом половина исходного бага) якорная
|
||||||
|
# дата передаётся явно предвычисленным `CURRENT_DATE - :health_window_days`.
|
||||||
|
# Само по себе равенство по дате -- ровно та семантика, что менялась в этом PR;
|
||||||
|
# зависимость от глобального max() по ВСЕЙ таблице сделала бы этот тест хрупким
|
||||||
|
# к данным других тестов в той же живой БД, никак не проверяя суть правки.
|
||||||
|
_OLD_EQUALITY_JOIN_FLOOR_SQL = text(
|
||||||
|
"""
|
||||||
|
SELECT percentile_disc(CAST(:revisit_quantile AS double precision))
|
||||||
|
WITHIN GROUP (
|
||||||
|
ORDER BY EXTRACT(epoch FROM (l.last_seen_at - prev.last_seen_at)) / 86400.0
|
||||||
|
)
|
||||||
|
FROM listings l
|
||||||
|
JOIN listing_sources ls
|
||||||
|
ON ls.listing_id = l.id
|
||||||
|
AND ls.ext_source = l.source
|
||||||
|
JOIN listing_source_snapshots prev
|
||||||
|
ON prev.listing_source_id = ls.id
|
||||||
|
AND prev.snapshot_date = CURRENT_DATE - CAST(:health_window_days AS integer)
|
||||||
|
WHERE l.source = :listing_source
|
||||||
|
AND l.last_seen_at > NOW() - CAST(:health_window_days || ' days' AS interval)
|
||||||
|
AND l.last_seen_at > prev.last_seen_at
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
|
||||||
|
_OLD_EQUALITY_JOIN_COUNT_SQL = text(
|
||||||
|
"""
|
||||||
|
SELECT count(*)
|
||||||
|
FROM listings l
|
||||||
|
JOIN listing_sources ls
|
||||||
|
ON ls.listing_id = l.id
|
||||||
|
AND ls.ext_source = l.source
|
||||||
|
JOIN listing_source_snapshots prev
|
||||||
|
ON prev.listing_source_id = ls.id
|
||||||
|
AND prev.snapshot_date = CURRENT_DATE - CAST(:health_window_days AS integer)
|
||||||
|
WHERE l.source = :listing_source
|
||||||
|
AND l.last_seen_at > NOW() - CAST(:health_window_days || ' days' AS interval)
|
||||||
|
AND l.last_seen_at > prev.last_seen_at
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _insert_pair(
|
||||||
|
db: Any,
|
||||||
|
*,
|
||||||
|
source: str,
|
||||||
|
ext_id: str,
|
||||||
|
listing_last_seen_at: datetime,
|
||||||
|
snapshots: list[tuple[date, datetime]],
|
||||||
|
) -> int:
|
||||||
|
"""Вставляет listings + listing_sources + N снимков listing_source_snapshots.
|
||||||
|
|
||||||
|
БЕЗ commit -- вся синтетика живёт в ОДНОЙ транзакции теста и видна собственным
|
||||||
|
же SELECT'ам (same-session read-your-writes), очистка -- rollback в конце теста,
|
||||||
|
отдельного DELETE не нужно.
|
||||||
|
"""
|
||||||
|
listing_id = db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
INSERT INTO listings
|
||||||
|
(source, source_url, source_id, dedup_hash, price_rub,
|
||||||
|
is_active, scraped_at, last_seen_at)
|
||||||
|
VALUES
|
||||||
|
(:source, :url, :ext_id, :dedup_hash, 5000000,
|
||||||
|
true, NOW(), :last_seen_at)
|
||||||
|
RETURNING id
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{
|
||||||
|
"source": source,
|
||||||
|
"url": f"https://example.test/pr2659-lateral/{ext_id}",
|
||||||
|
"ext_id": ext_id,
|
||||||
|
"dedup_hash": f"pr2659-lateral-{ext_id}",
|
||||||
|
"last_seen_at": listing_last_seen_at,
|
||||||
|
},
|
||||||
|
).scalar_one()
|
||||||
|
|
||||||
|
listing_source_id = db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
INSERT INTO listing_sources
|
||||||
|
(listing_id, ext_source, ext_id, confidence, matched_method)
|
||||||
|
VALUES
|
||||||
|
(:listing_id, :source, :ext_id, 1.0, 'test')
|
||||||
|
RETURNING id
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"listing_id": listing_id, "source": source, "ext_id": ext_id},
|
||||||
|
).scalar_one()
|
||||||
|
|
||||||
|
for snapshot_date, snapshot_last_seen_at in snapshots:
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
INSERT INTO listing_source_snapshots
|
||||||
|
(listing_source_id, snapshot_date, is_active, last_seen_at)
|
||||||
|
VALUES
|
||||||
|
(:lsid, :snapshot_date, true, :last_seen_at)
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{
|
||||||
|
"lsid": listing_source_id,
|
||||||
|
"snapshot_date": snapshot_date,
|
||||||
|
"last_seen_at": snapshot_last_seen_at,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
return int(listing_source_id)
|
||||||
|
|
||||||
|
|
||||||
|
def _floor(db: Any, *, source: str, quantile: float = 1.0) -> float | None:
|
||||||
|
result = db.execute(
|
||||||
|
task_mod._build_revisit_floor_sql("last_seen_at", with_segments=False),
|
||||||
|
{
|
||||||
|
"listing_source": source,
|
||||||
|
"health_window_days": _HEALTH_WINDOW_DAYS,
|
||||||
|
"revisit_quantile": quantile,
|
||||||
|
},
|
||||||
|
).scalar()
|
||||||
|
return float(result) if result is not None else None
|
||||||
|
|
||||||
|
|
||||||
|
def _n_pairs(db: Any, *, source: str) -> int:
|
||||||
|
result = db.execute(
|
||||||
|
task_mod._build_revisit_floor_pairs_count_sql("last_seen_at", with_segments=False),
|
||||||
|
{"listing_source": source, "health_window_days": _HEALTH_WINDOW_DAYS},
|
||||||
|
).scalar()
|
||||||
|
return int(result or 0)
|
||||||
|
|
||||||
|
|
||||||
|
def _old_floor(db: Any, *, source: str, quantile: float = 1.0) -> float | None:
|
||||||
|
result = db.execute(
|
||||||
|
_OLD_EQUALITY_JOIN_FLOOR_SQL,
|
||||||
|
{
|
||||||
|
"listing_source": source,
|
||||||
|
"health_window_days": _HEALTH_WINDOW_DAYS,
|
||||||
|
"revisit_quantile": quantile,
|
||||||
|
},
|
||||||
|
).scalar()
|
||||||
|
return float(result) if result is not None else None
|
||||||
|
|
||||||
|
|
||||||
|
def _old_n_pairs(db: Any, *, source: str) -> int:
|
||||||
|
result = db.execute(
|
||||||
|
_OLD_EQUALITY_JOIN_COUNT_SQL,
|
||||||
|
{"listing_source": source, "health_window_days": _HEALTH_WINDOW_DAYS},
|
||||||
|
).scalar()
|
||||||
|
return int(result or 0)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Живые тесты (реальный Postgres, транзакция + rollback) ────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB")
|
||||||
|
def test_lateral_ignores_gaps_and_matches_dense_daily_history() -> None:
|
||||||
|
"""Дырявая (change-only) история даёт ТОТ ЖЕ floor, что дала бы сплошная суточная.
|
||||||
|
|
||||||
|
Между изменениями last_seen_at по определению постоянен -- поэтому у одного и
|
||||||
|
того же listing_source_id снимок «ровно на дату якоря» и «последний снимок не
|
||||||
|
позже якоря» несут ОДНО И ТО ЖЕ last_seen_at, если между ними ничего не менялось.
|
||||||
|
Сначала считаем пол на СПЛОШНОЙ суточной истории (якорная дата присутствует),
|
||||||
|
затем удаляем все снимки, кроме самого раннего (change-only модель — сутки без
|
||||||
|
изменений не пишутся), и проверяем, что LATERAL даёт то же число.
|
||||||
|
"""
|
||||||
|
db = _live_session()
|
||||||
|
source = f"zzzlat_dense_{uuid.uuid4().hex[:8]}"
|
||||||
|
now = datetime.now(UTC)
|
||||||
|
anchor = date.today() - timedelta(days=_HEALTH_WINDOW_DAYS)
|
||||||
|
# last_seen_at не меняется все эти сутки -- то самое "между изменениями постоянен".
|
||||||
|
unchanged_last_seen_at = now - timedelta(days=_HEALTH_WINDOW_DAYS + 3)
|
||||||
|
try:
|
||||||
|
db.execute(text("BEGIN"))
|
||||||
|
lsid = _insert_pair(
|
||||||
|
db,
|
||||||
|
source=source,
|
||||||
|
ext_id="dense",
|
||||||
|
listing_last_seen_at=now,
|
||||||
|
snapshots=[
|
||||||
|
(anchor - timedelta(days=3), unchanged_last_seen_at),
|
||||||
|
(anchor - timedelta(days=2), unchanged_last_seen_at),
|
||||||
|
(anchor - timedelta(days=1), unchanged_last_seen_at),
|
||||||
|
(anchor, unchanged_last_seen_at), # якорная дата ПРИСУТСТВУЕТ
|
||||||
|
],
|
||||||
|
)
|
||||||
|
floor_dense = _floor(db, source=source)
|
||||||
|
assert floor_dense is not None
|
||||||
|
assert _n_pairs(db, source=source) == 1
|
||||||
|
|
||||||
|
# Change-only модель: сутки без изменений не пишутся -- удаляем всё, кроме
|
||||||
|
# самого раннего снимка (якорная дата больше НЕ присутствует ровно).
|
||||||
|
db.execute(
|
||||||
|
text(
|
||||||
|
"DELETE FROM listing_source_snapshots "
|
||||||
|
"WHERE listing_source_id = :lsid AND snapshot_date > :keep_date"
|
||||||
|
),
|
||||||
|
{"lsid": lsid, "keep_date": anchor - timedelta(days=3)},
|
||||||
|
)
|
||||||
|
floor_sparse = _floor(db, source=source)
|
||||||
|
assert floor_sparse is not None
|
||||||
|
assert _n_pairs(db, source=source) == 1
|
||||||
|
|
||||||
|
assert floor_sparse == pytest.approx(floor_dense, abs=0.01), (
|
||||||
|
f"LATERAL обязан игнорировать гэп: сплошная история дала {floor_dense}, "
|
||||||
|
f"дырявая -- {floor_sparse}, а между изменениями last_seen_at постоянен"
|
||||||
|
)
|
||||||
|
# Возраст известен точно (постоянный last_seen_at, N+3 суток разрыва) --
|
||||||
|
# пиним не только "совпадают", но и КАКОЕ именно число.
|
||||||
|
assert floor_dense == pytest.approx(_HEALTH_WINDOW_DAYS + 3, abs=0.01)
|
||||||
|
finally:
|
||||||
|
db.rollback()
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB")
|
||||||
|
def test_equality_join_loses_gapped_pairs_and_gives_a_smaller_floor() -> None:
|
||||||
|
"""На ОДНИХ И ТЕХ ЖЕ данных equality-join теряет пары и даёт МЕНЬШИЙ пол -- это
|
||||||
|
и есть суть правки PR-A: 3 строки, только у одной снимок ровно на дату якоря.
|
||||||
|
|
||||||
|
Group A -- снимок РОВНО на дату якоря, возраст ~5.5 сут (выживает в ОБЕИХ моделях).
|
||||||
|
Group B, C -- снимок ТОЛЬКО раньше якоря (change-only гэп ровно на дату якоря),
|
||||||
|
возраст 65 и 30 сут -- выживают ТОЛЬКО в LATERAL.
|
||||||
|
|
||||||
|
revisit_quantile=1.0 (максимум) -- детерминированно и без интерполяции percentile_disc
|
||||||
|
на маленькой выборке: LATERAL обязан дать max(5.5, 65, 30) = 65, equality-join --
|
||||||
|
только 5.5 (единственная сохранившаяся пара).
|
||||||
|
"""
|
||||||
|
db = _live_session()
|
||||||
|
source = f"zzzlat_gap_{uuid.uuid4().hex[:8]}"
|
||||||
|
now = datetime.now(UTC)
|
||||||
|
anchor = date.today() - timedelta(days=_HEALTH_WINDOW_DAYS)
|
||||||
|
try:
|
||||||
|
db.execute(text("BEGIN"))
|
||||||
|
# Group A: снимок ровно на дату якоря -- переживает equality-join.
|
||||||
|
_insert_pair(
|
||||||
|
db,
|
||||||
|
source=source,
|
||||||
|
ext_id="group-a-on-anchor",
|
||||||
|
listing_last_seen_at=now,
|
||||||
|
snapshots=[(anchor, now - timedelta(days=5, hours=12))], # возраст 5.5 сут
|
||||||
|
)
|
||||||
|
# Group B: снимок только на anchor-60 -- дыра ровно на дату якоря
|
||||||
|
# (типичный change-only разрыв: ничего не менялось 55+ суток подряд).
|
||||||
|
_insert_pair(
|
||||||
|
db,
|
||||||
|
source=source,
|
||||||
|
ext_id="group-b-gap-60d",
|
||||||
|
listing_last_seen_at=now,
|
||||||
|
snapshots=[(anchor - timedelta(days=60), now - timedelta(days=65))],
|
||||||
|
)
|
||||||
|
# Group C: тот же класс гэпа, гэп короче (30 сут) -- контроль, что LATERAL
|
||||||
|
# берёт максимум по ВСЕЙ выборке, а не просто "последнюю вставленную пару".
|
||||||
|
_insert_pair(
|
||||||
|
db,
|
||||||
|
source=source,
|
||||||
|
ext_id="group-c-gap-25d",
|
||||||
|
listing_last_seen_at=now,
|
||||||
|
snapshots=[(anchor - timedelta(days=25), now - timedelta(days=30))],
|
||||||
|
)
|
||||||
|
|
||||||
|
lateral_floor = _floor(db, source=source)
|
||||||
|
lateral_n_pairs = _n_pairs(db, source=source)
|
||||||
|
old_floor = _old_floor(db, source=source)
|
||||||
|
old_n_pairs = _old_n_pairs(db, source=source)
|
||||||
|
|
||||||
|
assert lateral_n_pairs == 3, "LATERAL обязан видеть все 3 пары, гэпы не теряют строк"
|
||||||
|
assert old_n_pairs == 1, "equality-join теряет Group B и C -- у них нет снимка на якоре"
|
||||||
|
|
||||||
|
assert lateral_floor == pytest.approx(65.0, abs=0.01), (
|
||||||
|
f"LATERAL max должен взять Group B (65 сут), получили {lateral_floor}"
|
||||||
|
)
|
||||||
|
assert old_floor == pytest.approx(5.5, abs=0.01), (
|
||||||
|
f"equality-join должен остаться только с Group A (5.5 сут), получили {old_floor}"
|
||||||
|
)
|
||||||
|
assert lateral_floor > old_floor, (
|
||||||
|
"равенство по дате даёт МЕНЬШИЙ пол на тех же данных -- это и есть регрессия, "
|
||||||
|
"которую эта правка чинит"
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
db.rollback()
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB")
|
||||||
|
def test_missing_history_before_anchor_still_yields_no_floor_via_lateral() -> None:
|
||||||
|
"""Контроль: если у listing_source_id вообще нет снимка не позже якоря (совсем
|
||||||
|
свежая строка), LATERAL закономерно не находит пару -- НЕ ломается на NULL/пусто,
|
||||||
|
а просто исключает строку (как и раньше исключал equality-join, если совсем
|
||||||
|
нет истории). Это НЕ регрессия, а ожидаемое поведение при отсутствии данных."""
|
||||||
|
db = _live_session()
|
||||||
|
source = f"zzzlat_nohist_{uuid.uuid4().hex[:8]}"
|
||||||
|
now = datetime.now(UTC)
|
||||||
|
anchor = date.today() - timedelta(days=_HEALTH_WINDOW_DAYS)
|
||||||
|
try:
|
||||||
|
db.execute(text("BEGIN"))
|
||||||
|
_insert_pair(
|
||||||
|
db,
|
||||||
|
source=source,
|
||||||
|
ext_id="only-future-snapshot",
|
||||||
|
listing_last_seen_at=now,
|
||||||
|
# Снимок ЕСТЬ, но он ПОЗЖЕ якоря -- LATERAL (snapshot_date <= якорь)
|
||||||
|
# его не видит, что и требуется: свежая строка без прошлого не должна
|
||||||
|
# выдумывать пол из снимка, который сам моложе окна здоровья.
|
||||||
|
snapshots=[(anchor + timedelta(days=1), now - timedelta(days=1))],
|
||||||
|
)
|
||||||
|
assert _floor(db, source=source) is None
|
||||||
|
assert _n_pairs(db, source=source) == 0
|
||||||
|
finally:
|
||||||
|
db.rollback()
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
|
||||||
|
# ── Pure-SQL проверки (без БД) -- форма LATERAL-запроса и общего среза ────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_lateral_sql_has_no_equality_join_on_snapshot_date() -> None:
|
||||||
|
"""Ключевой негативный инвариант правки: РАВЕНСТВО по дате запрещено."""
|
||||||
|
sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=False).text)
|
||||||
|
assert "JOIN LATERAL" in sql
|
||||||
|
assert "ORDER BY s.snapshot_date DESC" in sql
|
||||||
|
assert "LIMIT 1" in sql
|
||||||
|
assert not re.search(r"prev\.snapshot_date\s*=", sql), (
|
||||||
|
"equality-join по snapshot_date должен быть полностью удалён из пола переобхода"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_lateral_sql_still_psycopg_v3_safe() -> None:
|
||||||
|
sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=True).text)
|
||||||
|
assert "CAST(:revisit_quantile AS double precision)" in sql
|
||||||
|
assert "CAST(:health_window_days AS integer)" in sql
|
||||||
|
assert not re.search(r":\w+::", sql)
|
||||||
|
|
||||||
|
|
||||||
|
def test_pairs_count_sql_shares_the_exact_same_slice_as_the_floor_sql() -> None:
|
||||||
|
"""count(*) и percentile_disc обязаны идти по ОДНОМУ И ТОМУ ЖЕ FROM..WHERE --
|
||||||
|
иначе floor_n_pairs в counters не описывает реальную выборку percentile_disc."""
|
||||||
|
floor_sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=True).text)
|
||||||
|
count_sql = str(
|
||||||
|
task_mod._build_revisit_floor_pairs_count_sql("last_seen_at", with_segments=True).text
|
||||||
|
)
|
||||||
|
floor_from_where = floor_sql.split("FROM listings l", 1)[1]
|
||||||
|
count_from_where = count_sql.split("FROM listings l", 1)[1]
|
||||||
|
assert floor_from_where == count_from_where, (
|
||||||
|
"FROM..WHERE percentile_disc-запроса и count(*)-запроса разошлись -- "
|
||||||
|
"floor_n_pairs больше не описывает реальную выборку пола"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_comparison_predicate_against_prev_last_seen_at_is_untouched() -> None:
|
||||||
|
"""Предикат сравнения (l.<col> > prev.last_seen_at) -- НЕ часть этой правки,
|
||||||
|
LATERAL меняет только то, КАК ищется prev, а не что с ним сравнивается."""
|
||||||
|
sql = str(task_mod._build_revisit_floor_sql("last_seen_at", with_segments=False).text)
|
||||||
|
assert "l.last_seen_at > prev.last_seen_at" in sql
|
||||||
Loading…
Add table
Reference in a new issue