fix(tradein/snapshot): снимок источника пишется только при изменении (#2993, PR-C)
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
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
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>
This commit is contained in:
parent
a130303cc9
commit
e4d44025e5
3 changed files with 328 additions and 11 deletions
|
|
@ -1,7 +1,8 @@
|
|||
"""Daily per-source snapshot writer (#570).
|
||||
|
||||
Берёт текущее состояние listing_sources (последний снимок на canonical listing × source)
|
||||
и пишет ежедневный ряд в listing_source_snapshots + change-log price_change в
|
||||
и раз в сутки пишет в listing_source_snapshots строку на КАЖДОЕ ИЗМЕНЕНИЕ (#2993, см.
|
||||
_SNAPSHOT_SQL; раньше — полная копия каждые сутки) + change-log в
|
||||
listing_source_events. Так история per-source цены копится СРАЗУ, независимо от
|
||||
(сейчас DORMANT) скраперов — см. шапку data/sql/079_listing_source_history.sql.
|
||||
|
||||
|
|
@ -72,9 +73,28 @@ def _clamp_budget_sec(raw: Any) -> float:
|
|||
return max(_MIN_BUDGET_SEC, min(val, _MAX_BUDGET_SEC))
|
||||
|
||||
|
||||
# ── Daily snapshot upsert ─────────────────────────────────────────────────────
|
||||
# Снимок на (listing_source_id, CURRENT_DATE). ON CONFLICT → last-write-wins за день
|
||||
# (повторный прогон в те же сутки перезаписывает снимок свежими значениями).
|
||||
# ── 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).
|
||||
|
|
@ -85,15 +105,33 @@ _SNAPSHOT_SQL = text(
|
|||
last_seen_at, payload_hash, observed_at, run_id
|
||||
)
|
||||
SELECT
|
||||
id,
|
||||
cur.id,
|
||||
CURRENT_DATE,
|
||||
price_rub,
|
||||
(last_seen_at > now() - make_interval(days => :freshness_days)) AS is_active,
|
||||
last_seen_at,
|
||||
md5(raw_payload::text),
|
||||
cur.price_rub,
|
||||
cur.is_active,
|
||||
cur.last_seen_at,
|
||||
cur.payload_hash,
|
||||
now(),
|
||||
CAST(:run_id AS bigint)
|
||||
FROM listing_sources
|
||||
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,
|
||||
|
|
@ -108,6 +146,9 @@ _SNAPSHOT_SQL = text(
|
|||
# Для каждого источника сравниваем сегодняшний снимок (snapshot_date = CURRENT_DATE) с
|
||||
# самым свежим ПРЕДЫДУЩИМ (snapshot_date < CURRENT_DATE).
|
||||
# today — снимок за сегодня (только что записан _SNAPSHOT_SQL, в той же транзакции).
|
||||
# С #2993 здесь только изменившиеся и новые источники: у остальных все четыре
|
||||
# значения равны последнему снимку, а значит ни одно событие ниже сработать
|
||||
# не могло бы — набор событий тот же, что при суточной копии.
|
||||
# p — последний снимок строго ДО сегодня, per-row LATERAL point-lookup (#2607).
|
||||
#
|
||||
# #2674: схема (079) знает пять типов событий, писатель умел один — price_change,
|
||||
|
|
@ -249,7 +290,8 @@ def snapshot_listing_sources(
|
|||
|
||||
Sync (вызывается scheduler-триггером в executor, как import_rosreestr_dkp).
|
||||
Два set-based statement'а в одной транзакции:
|
||||
1. upsert снимка на (listing_source_id, CURRENT_DATE) — last-write-wins.
|
||||
1. upsert снимка на (listing_source_id, CURRENT_DATE) только для источников,
|
||||
изменившихся с последнего снимка (#2993) — last-write-wins.
|
||||
2. diff сегодняшнего снимка против последнего предыдущего → три события,
|
||||
выводимые из наших данных (#2674). delisted/relisted схема разрешает, но
|
||||
они НЕ выводимы при покрытии обхода 10-35% — см. _EVENT_DIFF_SQL.
|
||||
|
|
|
|||
|
|
@ -0,0 +1,29 @@
|
|||
-- 325_listing_source_snapshots_change_only_comment.sql
|
||||
-- #2993 — listing_source_snapshots переходит на «строку на изменение»
|
||||
-- (app/tasks/listing_source_snapshot.py, _SNAPSHOT_SQL; решение владельца 2026-08-23).
|
||||
--
|
||||
-- Схема не меняется, меняется смысл строк — поэтому только комментарии. Прежние
|
||||
-- («Ежедневный снимок…», «на дату снимка») после перехода стали бы неправдой:
|
||||
-- `WHERE snapshot_date = D` отдаёт теперь не состояние всех источников на D, а только
|
||||
-- изменившиеся в D. Состояние на дату — последний снимок с snapshot_date <= D.
|
||||
-- Суточные строки до перехода остаются в таблице как есть.
|
||||
--
|
||||
-- COMMENT берёт ShareUpdateExclusiveLock (конфликтует с autovacuum/ANALYZE, не с
|
||||
-- чтением и записью) — lock_timeout, чтобы деплой не повис за VACUUM 8-млн таблицы.
|
||||
|
||||
BEGIN;
|
||||
|
||||
SET LOCAL lock_timeout = '5s';
|
||||
|
||||
COMMENT ON TABLE listing_source_snapshots IS
|
||||
'Снимок listing_sources на источник (#570), строка только при изменении (#2993): '
|
||||
'пишется, если price_rub/is_active/last_seen_at/payload_hash отличаются от последнего '
|
||||
'снимка источника или снимка ещё нет. Состояние на дату D — последняя строка с '
|
||||
'snapshot_date <= D. До перехода (#2993) — полная суточная копия. '
|
||||
'PK (listing_source_id, snapshot_date), last-write-wins за день.';
|
||||
|
||||
COMMENT ON VIEW v_listing_source_price_on_date IS
|
||||
'Цена/активность источника с даты снимка (#570): listing_source_snapshots ⋈ listing_sources. '
|
||||
'С #2993 строки только на даты изменений — цена на дату D = последняя строка с snapshot_date <= D.';
|
||||
|
||||
COMMIT;
|
||||
246
tradein-mvp/backend/tests/test_2993_snapshot_change_only.py
Normal file
246
tradein-mvp/backend/tests/test_2993_snapshot_change_only.py
Normal file
|
|
@ -0,0 +1,246 @@
|
|||
"""#2993: listing_source_snapshots получает строку только при изменении источника.
|
||||
|
||||
Проверки идут по ЗНАЧЕНИЯМ в живом Postgres: какие источники получили строку за
|
||||
сегодня, с какими значениями, и какие события записал event-diff поверх них.
|
||||
Без реальной БД файл self-skip'ается, как соседние live-тесты (в CI есть
|
||||
postgres, см. ci-tradein.yml). Вся синтетика живёт в одной транзакции, очистка —
|
||||
rollback.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import uuid
|
||||
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 listing_source_snapshot as snap_mod
|
||||
|
||||
|
||||
def _live_session() -> Any | None:
|
||||
"""Тот же контракт, что у test_revisit_floor_lateral_lookup.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
|
||||
|
||||
|
||||
pytestmark = pytest.mark.skipif(_live_session() is None, reason="no reachable Postgres test DB")
|
||||
|
||||
_PRICE = 5_000_000
|
||||
|
||||
|
||||
def _source(db: Any, tag: str, *, price: int, payload: str, seen_days_ago: int) -> int:
|
||||
ext_id = f"zzz2993-{tag}-{uuid.uuid4().hex[:8]}"
|
||||
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 ('cian', :url, :ext_id, :ext_id, :price, true, now(), now())
|
||||
RETURNING id
|
||||
"""
|
||||
),
|
||||
{"url": f"https://example.test/2993/{ext_id}", "ext_id": ext_id, "price": price},
|
||||
).scalar_one()
|
||||
return int(
|
||||
db.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO listing_sources
|
||||
(listing_id, ext_source, ext_id, confidence, matched_method,
|
||||
price_rub, raw_payload, last_seen_at)
|
||||
VALUES (:listing_id, 'cian', :ext_id, 1.0, 'test',
|
||||
:price, CAST(:payload AS jsonb),
|
||||
now() - make_interval(days => :seen))
|
||||
RETURNING id
|
||||
"""
|
||||
),
|
||||
{
|
||||
"listing_id": listing_id,
|
||||
"ext_id": ext_id,
|
||||
"price": price,
|
||||
"payload": payload,
|
||||
"seen": seen_days_ago,
|
||||
},
|
||||
).scalar_one()
|
||||
)
|
||||
|
||||
|
||||
def _snapshot(
|
||||
db: Any,
|
||||
lsid: int,
|
||||
*,
|
||||
days_ago: int,
|
||||
price: int,
|
||||
payload: str,
|
||||
seen_days_ago: int,
|
||||
is_active: bool = True,
|
||||
) -> None:
|
||||
# now() постоянен внутри транзакции — last_seen_at снимка совпадает с
|
||||
# listing_sources.last_seen_at до микросекунды, если seen_days_ago тот же.
|
||||
db.execute(
|
||||
text(
|
||||
"""
|
||||
INSERT INTO listing_source_snapshots
|
||||
(listing_source_id, snapshot_date, price_rub, is_active,
|
||||
last_seen_at, payload_hash)
|
||||
VALUES (:lsid, CURRENT_DATE - :days_ago, :price, :is_active,
|
||||
now() - make_interval(days => :seen),
|
||||
md5(CAST(:payload AS jsonb)::text))
|
||||
"""
|
||||
),
|
||||
{
|
||||
"lsid": lsid,
|
||||
"days_ago": days_ago,
|
||||
"price": price,
|
||||
"is_active": is_active,
|
||||
"seen": seen_days_ago,
|
||||
"payload": payload,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def _run_writer(db: Any) -> None:
|
||||
db.execute(
|
||||
snap_mod._SNAPSHOT_SQL,
|
||||
{"freshness_days": snap_mod.FRESHNESS_WINDOW_DAYS, "run_id": None},
|
||||
)
|
||||
|
||||
|
||||
def _today(db: Any, ids: list[int]) -> dict[int, tuple[int, bool]]:
|
||||
rows = db.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT listing_source_id, price_rub, is_active
|
||||
FROM listing_source_snapshots
|
||||
WHERE snapshot_date = CURRENT_DATE AND listing_source_id = ANY(:ids)
|
||||
"""
|
||||
),
|
||||
{"ids": ids},
|
||||
).fetchall()
|
||||
return {int(r[0]): (int(r[1]), bool(r[2])) for r in rows}
|
||||
|
||||
|
||||
def _seed(db: Any) -> dict[str, int]:
|
||||
"""Семь источников: по одному на каждую ветку условия записи."""
|
||||
p = '{"v": 1}'
|
||||
ids = {
|
||||
"same": _source(db, "same", price=_PRICE, payload=p, seen_days_ago=1),
|
||||
"same_gap": _source(db, "same_gap", price=_PRICE, payload=p, seen_days_ago=1),
|
||||
"price": _source(db, "price", price=_PRICE, payload=p, seen_days_ago=1),
|
||||
"payload": _source(db, "payload", price=_PRICE, payload=p, seen_days_ago=1),
|
||||
"seen": _source(db, "seen", price=_PRICE, payload=p, seen_days_ago=1),
|
||||
"active": _source(db, "active", price=_PRICE, payload=p, seen_days_ago=10),
|
||||
"new": _source(db, "new", price=_PRICE, payload=p, seen_days_ago=1),
|
||||
}
|
||||
_snapshot(db, ids["same"], days_ago=1, price=_PRICE, payload=p, seen_days_ago=1)
|
||||
# Последний снимок пятидневной давности совпадает, а более старый — нет:
|
||||
# сравнивать обязаны с ПОСЛЕДНИМ, а не с любым из прошлых.
|
||||
_snapshot(db, ids["same_gap"], days_ago=9, price=_PRICE - 1, payload=p, seen_days_ago=9)
|
||||
_snapshot(db, ids["same_gap"], days_ago=5, price=_PRICE, payload=p, seen_days_ago=1)
|
||||
_snapshot(db, ids["price"], days_ago=1, price=4_900_000, payload=p, seen_days_ago=1)
|
||||
_snapshot(db, ids["payload"], days_ago=1, price=_PRICE, payload='{"v": 0}', seen_days_ago=1)
|
||||
_snapshot(db, ids["seen"], days_ago=1, price=_PRICE, payload=p, seen_days_ago=3)
|
||||
# last_seen_at тот же, но 10 суток > окна свежести: вчера снимок был активным.
|
||||
_snapshot(db, ids["active"], days_ago=1, price=_PRICE, payload=p, seen_days_ago=10)
|
||||
return ids
|
||||
|
||||
|
||||
def test_only_changed_or_new_sources_get_a_row_today() -> None:
|
||||
db = _live_session()
|
||||
try:
|
||||
ids = _seed(db)
|
||||
_run_writer(db)
|
||||
today = _today(db, list(ids.values()))
|
||||
by_tag = {tag: today.get(lsid) for tag, lsid in ids.items()}
|
||||
|
||||
assert by_tag == {
|
||||
"same": None,
|
||||
"same_gap": None,
|
||||
"price": (_PRICE, True),
|
||||
"payload": (_PRICE, True),
|
||||
"seen": (_PRICE, True),
|
||||
"active": (_PRICE, False),
|
||||
"new": (_PRICE, True),
|
||||
}
|
||||
finally:
|
||||
db.rollback()
|
||||
db.close()
|
||||
|
||||
|
||||
def test_events_survive_change_only_writes() -> None:
|
||||
"""event-diff видит только записанные сегодня строки — и не теряет ни одного события."""
|
||||
db = _live_session()
|
||||
try:
|
||||
ids = _seed(db)
|
||||
_run_writer(db)
|
||||
db.execute(snap_mod._EVENT_DIFF_SQL).fetchall()
|
||||
rows = db.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT listing_source_id, event_type, diff_percent
|
||||
FROM listing_source_events
|
||||
WHERE listing_source_id = ANY(:ids)
|
||||
"""
|
||||
),
|
||||
{"ids": list(ids.values())},
|
||||
).fetchall()
|
||||
tag_of = {lsid: tag for tag, lsid in ids.items()}
|
||||
events = {(tag_of[int(r[0])], r[1], None if r[2] is None else float(r[2])) for r in rows}
|
||||
|
||||
assert events == {
|
||||
("price", "price_change", 2.0408), # round((5.0-4.9)/4.9*100, 4)
|
||||
("payload", "edited", None),
|
||||
("new", "first_seen", None),
|
||||
}
|
||||
finally:
|
||||
db.rollback()
|
||||
db.close()
|
||||
|
||||
|
||||
def test_rerun_same_day_compares_with_todays_row() -> None:
|
||||
"""Второй прогон в те же сутки сравнивает с уже записанной сегодня строкой.
|
||||
|
||||
price вернулся к вчерашней цене — сегодняшняя строка обязана перезаписаться
|
||||
(иначе на сегодня осталась бы цена, которой у источника больше нет); same
|
||||
изменился после первого прогона — строка появляется; seen не менялся после
|
||||
первого прогона — строка остаётся как была.
|
||||
"""
|
||||
db = _live_session()
|
||||
try:
|
||||
ids = _seed(db)
|
||||
_run_writer(db)
|
||||
db.execute(
|
||||
text("UPDATE listing_sources SET price_rub = 4900000 WHERE id = :id"),
|
||||
{"id": ids["price"]},
|
||||
)
|
||||
db.execute(
|
||||
text("UPDATE listing_sources SET price_rub = 5100000 WHERE id = :id"),
|
||||
{"id": ids["same"]},
|
||||
)
|
||||
_run_writer(db)
|
||||
today = _today(db, list(ids.values()))
|
||||
|
||||
assert today.get(ids["price"]) == (4_900_000, True)
|
||||
assert today.get(ids["same"]) == (5_100_000, True)
|
||||
assert today.get(ids["seen"]) == (_PRICE, True)
|
||||
assert ids["same_gap"] not in today
|
||||
finally:
|
||||
db.rollback()
|
||||
db.close()
|
||||
Loading…
Add table
Reference in a new issue