Merge pull request 'Снимки источников МЕРА пишутся только при изменении — в 26–38 раз меньше строк за ночь (#2993, PR-C)' (#3577) from fix/snapshot-change-only into main
Some checks failed
Deploy Trade-In / test (push) Blocked by required conditions
Deploy Trade-In / build-backend (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Blocked by required conditions
Deploy Trade-In / build-browser (push) Blocked by required conditions
Deploy Trade-In / deploy (push) Blocked by required conditions
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / changes (push) Has been cancelled

This commit is contained in:
bot-backend 2026-09-17 11:23:00 +00:00
commit 5b6be30894
3 changed files with 328 additions and 11 deletions

View file

@ -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,
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),
now(),
CAST(:run_id AS bigint)
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.

View file

@ -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;

View 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()