fix(tradein/matching): listing_sources датируется построчно, а не стартом транзакции (#2731) #2743
2 changed files with 147 additions and 5 deletions
|
|
@ -246,7 +246,23 @@ def _upsert_listing_source(
|
|||
source_url: str | None,
|
||||
source_data: dict | None,
|
||||
) -> None:
|
||||
"""Insert or refresh listing_sources row for this source+ext_id."""
|
||||
"""Insert or refresh listing_sources row for this source+ext_id.
|
||||
|
||||
Отметки времени — statement_timestamp(), НЕ NOW() (#2731). Этот upsert вызывается
|
||||
ПОСТРОЧНО из save_listings (hook `_link_listing_to_house`), а транзакция batch'а
|
||||
коммитится один раз в конце, поэтому NOW() (== transaction_timestamp) давал одну
|
||||
метку на весь вызов: прод-замер 2026-08-06 — 219 строк на 1 метку в 11:00,
|
||||
235/1 в 10:00, 297/1 в 09:00, и так каждый час.
|
||||
|
||||
Чинится вместе с listings.scraped_at/last_seen_at, а не отдельно: сегодня
|
||||
listings.last_seen_at = listing_sources.last_seen_at у 100% пар (2407 из 2407 за
|
||||
сутки) именно потому, что обе колонки берут одну транзакционную метку. Почини
|
||||
только одну — вторая осталась бы замороженной на старте batch'а, и расхождение
|
||||
выросло бы с миллисекунд (честная разница двух записей) до длительности прогона.
|
||||
|
||||
Все три колонки пишутся ОДНИМ statement'ом, поэтому statement_timestamp() даёт им
|
||||
одинаковое значение; clock_timestamp() развёл бы их на микросекунды.
|
||||
"""
|
||||
raw = json.dumps(source_data) if source_data is not None else None
|
||||
db.execute(
|
||||
text("""
|
||||
|
|
@ -257,15 +273,15 @@ def _upsert_listing_source(
|
|||
price_rub, area_m2, floor, rooms_count, raw_payload
|
||||
) VALUES (
|
||||
CAST(:lid AS bigint), :s, :e,
|
||||
CAST(:c AS real), :m, NOW(), NOW(),
|
||||
:url, NOW(),
|
||||
CAST(:c AS real), :m, statement_timestamp(), statement_timestamp(),
|
||||
:url, statement_timestamp(),
|
||||
CAST(:p AS bigint), CAST(:a AS numeric), :fl, :rc,
|
||||
CAST(:raw AS jsonb)
|
||||
)
|
||||
ON CONFLICT (ext_source, ext_id) DO UPDATE SET
|
||||
confidence = GREATEST(EXCLUDED.confidence, listing_sources.confidence),
|
||||
last_seen_at = NOW(),
|
||||
last_scraped_at = NOW(),
|
||||
last_seen_at = statement_timestamp(),
|
||||
last_scraped_at = statement_timestamp(),
|
||||
price_rub = COALESCE(EXCLUDED.price_rub, listing_sources.price_rub),
|
||||
area_m2 = COALESCE(EXCLUDED.area_m2, listing_sources.area_m2),
|
||||
floor = COALESCE(EXCLUDED.floor, listing_sources.floor),
|
||||
|
|
|
|||
|
|
@ -0,0 +1,126 @@
|
|||
"""#2731 (вторая половина): listing_sources датируется построчно, как и listings.
|
||||
|
||||
`_upsert_listing_source` вызывается ПОСТРОЧНО из save_listings (hook
|
||||
`_link_listing_to_house`), а транзакция batch'а коммитится один раз в конце. С `NOW()`
|
||||
(== `transaction_timestamp()`) все строки вызова получали одну метку — прод-замер
|
||||
2026-08-06 по часам: 219 строк / 1 метка, 235 / 1, 297 / 1, 150 / 1 …
|
||||
|
||||
Почему это чинится ВМЕСТЕ с listings, а не отдельно: сегодня
|
||||
`listings.last_seen_at = listing_sources.last_seen_at` у 100% пар (прод: 2407 из 2407 за
|
||||
сутки) — ровно потому, что обе колонки берут одну транзакционную метку. Если починить
|
||||
только listings, вторая колонка осталась бы замороженной на старте batch'а, и расхождение
|
||||
выросло бы с миллисекунд (честная разница двух соседних записей) до длительности прогона.
|
||||
|
||||
Фальсификация: на старом коде (`NOW()`) оба теста дают одну метку на все строки.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import re
|
||||
from typing import Any
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||
|
||||
from app.services.matching.listings import _upsert_listing_source
|
||||
|
||||
# Колонки времени этого писателя. Пишутся ОДНИМ statement'ом → обязаны совпадать.
|
||||
_COLS = ("matched_at", "last_seen_at", "last_scraped_at")
|
||||
_STEP_SEC = 0.6 # три строки → 1.8 с работы одной транзакции
|
||||
|
||||
|
||||
class _FakePg:
|
||||
"""Модель PostgreSQL на две функции времени (см. test_2731_observation_timestamps).
|
||||
|
||||
`now()` замерзает на старте транзакции, `statement_timestamp()` — на старте запроса;
|
||||
`clock_timestamp()` вычисляется на каждый вызов, поэтому три колонки одного
|
||||
statement'а разъехались бы.
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.wall = 0.0
|
||||
self.tx_start: float | None = None
|
||||
self.stmt_start = 0.0
|
||||
self._clock_calls = 0
|
||||
self.rows: list[dict[str, float]] = []
|
||||
|
||||
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 execute(self, stmt: Any, params: Any = None) -> None:
|
||||
if self.tx_start is None:
|
||||
self.tx_start = self.wall
|
||||
self.stmt_start = self.wall
|
||||
self.wall += _STEP_SEC
|
||||
sql = re.sub(r"--[^\n]*", "", str(stmt))
|
||||
if "INSERT INTO listing_sources" not in sql:
|
||||
return None
|
||||
# VALUES: отметки позиционные; ON CONFLICT: именованные. Берём обе формы.
|
||||
cols = re.search(
|
||||
r"INSERT INTO listing_sources\s*\((.*?)\)\s*VALUES\s*\((.*?)\)\s*\n", sql, re.S
|
||||
)
|
||||
row: dict[str, float] = {}
|
||||
if cols is not None:
|
||||
names = [c.strip() for c in cols.group(1).replace("\n", " ").split(",")]
|
||||
vals = [v.strip() for v in cols.group(2).replace("\n", " ").split(",")]
|
||||
for name, val in zip(names, vals, strict=False):
|
||||
f = re.fullmatch(r"(\w+)\s*\(\s*\)", val)
|
||||
if name in _COLS and f is not None:
|
||||
row[name] = self._value(f.group(1).lower())
|
||||
for col in _COLS: # ON CONFLICT DO UPDATE SET <col> = <func>()
|
||||
m = re.search(rf"\b{col}\s*=\s*(\w+)\s*\(\s*\)", sql)
|
||||
if m is not None:
|
||||
row.setdefault(col, self._value(m.group(1).lower()))
|
||||
self.rows.append(row)
|
||||
return None
|
||||
|
||||
|
||||
def _write(db: _FakePg, n: int = 3) -> None:
|
||||
for i in range(n):
|
||||
_upsert_listing_source(
|
||||
db,
|
||||
listing_id=i,
|
||||
ext_source="cian",
|
||||
ext_id=str(i),
|
||||
method="source_link",
|
||||
confidence=1.0,
|
||||
price_rub=5_000_000,
|
||||
area_m2=42.0,
|
||||
floor=3,
|
||||
rooms_count=1,
|
||||
source_url=f"https://ekb.cian.ru/sale/flat/{i}/",
|
||||
source_data=None,
|
||||
)
|
||||
|
||||
|
||||
def test_listing_source_marks_differ_per_row() -> None:
|
||||
"""Писатель, отработавший дольше секунды, оставляет РАЗЛИЧНЫЕ метки у строк."""
|
||||
db = _FakePg()
|
||||
_write(db)
|
||||
|
||||
marks = [row["last_seen_at"] for row in db.rows]
|
||||
assert len(marks) == 3
|
||||
assert max(marks) - min(marks) > 1.0, "писатель отработал дольше секунды"
|
||||
assert len(set(marks)) == len(marks), "у строк одной транзакции обязаны быть разные метки"
|
||||
|
||||
|
||||
def test_all_three_marks_of_one_statement_agree() -> None:
|
||||
"""matched_at / last_seen_at / last_scraped_at пишутся одним statement'ом — и совпадают.
|
||||
|
||||
`clock_timestamp()` развёл бы их на микросекунды: модель БД это воспроизводит,
|
||||
поэтому такая «починка» провалит тест.
|
||||
"""
|
||||
db = _FakePg()
|
||||
_write(db)
|
||||
|
||||
for row in db.rows:
|
||||
assert set(row) == set(_COLS), "все три отметки обязаны быть записаны"
|
||||
assert len(set(row.values())) == 1, "отметки одного statement'а обязаны совпадать"
|
||||
Loading…
Add table
Reference in a new issue