fix(tradein/scrapers): метка наблюдения — время строки, а не старта транзакции (#2731)
All checks were successful
CI / changes (pull_request) Successful in 13s
CI Trade-In / changes (pull_request) Successful in 13s
CI Trade-In / browser-tests (pull_request) Has been skipped
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 3m36s
All checks were successful
CI / changes (pull_request) Successful in 13s
CI Trade-In / changes (pull_request) Successful in 13s
CI Trade-In / browser-tests (pull_request) Has been skipped
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 3m36s
save_listings и upsert_listing_snapshot писали observed_at/scraped_at/last_seen_at через NOW() == transaction_timestamp(). Коммит у писателя один, в конце всего batch'а, поэтому все строки вызова несли одну метку: прод 2026-08-06 — run 3303 219 строк на 1 метку при семи минутах работы, run 3229 2643 строки на 55 меток (метка на вызов save_listings, не на строку). Цена — не фильтры свежести: смещение равно длительности прогона (минуты; худший случай 17.3 ч у cian_full_load) против окна в 14 суток, то есть 0.03-5%. Теряется разрешение во времени — темп сбора и порядок строк внутри прогона по данным не восстановить, отсюда же следовала невозможность бэкфилла run_id (#2701). statement_timestamp(), а НЕ clock_timestamp() как в #2702/#2718: здесь одним statement'ом пишутся ДВЕ колонки, обязанные совпадать (после #2206 scraped_at = last_seen_at у 100% строк, на этом равенстве стоит предикат миграции 161). Прод-проверка: clock_timestamp() = clock_timestamp() → false, а statement_timestamp() = statement_timestamp() → true; при этом statement_timestamp() двигается от запроса к запросу внутри одной транзакции, а каждый upsert — свой запрос. Set-based писатели (listing_source_snapshots 89827 строк одним statement'ом, deactivate_stale_avito) намеренно НЕ тронуты: там одна метка — правда о statement'е, а clock_timestamp() дал бы ложное разрешение (прод: 27783 разных значения на 200000 строк одного statement'а, кодирующих порядок обработки, а не наблюдения). Миграция 232 — только COMMENT ON COLUMN: фиксирует дату перехода и то, что исторические строки невосстановимы (внутрипрогонное время нигде не сохранялось). Refs #2731
This commit is contained in:
parent
c86a5378ef
commit
54714f0aa2
5 changed files with 423 additions and 13 deletions
|
|
@ -0,0 +1,84 @@
|
||||||
|
-- 232_listings_observation_time_meaning.sql
|
||||||
|
-- Purpose (#2731): зафиксировать в схеме, что содержательные метки наблюдения до
|
||||||
|
-- этой правки несли время НАЧАЛА транзакции сбора, а не момент наблюдения строки,
|
||||||
|
-- и назвать дату, с которой их смысл изменился.
|
||||||
|
--
|
||||||
|
-- Dependencies: 002_core_tables.sql (listings.scraped_at/last_seen_at),
|
||||||
|
-- 016_listings_snapshots.sql (listings_snapshots.observed_at),
|
||||||
|
-- 161_backfill_scraped_at_active_recent.sql (ретро-выравнивание scraped_at по last_seen_at).
|
||||||
|
-- Apply after: 230_house_merge_log.sql
|
||||||
|
-- Идемпотентно: только COMMENT ON COLUMN (перезаписывает сам себя), данных не трогает.
|
||||||
|
--
|
||||||
|
-- ── ЧТО БЫЛО ─────────────────────────────────────────────────────────────────
|
||||||
|
-- Писатель объявлений (packages/scraper-kit/.../base.py::save_listings + snapshot_writer)
|
||||||
|
-- ставил NOW(). В PostgreSQL NOW() == transaction_timestamp() — время старта транзакции.
|
||||||
|
-- save_listings коммитит ОДИН раз в конце всего batch'а, поэтому одну метку получали все
|
||||||
|
-- строки одного вызова, сколько бы он ни работал.
|
||||||
|
--
|
||||||
|
-- Прод-замер 2026-08-06 (listings_snapshots, строк / различных меток на прогон):
|
||||||
|
-- run 3303 — 219 / 1 (прогон шёл 7 минут, все метки на нулевой секунде);
|
||||||
|
-- run 3299 — 235 / 1; run 3293 — 297 / 1;
|
||||||
|
-- run 3229 — 2643 / 55 (метка на вызов save_listings, а не на строку).
|
||||||
|
-- listings за те же сутки: 219/1, 235/1, 297/1 … — и scraped_at, и last_seen_at, причём
|
||||||
|
-- у 100% строк они РАВНЫ между собой (замер по часам: eq == rows во всех корзинах).
|
||||||
|
--
|
||||||
|
-- ── ЧЕГО ЭТО НЕ ЛОМАЛО ───────────────────────────────────────────────────────
|
||||||
|
-- Фильтры свежести и TTL — НЕ искажены. Смещение равно длительности прогона: типично
|
||||||
|
-- минуты, худший случай на проде 17.3 ч (cian_full_load). Против окна свежести в 14 суток
|
||||||
|
-- это 0.03-5%. Утверждение «долгие прогоны ломают фильтр свежести» проверено и снято.
|
||||||
|
--
|
||||||
|
-- ── ЧТО ЭТО ЛОМАЛО ───────────────────────────────────────────────────────────
|
||||||
|
-- Разрешение во времени. По данным нельзя восстановить ни темп сбора, ни порядок строк
|
||||||
|
-- внутри прогона: все они выглядят одномоментными. Отсюда же следовала невозможность
|
||||||
|
-- бэкфилла run_id (#2701). Ошибка тихая — значения правдоподобны.
|
||||||
|
--
|
||||||
|
-- ── ПОЧЕМУ ИСТОРИЮ НЕ ЧИНИМ ──────────────────────────────────────────────────
|
||||||
|
-- Внутрипрогонное время НИГДЕ БОЛЬШЕ НЕ СОХРАНЯЛОСЬ: у прогона есть только started_at и
|
||||||
|
-- finished_at, а распределение строк между ними неизвестно. Строки ДО перехода
|
||||||
|
-- невосстановимы — их метки помечаются, а не переписываются.
|
||||||
|
--
|
||||||
|
-- ── ПОЧЕМУ statement_timestamp(), А НЕ clock_timestamp() ─────────────────────
|
||||||
|
-- В #2702/#2718 (служебные колонки прогона) взяли clock_timestamp() — там колонка одна.
|
||||||
|
-- Здесь в одном statement'е пишутся ДВЕ колонки, и они обязаны совпадать: после #2206
|
||||||
|
-- scraped_at и last_seen_at равны, и на этом равенстве стоит предикат миграции 161
|
||||||
|
-- (`WHERE last_seen_at > scraped_at` как признак «видели живым, но не пере-скрейпили»).
|
||||||
|
-- Проверка на проде: `clock_timestamp() = clock_timestamp()` → false,
|
||||||
|
-- `statement_timestamp() = statement_timestamp()` → true. При этом statement_timestamp()
|
||||||
|
-- двигается ОТ STATEMENT'А К STATEMENT'У внутри одной транзакции (проверено: 1.2 с между
|
||||||
|
-- соседними запросами при неподвижном now()), а каждый upsert объявления — свой statement.
|
||||||
|
-- Итог: построчная метка без ложного расхождения колонок.
|
||||||
|
--
|
||||||
|
-- Set-based писатели (listing_source_snapshots, deactivate_stale_avito) НАМЕРЕННО оставлены
|
||||||
|
-- на now(): там один statement пишет десятки тысяч строк, и одна метка — это правда о нём.
|
||||||
|
-- clock_timestamp() выдал бы там ложное разрешение: на проде
|
||||||
|
-- `count(DISTINCT clock_timestamp())` по 200 000 строк одного statement'а = 27 783 разных
|
||||||
|
-- значения, кодирующих порядок обработки строк планировщиком, а не порядок наблюдения.
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
COMMENT ON COLUMN listings_snapshots.observed_at IS
|
||||||
|
'Момент наблюдения снимка. С #2731 (2026-08-06) — statement_timestamp(), то есть '
|
||||||
|
'время записи КОНКРЕТНОЙ строки. У строк ДО этой даты значение общее на весь вызов '
|
||||||
|
'save_listings (писалось now() = старт транзакции batch''а): на проде 219 снимков '
|
||||||
|
'семиминутного прогона несли одну метку. Для строк до перехода observed_at читать как '
|
||||||
|
'«прогон, в котором строку увидели», а НЕ как момент; темп сбора и порядок строк '
|
||||||
|
'внутри прогона по историческим данным невосстановимы — внутрипрогонное время нигде '
|
||||||
|
'не сохранялось.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN listings.scraped_at IS
|
||||||
|
'Момент последнего скрейпа объявления. С #2206 двигается и при ре-подтверждении живым '
|
||||||
|
'(не только при вставке), с #2731 (2026-08-06) пишется statement_timestamp() — '
|
||||||
|
'построчно. У строк ДО этой даты — время старта транзакции batch''а, общее на все '
|
||||||
|
'строки вызова save_listings. На фильтры свежести это не влияло: смещение равно '
|
||||||
|
'длительности прогона (минуты; худший случай 17.3 ч у cian_full_load) против окна в '
|
||||||
|
'14 суток. Невосстановимо потеряно другое — разрешение внутри прогона.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN listings.last_seen_at IS
|
||||||
|
'Момент последнего подтверждения, что объявление живо. Пишется тем же statement''ом, '
|
||||||
|
'что и scraped_at, и потому РАВЕН ему — это свойство сохранено намеренно: '
|
||||||
|
'с #2731 (2026-08-06) обе колонки берут statement_timestamp(), стабильный в statement''е '
|
||||||
|
'(clock_timestamp() развёл бы их на микросекунды и сделал бы предикат '
|
||||||
|
'`last_seen_at > scraped_at` из миграции 161 истинным почти для всех строк). '
|
||||||
|
'У строк ДО перехода — время старта транзакции batch''а, см. комментарий к scraped_at.';
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
292
tradein-mvp/backend/tests/test_2731_observation_timestamps.py
Normal file
292
tradein-mvp/backend/tests/test_2731_observation_timestamps.py
Normal file
|
|
@ -0,0 +1,292 @@
|
||||||
|
"""#2731: метка наблюдения — время строки, а не старта транзакции сбора.
|
||||||
|
|
||||||
|
Что было. `save_listings` и `upsert_listing_snapshot` писали содержательные колонки
|
||||||
|
времени через `NOW()`, а `NOW()` в PostgreSQL — синоним `transaction_timestamp()`:
|
||||||
|
он замерзает на СТАРТЕ транзакции. Коммит у писателя объявлений ОДИН, в конце всего
|
||||||
|
batch'а, поэтому все строки вызова получали одну и ту же метку — сколько бы он ни работал.
|
||||||
|
|
||||||
|
Прод-замер 2026-08-06 (listings_snapshots, строк / различных меток на прогон):
|
||||||
|
run 3303 — 219 / 1, при семи минутах работы (11:59:59 → 12:07:19);
|
||||||
|
run 3299 — 235 / 1; run 3293 — 297 / 1;
|
||||||
|
run 3229 — 2643 / 55 — метка на ВЫЗОВ save_listings, а не на строку.
|
||||||
|
`listings.scraped_at` и `last_seen_at` за те же сутки схлопнуты так же и при этом равны
|
||||||
|
друг другу у 100% строк (по часам: eq == rows во всех корзинах).
|
||||||
|
|
||||||
|
Цена дефекта — НЕ фильтры свежести (смещение равно длительности прогона: минуты, худший
|
||||||
|
случай 17.3 ч, против окна 14 суток), а разрешение во времени: темп сбора и порядок строк
|
||||||
|
внутри прогона по данным не восстановить (отсюда же невозможность бэкфилла run_id, #2701).
|
||||||
|
|
||||||
|
Почему `statement_timestamp()`, а не `clock_timestamp()` как в #2702/#2718: здесь в одном
|
||||||
|
statement'е пишутся ДВЕ колонки, обязанные совпадать (после #2206 scraped_at = last_seen_at,
|
||||||
|
на этом равенстве стоит предикат миграции 161). На проде проверено:
|
||||||
|
`clock_timestamp() = clock_timestamp()` → false, `statement_timestamp() = ...` → true.
|
||||||
|
Модель БД ниже воспроизводит ровно это различие, поэтому «починка» через clock_timestamp()
|
||||||
|
провалит тест на согласованность.
|
||||||
|
|
||||||
|
Фальсификация: на старом коде (`NOW()`) тесты 1 и 2 дают одну метку на все строки.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import inspect
|
||||||
|
import os
|
||||||
|
import re
|
||||||
|
from contextlib import contextmanager
|
||||||
|
from pathlib import Path
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||||
|
|
||||||
|
from scraper_kit import base as kit_base
|
||||||
|
from scraper_kit import snapshot_writer as kit_snapshot
|
||||||
|
from scraper_kit.base import ScrapedLot, save_listings
|
||||||
|
|
||||||
|
# Сколько «работает» писатель на одну строку. Пять лотов → 2 с работы одной транзакции:
|
||||||
|
# требование задачи — писатель, отработавший ДОЛЬШЕ СЕКУНДЫ, оставляет различные метки.
|
||||||
|
_STEP_SEC = 0.4
|
||||||
|
|
||||||
|
# Содержательные колонки времени, за которыми следим, по таблицам-писателям.
|
||||||
|
_WATCHED: dict[str, tuple[str, ...]] = {
|
||||||
|
"listings": ("scraped_at", "last_seen_at"),
|
||||||
|
"listings_snapshots": ("observed_at",),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _strip_sql_comments(sql: str) -> str:
|
||||||
|
"""Убрать `--`-комментарии: в них есть и скобки, и слово NOW() в объяснительной прозе."""
|
||||||
|
return re.sub(r"--[^\n]*", "", sql)
|
||||||
|
|
||||||
|
|
||||||
|
def _split_top(items: str) -> list[str]:
|
||||||
|
"""Разбить список SQL-элементов по запятым ВЕРХНЕГО уровня (CAST(:x AS t) — один)."""
|
||||||
|
out: list[str] = []
|
||||||
|
depth = 0
|
||||||
|
cur = ""
|
||||||
|
for ch in items:
|
||||||
|
if ch == "," and depth == 0:
|
||||||
|
out.append(cur.strip())
|
||||||
|
cur = ""
|
||||||
|
continue
|
||||||
|
depth += (ch == "(") - (ch == ")")
|
||||||
|
cur += ch
|
||||||
|
out.append(cur.strip())
|
||||||
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
def _balanced(sql: str, open_idx: int) -> tuple[str, int]:
|
||||||
|
"""Содержимое скобки, открытой на open_idx, и индекс её закрывающей пары."""
|
||||||
|
depth = 0
|
||||||
|
for i in range(open_idx, len(sql)):
|
||||||
|
depth += (sql[i] == "(") - (sql[i] == ")")
|
||||||
|
if depth == 0:
|
||||||
|
return sql[open_idx + 1 : i], i
|
||||||
|
raise AssertionError("несбалансированные скобки в SQL")
|
||||||
|
|
||||||
|
|
||||||
|
class _Row:
|
||||||
|
"""Строка ответа: id/inserted для INSERT INTO listings ... RETURNING."""
|
||||||
|
|
||||||
|
id = 42
|
||||||
|
inserted = False
|
||||||
|
card_hash = None
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeResult:
|
||||||
|
def __init__(self, row: _Row | None) -> None:
|
||||||
|
self._row = row
|
||||||
|
|
||||||
|
def fetchone(self) -> _Row | None:
|
||||||
|
return self._row
|
||||||
|
|
||||||
|
def first(self) -> _Row | None:
|
||||||
|
return self._row
|
||||||
|
|
||||||
|
def scalar_one_or_none(self) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
class _FakePg:
|
||||||
|
"""Мини-модель PostgreSQL на три функции времени и ленивую транзакцию.
|
||||||
|
|
||||||
|
`now()` == `transaction_timestamp()` — замерзает на старте транзакции (autobegin на
|
||||||
|
первом execute, сброс на commit/rollback);
|
||||||
|
`statement_timestamp()` — момент начала ТЕКУЩЕГО statement'а: двигается от запроса к
|
||||||
|
запросу и СТАБИЛЕН внутри запроса;
|
||||||
|
`clock_timestamp()` — настоящие часы, вычисляются на КАЖДЫЙ вызов отдельно, поэтому
|
||||||
|
два вызова в одном statement'е дают разные значения (проверено на проде).
|
||||||
|
|
||||||
|
Каждый execute «работает» _STEP_SEC — так модель отличает построчную метку от общей.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.wall: float = 0.0
|
||||||
|
self.tx_start: float | None = None
|
||||||
|
self.stmt_start: float = 0.0
|
||||||
|
self._clock_calls: int = 0
|
||||||
|
# Что осело в таблицах: список строк {колонка: значение}.
|
||||||
|
self.rows: dict[str, list[dict[str, float]]] = {t: [] for t in _WATCHED}
|
||||||
|
|
||||||
|
# ── модель времени ────────────────────────────────────────────────────────
|
||||||
|
def _value(self, func: str) -> float:
|
||||||
|
if func == "now":
|
||||||
|
assert self.tx_start is not None
|
||||||
|
return self.tx_start
|
||||||
|
if func == "statement_timestamp":
|
||||||
|
return self.stmt_start
|
||||||
|
if func == "clock_timestamp":
|
||||||
|
self._clock_calls += 1 # каждый вызов — своё значение
|
||||||
|
return self.stmt_start + self._clock_calls * 1e-6
|
||||||
|
raise AssertionError(f"неизвестная функция времени: {func}")
|
||||||
|
|
||||||
|
# ── разбор писателя ───────────────────────────────────────────────────────
|
||||||
|
def _record(self, sql: str) -> None:
|
||||||
|
ins = re.search(r"INSERT INTO (listings|listings_snapshots)\s*\(", sql)
|
||||||
|
if ins is not None:
|
||||||
|
table = ins.group(1)
|
||||||
|
cols, end = _balanced(sql, ins.end() - 1)
|
||||||
|
vals_kw = re.compile(r"\bVALUES\s*\(").search(sql, end)
|
||||||
|
assert vals_kw is not None, "INSERT без VALUES"
|
||||||
|
vals, _ = _balanced(sql, vals_kw.end() - 1)
|
||||||
|
row: dict[str, float] = {}
|
||||||
|
for col, val in zip(_split_top(cols), _split_top(vals), strict=False):
|
||||||
|
f = re.fullmatch(r"(\w+)\s*\(\s*\)", val)
|
||||||
|
if col in _WATCHED[table] and f is not None:
|
||||||
|
row[col] = self._value(f.group(1).lower())
|
||||||
|
if row:
|
||||||
|
self.rows[table].append(row)
|
||||||
|
return
|
||||||
|
upd = re.search(r"UPDATE (listings)\b", sql)
|
||||||
|
if upd is not None: # reconcile-путь при дрейфе dedup_hash
|
||||||
|
table = upd.group(1)
|
||||||
|
row = {}
|
||||||
|
for col in _WATCHED[table]:
|
||||||
|
m = re.search(rf"\b{col}\s*=\s*(\w+)\s*\(\s*\)", sql)
|
||||||
|
if m is not None:
|
||||||
|
row[col] = self._value(m.group(1).lower())
|
||||||
|
if row:
|
||||||
|
self.rows[table].append(row)
|
||||||
|
|
||||||
|
# ── интерфейс сессии ──────────────────────────────────────────────────────
|
||||||
|
def execute(self, stmt: Any, params: Any = None) -> _FakeResult:
|
||||||
|
if self.tx_start is None:
|
||||||
|
self.tx_start = self.wall
|
||||||
|
self.stmt_start = self.wall
|
||||||
|
self.wall += _STEP_SEC # statement отработал
|
||||||
|
sql = _strip_sql_comments(str(stmt))
|
||||||
|
self._record(sql)
|
||||||
|
if "INSERT INTO listings (" in sql:
|
||||||
|
return _FakeResult(_Row())
|
||||||
|
return _FakeResult(None)
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
self.tx_start = None
|
||||||
|
|
||||||
|
def rollback(self) -> None:
|
||||||
|
self.tx_start = None
|
||||||
|
|
||||||
|
def begin_nested(self) -> Any:
|
||||||
|
@contextmanager
|
||||||
|
def _ctx() -> Any:
|
||||||
|
yield MagicMock()
|
||||||
|
|
||||||
|
return _ctx()
|
||||||
|
|
||||||
|
|
||||||
|
def _matcher() -> MagicMock:
|
||||||
|
matcher = MagicMock()
|
||||||
|
matcher.match_or_create_house.return_value = (101, 1.0, "new")
|
||||||
|
matcher.upsert_listing_source.return_value = None
|
||||||
|
return matcher
|
||||||
|
|
||||||
|
|
||||||
|
def _lots(n: int = 5) -> list[ScrapedLot]:
|
||||||
|
return [
|
||||||
|
ScrapedLot(
|
||||||
|
source="cian",
|
||||||
|
source_url=f"https://ekb.cian.ru/sale/flat/{i}/",
|
||||||
|
source_id=str(i),
|
||||||
|
price_rub=5_000_000 + i,
|
||||||
|
)
|
||||||
|
for i in range(n)
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def _run_writer() -> _FakePg:
|
||||||
|
db = _FakePg()
|
||||||
|
save_listings(db, _lots(), matcher=_matcher(), region_code=66)
|
||||||
|
return db
|
||||||
|
|
||||||
|
|
||||||
|
# ── 1. Снимки одного batch'а датируются построчно ────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_snapshot_marks_differ_per_row() -> None:
|
||||||
|
"""Прогон 3303 дословно: строки одного вызова — разные моменты, а не один.
|
||||||
|
|
||||||
|
На старом коде (`NOW()`) все снимки batch'а несут метку старта транзакции — ровно
|
||||||
|
219 строк на одну метку, ради чего заведена задача.
|
||||||
|
"""
|
||||||
|
db = _run_writer()
|
||||||
|
marks = [row["observed_at"] for row in db.rows["listings_snapshots"]]
|
||||||
|
|
||||||
|
assert len(marks) == 5, "снимок пишется на каждый сохранённый лот"
|
||||||
|
assert max(marks) - min(marks) > 1.0, "писатель отработал дольше секунды"
|
||||||
|
assert len(set(marks)) == len(marks), "у строк одного batch'а обязаны быть разные метки"
|
||||||
|
|
||||||
|
|
||||||
|
# ── 2. scraped_at/last_seen_at: построчно, но по-прежнему равны между собой ──
|
||||||
|
|
||||||
|
|
||||||
|
def test_listing_marks_differ_per_row_and_stay_equal_to_each_other() -> None:
|
||||||
|
"""Две колонки одного statement'а: построчная метка без нового расхождения.
|
||||||
|
|
||||||
|
Равенство scraped_at = last_seen_at — не косметика: на нём стоит предикат миграции
|
||||||
|
161 (`last_seen_at > scraped_at` как признак «видели живым, но не пере-скрейпили»).
|
||||||
|
`clock_timestamp()` вычисляется на каждый вызов отдельно и развёл бы колонки на
|
||||||
|
микросекунды — модель БД это воспроизводит, поэтому такая «починка» провалит тест.
|
||||||
|
"""
|
||||||
|
db = _run_writer()
|
||||||
|
rows = db.rows["listings"]
|
||||||
|
|
||||||
|
assert len(rows) == 5
|
||||||
|
for row in rows:
|
||||||
|
assert row["scraped_at"] == row["last_seen_at"], "колонки statement'а обязаны совпадать"
|
||||||
|
scraped = [row["scraped_at"] for row in rows]
|
||||||
|
assert max(scraped) - min(scraped) > 1.0
|
||||||
|
assert len(set(scraped)) == len(scraped), "метки строк одного batch'а обязаны различаться"
|
||||||
|
|
||||||
|
|
||||||
|
# ── 3. Инвариант источника: содержательная колонка не пишется now() ─────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("mod", [kit_base, kit_snapshot], ids=["base", "snapshot_writer"])
|
||||||
|
def test_no_content_timestamp_is_written_with_now(mod: Any) -> None:
|
||||||
|
"""`<колонка> = NOW()` в этих писателях не имеет корректного применения.
|
||||||
|
|
||||||
|
Комментарии вырезаются: `NOW()` там остаётся в объяснительной прозе (и должен —
|
||||||
|
без неё следующий читатель повторит дефект).
|
||||||
|
"""
|
||||||
|
src = _strip_sql_comments(inspect.getsource(mod))
|
||||||
|
cols = "|".join(sorted({c for t in _WATCHED.values() for c in t}))
|
||||||
|
assert re.search(rf"\b({cols})\s*=\s*NOW\s*\(\s*\)", src, re.I) is None
|
||||||
|
assert "statement_timestamp()" in src
|
||||||
|
|
||||||
|
|
||||||
|
# ── 4. Историческая граница зафиксирована в схеме ────────────────────────────
|
||||||
|
|
||||||
|
_SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql"
|
||||||
|
_MIGRATION = _SQL_DIR / "232_listings_observation_time_meaning.sql"
|
||||||
|
|
||||||
|
|
||||||
|
def test_migration_232_marks_the_transition_date_and_unrecoverable_history() -> None:
|
||||||
|
"""Аналитике нужно знать, где сменился смысл колонки и что до него не чинится."""
|
||||||
|
sql = _MIGRATION.read_text("utf-8")
|
||||||
|
assert "BEGIN;" in sql and "COMMIT;" in sql
|
||||||
|
for col in ("listings_snapshots.observed_at", "listings.scraped_at", "listings.last_seen_at"):
|
||||||
|
assert f"COMMENT ON COLUMN {col} IS" in sql
|
||||||
|
assert sql.count("#2731 (2026-08-06)") >= 3, "дата перехода — в каждом комментарии"
|
||||||
|
assert "невосстановим" in sql
|
||||||
|
assert not re.search(r":\w+::", sql), "psycopg v3: никаких :param::type"
|
||||||
|
|
@ -7,7 +7,8 @@
|
||||||
эстиматору было видно лишь ~45.6% активного инвентаря (CIAN — 5.9%).
|
эстиматору было видно лишь ~45.6% активного инвентаря (CIAN — 5.9%).
|
||||||
|
|
||||||
Фикс: и ON CONFLICT DO UPDATE, и dedup-drift reconcile UPDATE теперь выставляют
|
Фикс: и ON CONFLICT DO UPDATE, и dedup-drift reconcile UPDATE теперь выставляют
|
||||||
scraped_at = NOW() рядом с last_seen_at = NOW(). Проверяем оба пути в
|
scraped_at рядом с last_seen_at (с #2731 — statement_timestamp(), см.
|
||||||
|
test_2731_observation_timestamps.py; до него NOW()). Проверяем оба пути в
|
||||||
`scraper_kit.base` (единственный живой модуль — `app.services.scrapers.base`
|
`scraper_kit.base` (единственный живой модуль — `app.services.scrapers.base`
|
||||||
удалён #2397 финальный шаг E, mirror-тесты через legacy убраны) и свойства
|
удалён #2397 финальный шаг E, mirror-тесты через legacy убраны) и свойства
|
||||||
ретро-бэкфилл-миграции 161.
|
ретро-бэкфилл-миграции 161.
|
||||||
|
|
@ -120,7 +121,12 @@ def _kit_matcher() -> MagicMock:
|
||||||
|
|
||||||
|
|
||||||
def test_kit_on_conflict_bumps_scraped_at() -> None:
|
def test_kit_on_conflict_bumps_scraped_at() -> None:
|
||||||
"""kit ON CONFLICT DO UPDATE двигает scraped_at = NOW()."""
|
"""kit ON CONFLICT DO UPDATE двигает scraped_at рядом с last_seen_at.
|
||||||
|
|
||||||
|
#2731: обе метки пишутся statement_timestamp() (построчно, но одинаково внутри
|
||||||
|
statement'а) вместо NOW() — инвариант #2206 «scraped_at двигается вместе с
|
||||||
|
last_seen_at» от этого не меняется, меняется только источник времени.
|
||||||
|
"""
|
||||||
db = _mock_db_update_path(inserted=False)
|
db = _mock_db_update_path(inserted=False)
|
||||||
lot = KitLot(
|
lot = KitLot(
|
||||||
source="cian",
|
source="cian",
|
||||||
|
|
@ -134,12 +140,12 @@ def test_kit_on_conflict_bumps_scraped_at() -> None:
|
||||||
|
|
||||||
assert (inserted, updated) == (0, 1)
|
assert (inserted, updated) == (0, 1)
|
||||||
sql = _find_sql(db, "INSERT INTO listings (")
|
sql = _find_sql(db, "INSERT INTO listings (")
|
||||||
assert "scraped_at = NOW()" in sql
|
assert "scraped_at = statement_timestamp()" in sql
|
||||||
assert "last_seen_at = NOW()" in sql
|
assert "last_seen_at = statement_timestamp()" in sql
|
||||||
|
|
||||||
|
|
||||||
def test_kit_reconcile_bumps_scraped_at() -> None:
|
def test_kit_reconcile_bumps_scraped_at() -> None:
|
||||||
"""kit dedup-drift reconcile UPDATE двигает scraped_at = NOW()."""
|
"""kit dedup-drift reconcile UPDATE двигает scraped_at рядом с last_seen_at (#2731)."""
|
||||||
db = _mock_db_reconcile(reconcile_id=88)
|
db = _mock_db_reconcile(reconcile_id=88)
|
||||||
lot = KitLot(
|
lot = KitLot(
|
||||||
source="avito",
|
source="avito",
|
||||||
|
|
@ -152,8 +158,8 @@ def test_kit_reconcile_bumps_scraped_at() -> None:
|
||||||
kit_save_listings(db, [lot], matcher=_kit_matcher(), region_code=66)
|
kit_save_listings(db, [lot], matcher=_kit_matcher(), region_code=66)
|
||||||
|
|
||||||
sql = _find_sql(db, "SET dedup_hash")
|
sql = _find_sql(db, "SET dedup_hash")
|
||||||
assert "scraped_at = NOW()" in sql
|
assert "scraped_at = statement_timestamp()" in sql
|
||||||
assert "last_seen_at = NOW()" in sql
|
assert "last_seen_at = statement_timestamp()" in sql
|
||||||
|
|
||||||
|
|
||||||
# ── Migration 161: retro-backfill scraped_at ──────────────────────────────────
|
# ── Migration 161: retro-backfill scraped_at ──────────────────────────────────
|
||||||
|
|
|
||||||
|
|
@ -553,7 +553,16 @@ def save_listings(
|
||||||
:predicted_price_min, :predicted_price_max,
|
:predicted_price_min, :predicted_price_max,
|
||||||
:price_trend, :price_previous_rub,
|
:price_trend, :price_previous_rub,
|
||||||
:geo_precision, :card_hash,
|
:geo_precision, :card_hash,
|
||||||
NOW(), NOW()
|
-- scraped_at / last_seen_at — statement_timestamp(), НЕ NOW() (#2731).
|
||||||
|
-- NOW() == transaction_timestamp() замерзает на старте транзакции, а
|
||||||
|
-- save_listings коммитит один раз в конце всего batch'а → все строки
|
||||||
|
-- прогона несли ОДНУ метку (прод 2026-08-06: 219 строк — 1 метка).
|
||||||
|
-- statement_timestamp() двигается построчно (каждый upsert — свой
|
||||||
|
-- statement) и СТАБИЛЕН внутри statement'а, поэтому обе колонки
|
||||||
|
-- получают одно и то же значение. clock_timestamp() здесь нельзя:
|
||||||
|
-- он вычисляется на каждый вызов отдельно и развёл бы scraped_at и
|
||||||
|
-- last_seen_at на микросекунды — расхождение там, где его нет.
|
||||||
|
statement_timestamp(), statement_timestamp()
|
||||||
)
|
)
|
||||||
-- Конфликт-арбитр — dedup_hash (sha256(source|source_id)).
|
-- Конфликт-арбитр — dedup_hash (sha256(source|source_id)).
|
||||||
-- Для одного (source, source_id) формула даёт ОДИН dedup_hash,
|
-- Для одного (source, source_id) формула даёт ОДИН dedup_hash,
|
||||||
|
|
@ -563,10 +572,12 @@ def save_listings(
|
||||||
-- (source, source_id) НЕЛЬЗЯ: yandex/url-only дают source_id=NULL,
|
-- (source, source_id) НЕЛЬЗЯ: yandex/url-only дают source_id=NULL,
|
||||||
-- а NULL не годится как conflict-target → дубли вернулись бы (#1773).
|
-- а NULL не годится как conflict-target → дубли вернулись бы (#1773).
|
||||||
ON CONFLICT (dedup_hash) DO UPDATE
|
ON CONFLICT (dedup_hash) DO UPDATE
|
||||||
SET last_seen_at = NOW(),
|
SET last_seen_at = statement_timestamp(),
|
||||||
-- #2206: ре-подтверждение живым = свежий скрейп; без bump'а
|
-- #2206: ре-подтверждение живым = свежий скрейп; без bump'а
|
||||||
-- эстиматор (scraped_at > NOW()-14d) терял ~55% живого инвентаря.
|
-- эстиматор (scraped_at > NOW()-14d) терял ~55% живого инвентаря.
|
||||||
scraped_at = NOW(),
|
-- #2731: обе метки — statement_timestamp() (см. VALUES выше);
|
||||||
|
-- их равенство сохраняется, потому что оно стабильно в statement'е.
|
||||||
|
scraped_at = statement_timestamp(),
|
||||||
is_active = true,
|
is_active = true,
|
||||||
-- если цена изменилась — обновляем
|
-- если цена изменилась — обновляем
|
||||||
price_rub = EXCLUDED.price_rub,
|
price_rub = EXCLUDED.price_rub,
|
||||||
|
|
@ -679,10 +690,12 @@ def save_listings(
|
||||||
"""
|
"""
|
||||||
UPDATE listings
|
UPDATE listings
|
||||||
SET dedup_hash = :dedup,
|
SET dedup_hash = :dedup,
|
||||||
last_seen_at = NOW(),
|
last_seen_at = statement_timestamp(),
|
||||||
-- #2206: ре-подтверждение живым = свежий скрейп;
|
-- #2206: ре-подтверждение живым = свежий скрейп;
|
||||||
-- без bump'а эстиматор терял ~55% живого инвентаря.
|
-- без bump'а эстиматор терял ~55% живого инвентаря.
|
||||||
scraped_at = NOW(),
|
-- #2731: statement_timestamp() — построчная метка,
|
||||||
|
-- одинаковая для обеих колонок (см. INSERT выше).
|
||||||
|
scraped_at = statement_timestamp(),
|
||||||
is_active = true,
|
is_active = true,
|
||||||
price_rub = :price_rub,
|
price_rub = :price_rub,
|
||||||
price_per_m2 = :ppm2,
|
price_per_m2 = :ppm2,
|
||||||
|
|
|
||||||
|
|
@ -15,6 +15,21 @@
|
||||||
(data-modifying CTE в app/tasks/deactivate_stale_avito.py), не через этот
|
(data-modifying CTE в app/tasks/deactivate_stale_avito.py), не через этот
|
||||||
per-row хелпер.
|
per-row хелпер.
|
||||||
|
|
||||||
|
observed_at пишется `statement_timestamp()`, а НЕ `NOW()` (#2731). `NOW()` в PostgreSQL —
|
||||||
|
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции, а save_listings
|
||||||
|
коммитит ОДИН раз в конце всего batch'а — то есть все снимки прогона получали ОДНУ метку.
|
||||||
|
Прод-замер 2026-08-06: run 3303 — 219 строк, 1 различная метка, при семи минутах работы;
|
||||||
|
run 3229 — 2643 строки на 55 меток (метка на save_listings-вызов, не на строку).
|
||||||
|
Цена — не фильтры свежести (14 суток против минут прогона), а разрешение во времени:
|
||||||
|
темп сбора и порядок внутри прогона по данным не восстановить (отсюда же #2701).
|
||||||
|
|
||||||
|
Почему statement_timestamp(), а не clock_timestamp() как в #2702/#2718: этот INSERT — один
|
||||||
|
statement на строку (вызов в цикле save_listings), поэтому statement_timestamp() двигается
|
||||||
|
ПОСТРОЧНО и при этом СТАБИЛЕН внутри строки. clock_timestamp() дал бы разные значения даже
|
||||||
|
двум колонкам одного statement'а (прод-проверка: `clock_timestamp() = clock_timestamp()` → f,
|
||||||
|
`statement_timestamp() = statement_timestamp()` → t) — это сломало бы равенство
|
||||||
|
listings.scraped_at = last_seen_at, на которое опирается предикат миграции 161.
|
||||||
|
|
||||||
position_in_serp (#2674): параметр УДАЛЁН, колонка осталась и остаётся NULL.
|
position_in_serp (#2674): параметр УДАЛЁН, колонка осталась и остаётся NULL.
|
||||||
Это не оборванная проводка — задуманное отношение НЕВЫРАЗИМО в этой таблице.
|
Это не оборванная проводка — задуманное отношение НЕВЫРАЗИМО в этой таблице.
|
||||||
Позиция — свойство пары (объявление, конкретный прогон выдачи с конкретными
|
Позиция — свойство пары (объявление, конкретный прогон выдачи с конкретными
|
||||||
|
|
@ -87,7 +102,7 @@ def upsert_listing_snapshot(
|
||||||
CAST(:price AS bigint),
|
CAST(:price AS bigint),
|
||||||
CAST(:ppm2 AS int),
|
CAST(:ppm2 AS int),
|
||||||
:status,
|
:status,
|
||||||
NOW()
|
statement_timestamp()
|
||||||
)
|
)
|
||||||
ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET
|
ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET
|
||||||
run_id = COALESCE(EXCLUDED.run_id, listings_snapshots.run_id),
|
run_id = COALESCE(EXCLUDED.run_id, listings_snapshots.run_id),
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue