fix(tradein/scrapers): метка наблюдения — время строки, а не старта транзакции (#2731) (#2742)
Some checks failed
Deploy Trade-In / changes (push) Successful in 24s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m41s
Deploy Trade-In / build-backend (push) Successful in 1m56s
Deploy Trade-In / deploy (push) Has been cancelled
Some checks failed
Deploy Trade-In / changes (push) Successful in 24s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m41s
Deploy Trade-In / build-backend (push) Successful in 1m56s
Deploy Trade-In / deploy (push) Has been cancelled
This commit is contained in:
parent
a398b17e6d
commit
91423e0b53
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%).
|
||||
|
||||
Фикс: и 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`
|
||||
удалён #2397 финальный шаг E, mirror-тесты через legacy убраны) и свойства
|
||||
ретро-бэкфилл-миграции 161.
|
||||
|
|
@ -120,7 +121,12 @@ def _kit_matcher() -> MagicMock:
|
|||
|
||||
|
||||
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)
|
||||
lot = KitLot(
|
||||
source="cian",
|
||||
|
|
@ -134,12 +140,12 @@ def test_kit_on_conflict_bumps_scraped_at() -> None:
|
|||
|
||||
assert (inserted, updated) == (0, 1)
|
||||
sql = _find_sql(db, "INSERT INTO listings (")
|
||||
assert "scraped_at = NOW()" in sql
|
||||
assert "last_seen_at = NOW()" in sql
|
||||
assert "scraped_at = statement_timestamp()" in sql
|
||||
assert "last_seen_at = statement_timestamp()" in sql
|
||||
|
||||
|
||||
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)
|
||||
lot = KitLot(
|
||||
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)
|
||||
|
||||
sql = _find_sql(db, "SET dedup_hash")
|
||||
assert "scraped_at = NOW()" in sql
|
||||
assert "last_seen_at = NOW()" in sql
|
||||
assert "scraped_at = statement_timestamp()" in sql
|
||||
assert "last_seen_at = statement_timestamp()" in sql
|
||||
|
||||
|
||||
# ── Migration 161: retro-backfill scraped_at ──────────────────────────────────
|
||||
|
|
|
|||
|
|
@ -553,7 +553,16 @@ def save_listings(
|
|||
:predicted_price_min, :predicted_price_max,
|
||||
:price_trend, :price_previous_rub,
|
||||
: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)).
|
||||
-- Для одного (source, source_id) формула даёт ОДИН dedup_hash,
|
||||
|
|
@ -563,10 +572,12 @@ def save_listings(
|
|||
-- (source, source_id) НЕЛЬЗЯ: yandex/url-only дают source_id=NULL,
|
||||
-- а NULL не годится как conflict-target → дубли вернулись бы (#1773).
|
||||
ON CONFLICT (dedup_hash) DO UPDATE
|
||||
SET last_seen_at = NOW(),
|
||||
SET last_seen_at = statement_timestamp(),
|
||||
-- #2206: ре-подтверждение живым = свежий скрейп; без bump'а
|
||||
-- эстиматор (scraped_at > NOW()-14d) терял ~55% живого инвентаря.
|
||||
scraped_at = NOW(),
|
||||
-- #2731: обе метки — statement_timestamp() (см. VALUES выше);
|
||||
-- их равенство сохраняется, потому что оно стабильно в statement'е.
|
||||
scraped_at = statement_timestamp(),
|
||||
is_active = true,
|
||||
-- если цена изменилась — обновляем
|
||||
price_rub = EXCLUDED.price_rub,
|
||||
|
|
@ -679,10 +690,12 @@ def save_listings(
|
|||
"""
|
||||
UPDATE listings
|
||||
SET dedup_hash = :dedup,
|
||||
last_seen_at = NOW(),
|
||||
last_seen_at = statement_timestamp(),
|
||||
-- #2206: ре-подтверждение живым = свежий скрейп;
|
||||
-- без bump'а эстиматор терял ~55% живого инвентаря.
|
||||
scraped_at = NOW(),
|
||||
-- #2731: statement_timestamp() — построчная метка,
|
||||
-- одинаковая для обеих колонок (см. INSERT выше).
|
||||
scraped_at = statement_timestamp(),
|
||||
is_active = true,
|
||||
price_rub = :price_rub,
|
||||
price_per_m2 = :ppm2,
|
||||
|
|
|
|||
|
|
@ -15,6 +15,21 @@
|
|||
(data-modifying CTE в app/tasks/deactivate_stale_avito.py), не через этот
|
||||
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.
|
||||
Это не оборванная проводка — задуманное отношение НЕВЫРАЗИМО в этой таблице.
|
||||
Позиция — свойство пары (объявление, конкретный прогон выдачи с конкретными
|
||||
|
|
@ -87,7 +102,7 @@ def upsert_listing_snapshot(
|
|||
CAST(:price AS bigint),
|
||||
CAST(:ppm2 AS int),
|
||||
:status,
|
||||
NOW()
|
||||
statement_timestamp()
|
||||
)
|
||||
ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET
|
||||
run_id = COALESCE(EXCLUDED.run_id, listings_snapshots.run_id),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue