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_url: str | None,
|
||||||
source_data: dict | None,
|
source_data: dict | None,
|
||||||
) -> 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
|
raw = json.dumps(source_data) if source_data is not None else None
|
||||||
db.execute(
|
db.execute(
|
||||||
text("""
|
text("""
|
||||||
|
|
@ -257,15 +273,15 @@ def _upsert_listing_source(
|
||||||
price_rub, area_m2, floor, rooms_count, raw_payload
|
price_rub, area_m2, floor, rooms_count, raw_payload
|
||||||
) VALUES (
|
) VALUES (
|
||||||
CAST(:lid AS bigint), :s, :e,
|
CAST(:lid AS bigint), :s, :e,
|
||||||
CAST(:c AS real), :m, NOW(), NOW(),
|
CAST(:c AS real), :m, statement_timestamp(), statement_timestamp(),
|
||||||
:url, NOW(),
|
:url, statement_timestamp(),
|
||||||
CAST(:p AS bigint), CAST(:a AS numeric), :fl, :rc,
|
CAST(:p AS bigint), CAST(:a AS numeric), :fl, :rc,
|
||||||
CAST(:raw AS jsonb)
|
CAST(:raw AS jsonb)
|
||||||
)
|
)
|
||||||
ON CONFLICT (ext_source, ext_id) DO UPDATE SET
|
ON CONFLICT (ext_source, ext_id) DO UPDATE SET
|
||||||
confidence = GREATEST(EXCLUDED.confidence, listing_sources.confidence),
|
confidence = GREATEST(EXCLUDED.confidence, listing_sources.confidence),
|
||||||
last_seen_at = NOW(),
|
last_seen_at = statement_timestamp(),
|
||||||
last_scraped_at = NOW(),
|
last_scraped_at = statement_timestamp(),
|
||||||
price_rub = COALESCE(EXCLUDED.price_rub, listing_sources.price_rub),
|
price_rub = COALESCE(EXCLUDED.price_rub, listing_sources.price_rub),
|
||||||
area_m2 = COALESCE(EXCLUDED.area_m2, listing_sources.area_m2),
|
area_m2 = COALESCE(EXCLUDED.area_m2, listing_sources.area_m2),
|
||||||
floor = COALESCE(EXCLUDED.floor, listing_sources.floor),
|
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