fix(tradein): гейт против сдвига разряда в offer_price_history (#3376)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 11s
CI / changes (pull_request) Successful in 14s
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 5m4s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 11s
CI / changes (pull_request) Successful in 14s
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 5m4s
В priceHistory источников встречаются точки ровно ×10/÷10 к соседям с возвратом к базе следующей же точкой — это потерянный разряд у источника, а не рынок. Прод-замер 06.09.2026: 76 таких строк у 55 объявлений (domklik 48, cian 21, yandex 6; единственная avito-строка — триггерная). Читатели колонки — медианный торг лендинга (#3223), админка, /scrapers. drop_decimal_slips живёт рядом с validate_diff_percent — на той же единой границе записи, что и гейт #3225, и подключён ко ВСЕМ писателям истории (domclick/detail.py, cian/detail.py, yandex_price_history.py). Критерий требует двух свидетелей: скачок ×10 к предыдущей точке И возврат к базе у следующей; у последней точки серии свидетель — текущая цена объявления, нет и её → точку не трогаем (без второго свидетеля ×10 может быть честной сменой цены). У domklik после выброса пересчитывается diff_percent соседа: парсер считал его от базы, которой больше нет. Миграция 286 чистит уже собранные строки тем же критерием и только у загрузчика (change_time <> recorded_at): триггерные строки — это живые смены listings.price_rub, у них другая база отсчёта (см. 285). Падает, если кандидатов больше 200.
This commit is contained in:
parent
1cff12cc71
commit
142967d064
7 changed files with 419 additions and 36 deletions
|
|
@ -34,6 +34,7 @@ from datetime import UTC, datetime, timedelta
|
||||||
# isinstance checks — так что любой ScrapedLot-совместимый объект (source, source_id,
|
# isinstance checks — так что любой ScrapedLot-совместимый объект (source, source_id,
|
||||||
# price_rub, price_previous_rub, price_trend) остаётся взаимозаменяем для этой функции.
|
# price_rub, price_previous_rub, price_trend) остаётся взаимозаменяем для этой функции.
|
||||||
from scraper_kit.base import ScrapedLot
|
from scraper_kit.base import ScrapedLot
|
||||||
|
from scraper_kit.offer_price_history import drop_decimal_slips
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
|
@ -116,6 +117,7 @@ def record_yandex_price_history(db: Session, lots: list[ScrapedLot]) -> int:
|
||||||
inserted = 0
|
inserted = 0
|
||||||
skipped = 0
|
skipped = 0
|
||||||
errors = 0
|
errors = 0
|
||||||
|
slips = 0
|
||||||
|
|
||||||
for lot in lots:
|
for lot in lots:
|
||||||
if lot.source != _SOURCE:
|
if lot.source != _SOURCE:
|
||||||
|
|
@ -140,36 +142,47 @@ def record_yandex_price_history(db: Session, lots: list[ScrapedLot]) -> int:
|
||||||
{"listing_id": listing_id},
|
{"listing_id": listing_id},
|
||||||
).fetchone()
|
).fetchone()
|
||||||
|
|
||||||
|
# Серия для гейта сдвигов разряда (#3376): (цена, change_time).
|
||||||
|
# change_time=None — якорь, последняя лежащая в БД точка: гейту нужна
|
||||||
|
# как база отсчёта, вставлять её повторно не надо.
|
||||||
|
points: list[tuple[float, datetime | None]] = []
|
||||||
if latest is None:
|
if latest is None:
|
||||||
# История пуста — seed (+ опционально previous-точка).
|
# История пуста — seed (+ опционально previous-точка).
|
||||||
prev = lot.price_previous_rub
|
prev = lot.price_previous_rub
|
||||||
if prev is not None and prev != lot.price_rub:
|
if prev is not None and prev != lot.price_rub:
|
||||||
_insert_point(
|
points.append((float(prev), prev_ts))
|
||||||
db,
|
points.append((float(lot.price_rub), now))
|
||||||
listing_id=listing_id,
|
else:
|
||||||
price_rub=prev,
|
latest_price = float(latest[0])
|
||||||
change_time=prev_ts,
|
if latest_price == float(lot.price_rub):
|
||||||
)
|
skipped += 1
|
||||||
inserted += 1
|
continue
|
||||||
|
points.append((latest_price, None))
|
||||||
|
points.append((float(lot.price_rub), now))
|
||||||
|
|
||||||
|
# Потолок гейта на этом пути: у yandex последняя точка серии И есть
|
||||||
|
# текущая цена лота, то есть второй свидетель совпадает с самой
|
||||||
|
# точкой и она не выбрасывается никогда. Ловится только сдвиг ВНУТРИ
|
||||||
|
# серии (seed из price.previous + current). Уже лежащие в БД
|
||||||
|
# yandex-сдвиги чистит миграция 286, живой отлов «на следующем
|
||||||
|
# наблюдении» потребовал бы DELETE — отдельная задача.
|
||||||
|
kept, dropped = drop_decimal_slips(
|
||||||
|
points,
|
||||||
|
lambda point: point[0],
|
||||||
|
current_price=lot.price_rub,
|
||||||
|
listing_id=listing_id,
|
||||||
|
)
|
||||||
|
slips += dropped
|
||||||
|
for price_rub, change_time in kept:
|
||||||
|
if change_time is None:
|
||||||
|
continue
|
||||||
_insert_point(
|
_insert_point(
|
||||||
db,
|
db,
|
||||||
listing_id=listing_id,
|
listing_id=listing_id,
|
||||||
price_rub=lot.price_rub,
|
price_rub=price_rub,
|
||||||
change_time=now,
|
change_time=change_time,
|
||||||
)
|
)
|
||||||
inserted += 1
|
inserted += 1
|
||||||
else:
|
|
||||||
latest_price = float(latest[0])
|
|
||||||
if latest_price != float(lot.price_rub):
|
|
||||||
_insert_point(
|
|
||||||
db,
|
|
||||||
listing_id=listing_id,
|
|
||||||
price_rub=lot.price_rub,
|
|
||||||
change_time=now,
|
|
||||||
)
|
|
||||||
inserted += 1
|
|
||||||
else:
|
|
||||||
skipped += 1
|
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
# Per-lot SAVEPOINT откатился — продолжаем батч, не роняем остальные lot'ы.
|
# Per-lot SAVEPOINT откатился — продолжаем батч, не роняем остальные lot'ы.
|
||||||
errors += 1
|
errors += 1
|
||||||
|
|
@ -181,10 +194,11 @@ def record_yandex_price_history(db: Session, lots: list[ScrapedLot]) -> int:
|
||||||
|
|
||||||
db.commit()
|
db.commit()
|
||||||
logger.info(
|
logger.info(
|
||||||
"yandex_price_history: lots=%d inserted=%d skipped=%d errors=%d",
|
"yandex_price_history: lots=%d inserted=%d skipped=%d errors=%d decimal_slips_dropped=%d",
|
||||||
len(lots),
|
len(lots),
|
||||||
inserted,
|
inserted,
|
||||||
skipped,
|
skipped,
|
||||||
errors,
|
errors,
|
||||||
|
slips,
|
||||||
)
|
)
|
||||||
return inserted
|
return inserted
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,155 @@
|
||||||
|
-- 286_offer_price_history_decimal_slips.sql
|
||||||
|
-- Удалить из offer_price_history точки с потерянным разрядом и починить соседа (#3376).
|
||||||
|
--
|
||||||
|
-- ЧТО НЕ ТАК
|
||||||
|
-- В priceHistory источников встречаются точки ровно ×10 / ÷10 к соседям:
|
||||||
|
-- 377 000 → 3 770 000 → 3 720 000; 19 064 000 → 190 150 000 → 19 150 000. Цена
|
||||||
|
-- возвращается к базе следующей же точкой — это потерянный разряд у источника,
|
||||||
|
-- а не рынок. Замер на проде 06.09.2026 (LAG по listing_id, отношение к
|
||||||
|
-- предыдущей точке в [9.8, 10.2] либо [0.098, 0.102]):
|
||||||
|
-- yandex 30 961 строк → 6 вверх, 0 вниз, 6 объявлений (все — строки загрузчика)
|
||||||
|
-- domklik 22 020 строк → 30 вверх, 18 вниз, 35 объявлений (все — строки загрузчика)
|
||||||
|
-- cian 18 266 строк → 9 вверх, 12 вниз, 13 объявлений (все — строки загрузчика)
|
||||||
|
-- avito 6 406 строк → 0 вверх, 1 вниз, 1 объявление (строка ТРИГГЕРА)
|
||||||
|
-- Читатели колонки — медианный торг лендинга (#3223), админка, /scrapers,
|
||||||
|
-- публичный API. Оценщик (services/estimator.py) историю не читает, на оценку
|
||||||
|
-- эти точки не влияли.
|
||||||
|
--
|
||||||
|
-- Код починен в том же PR: гейт drop_decimal_slips в scraper_kit/offer_price_history.py
|
||||||
|
-- подключён ко всем писателям истории (domclick/detail.py, cian/detail.py,
|
||||||
|
-- backend/app/services/yandex_price_history.py). Эта миграция отрабатывает
|
||||||
|
-- задним числом по уже собранным строкам.
|
||||||
|
--
|
||||||
|
-- КРИТЕРИЙ — ТОТ ЖЕ, ЧТО В КОДЕ, И ТРЕБУЕТ ДВУХ СВИДЕТЕЛЕЙ
|
||||||
|
-- 1) скачок к предыдущей точке: price/prev ∈ [9.5, 10.5] ИЛИ prev/price ∈ [9.5, 10.5];
|
||||||
|
-- 2) возврат к базе у следующей точки: next/prev ∈ [0.9, 1.1].
|
||||||
|
-- Окно шире чистой десятки, потому что сдвиг разряда часто идёт вместе с настоящим
|
||||||
|
-- мелким изменением цены (377 000 → 3 720 000 это ×9.87). Одного скачка НЕ хватает:
|
||||||
|
-- ×10 без возврата к базе — это возможное честное изменение цены, и такие точки
|
||||||
|
-- остаются. У последней точки листинга следующей нет — свидетелем берём текущую
|
||||||
|
-- цену объявления (listings.price_rub); нет и её → точку не трогаем.
|
||||||
|
--
|
||||||
|
-- ТОЛЬКО СТРОКИ ЗАГРУЗЧИКА: change_time <> recorded_at
|
||||||
|
-- Признак происхождения тот же, что в 285 (там же его доказательство): триггер
|
||||||
|
-- record_listing_price_change пишет now() и в change_time, и в recorded_at одним
|
||||||
|
-- INSERT'ом → метки равны побайтово; загрузчик кладёт в change_time дату источника,
|
||||||
|
-- а recorded_at оставляет DEFAULT NOW(). Триггерные строки (миграция 131) — это
|
||||||
|
-- зафиксированные ЖИВЫЕ смены listings.price_rub, а не разбор чужого массива; их
|
||||||
|
-- не удаляем и в окно lag()/lead() не берём. Поэтому единственная avito-строка из
|
||||||
|
-- замера остаётся: она триггерная, и если цена в карточке действительно менялась
|
||||||
|
-- ×10, то это наблюдение, а не дефект разбора.
|
||||||
|
--
|
||||||
|
-- ПЕРЕСЧЁТ diff_percent У СЛЕДУЮЩЕЙ СТРОКИ
|
||||||
|
-- У строки, шедшей за удалённой, предыдущая цена сменилась — процент, посчитанный
|
||||||
|
-- от фантомной базы, надо пересчитать по формуле 285: (price − prev)/prev*100,
|
||||||
|
-- prev отсутствует/нулевой → NULL (процента не существует, ноль читался бы как
|
||||||
|
-- «цена не менялась»). |x| > 100 → NULL — ровно то, что делает validate_diff_percent
|
||||||
|
-- на записи: такое значение уже не процент.
|
||||||
|
--
|
||||||
|
-- ИДЕМПОТЕНТНОСТЬ
|
||||||
|
-- Второй прогон: удалённых точек больше нет → выборка пуста → 0 удалений и 0
|
||||||
|
-- пересчётов (UPDATE вдобавок ограничен IS DISTINCT FROM). Новых объектов схемы нет,
|
||||||
|
-- временная таблица уходит по ON COMMIT DROP.
|
||||||
|
--
|
||||||
|
-- СТРАХОВКА
|
||||||
|
-- Замер дал 76 строк-кандидатов. Если критерий поймал больше 200 — он ловит не то
|
||||||
|
-- (или данные разъехались), и миграция обязана упасть, а не молча вычистить историю.
|
||||||
|
BEGIN;
|
||||||
|
-- Конвенция проекта (#2752): DELETE/UPDATE берут блокировки на строках
|
||||||
|
-- offer_price_history, а CREATE TEMP TABLE — DDL; без lock_timeout деплой встанет
|
||||||
|
-- в очередь за чужой сессией и утащит за собой запросы приложения.
|
||||||
|
SET LOCAL lock_timeout = '5s';
|
||||||
|
|
||||||
|
CREATE TEMP TABLE oph_decimal_slips ON COMMIT DROP AS
|
||||||
|
WITH pts AS (
|
||||||
|
SELECT id,
|
||||||
|
listing_id,
|
||||||
|
source,
|
||||||
|
price_rub,
|
||||||
|
lag(price_rub) OVER w AS prev_price,
|
||||||
|
lead(price_rub) OVER w AS next_price,
|
||||||
|
lead(id) OVER w AS next_id
|
||||||
|
FROM offer_price_history
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
WINDOW w AS (PARTITION BY listing_id ORDER BY change_time, id)
|
||||||
|
)
|
||||||
|
SELECT p.id,
|
||||||
|
p.listing_id,
|
||||||
|
p.source,
|
||||||
|
p.next_id,
|
||||||
|
-- Свидетель: следующая точка серии, а если её нет — текущая цена объявления.
|
||||||
|
COALESCE(p.next_price, l.price_rub) AS witness_price
|
||||||
|
FROM pts p
|
||||||
|
LEFT JOIN listings l ON l.id = p.listing_id
|
||||||
|
WHERE p.prev_price > 0
|
||||||
|
AND p.price_rub > 0
|
||||||
|
AND (p.price_rub / p.prev_price BETWEEN 9.5 AND 10.5
|
||||||
|
OR p.prev_price / p.price_rub BETWEEN 9.5 AND 10.5)
|
||||||
|
AND COALESCE(p.next_price, l.price_rub) > 0
|
||||||
|
AND COALESCE(p.next_price, l.price_rub) / p.prev_price BETWEEN 0.9 AND 1.1;
|
||||||
|
|
||||||
|
DO $$
|
||||||
|
DECLARE
|
||||||
|
to_delete bigint;
|
||||||
|
per_source text;
|
||||||
|
deleted_rows bigint;
|
||||||
|
updated_rows bigint;
|
||||||
|
BEGIN
|
||||||
|
SELECT count(*) INTO to_delete FROM oph_decimal_slips;
|
||||||
|
|
||||||
|
SELECT string_agg(source || '=' || cnt, ', ' ORDER BY source)
|
||||||
|
INTO per_source
|
||||||
|
FROM (
|
||||||
|
SELECT source, count(*) AS cnt
|
||||||
|
FROM oph_decimal_slips
|
||||||
|
GROUP BY source
|
||||||
|
) s;
|
||||||
|
|
||||||
|
RAISE NOTICE 'offer_price_history: сдвигов разряда найдено = % (%)',
|
||||||
|
to_delete, COALESCE(per_source, 'ни одного');
|
||||||
|
|
||||||
|
IF to_delete > 200 THEN
|
||||||
|
RAISE EXCEPTION
|
||||||
|
'offer_price_history: к удалению % строк при замере 06.09.2026 = 76 — '
|
||||||
|
'критерий ловит не то, миграция остановлена', to_delete;
|
||||||
|
END IF;
|
||||||
|
|
||||||
|
DELETE FROM offer_price_history oph
|
||||||
|
USING oph_decimal_slips s
|
||||||
|
WHERE oph.id = s.id;
|
||||||
|
GET DIAGNOSTICS deleted_rows = ROW_COUNT;
|
||||||
|
|
||||||
|
-- Пересчёт СТРОГО у строк, шедших за удалёнными: только у них сменился prev.
|
||||||
|
WITH neighbours AS (
|
||||||
|
-- Окно сужено до затронутых листингов: без этого lag() гонится по всей
|
||||||
|
-- таблице, а под lock_timeout = 5s это лишний риск на деплое.
|
||||||
|
SELECT id,
|
||||||
|
price_rub,
|
||||||
|
lag(price_rub) OVER (PARTITION BY listing_id ORDER BY change_time, id)
|
||||||
|
AS prev_price
|
||||||
|
FROM offer_price_history
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
AND listing_id IN (SELECT DISTINCT listing_id FROM oph_decimal_slips)
|
||||||
|
),
|
||||||
|
recomputed AS (
|
||||||
|
SELECT n.id,
|
||||||
|
CASE
|
||||||
|
WHEN n.prev_price > 0
|
||||||
|
AND abs((n.price_rub - n.prev_price) / n.prev_price * 100) <= 100
|
||||||
|
THEN round((n.price_rub - n.prev_price) / n.prev_price * 100, 2)
|
||||||
|
END AS new_diff
|
||||||
|
FROM neighbours n
|
||||||
|
WHERE n.id IN (SELECT next_id FROM oph_decimal_slips WHERE next_id IS NOT NULL)
|
||||||
|
)
|
||||||
|
UPDATE offer_price_history oph
|
||||||
|
SET diff_percent = r.new_diff
|
||||||
|
FROM recomputed r
|
||||||
|
WHERE oph.id = r.id
|
||||||
|
AND oph.diff_percent IS DISTINCT FROM r.new_diff;
|
||||||
|
GET DIAGNOSTICS updated_rows = ROW_COUNT;
|
||||||
|
|
||||||
|
RAISE NOTICE 'offer_price_history: удалено строк = %, пересчитано diff_percent у соседей = %',
|
||||||
|
deleted_rows, updated_rows;
|
||||||
|
END $$;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
|
|
@ -41,7 +41,7 @@ import httpx
|
||||||
import pytest
|
import pytest
|
||||||
from scraper_kit.browser_fetcher import SidecarBanPageError
|
from scraper_kit.browser_fetcher import SidecarBanPageError
|
||||||
from scraper_kit.domclick_exceptions import DomClickBlockedError, DomClickParseError
|
from scraper_kit.domclick_exceptions import DomClickBlockedError, DomClickParseError
|
||||||
from scraper_kit.offer_price_history import validate_diff_percent
|
from scraper_kit.offer_price_history import drop_decimal_slips, validate_diff_percent
|
||||||
from scraper_kit.providers.domclick.detail import (
|
from scraper_kit.providers.domclick.detail import (
|
||||||
DomClickDetailEnrichment,
|
DomClickDetailEnrichment,
|
||||||
_extract_ssr_state,
|
_extract_ssr_state,
|
||||||
|
|
@ -731,6 +731,81 @@ def test_validate_diff_percent_bool_treated_as_none() -> None:
|
||||||
assert validate_diff_percent(False) is None
|
assert validate_diff_percent(False) is None
|
||||||
|
|
||||||
|
|
||||||
|
# ── drop_decimal_slips (#3376) — сдвиг разряда в priceHistory источника ───────
|
||||||
|
# Прод-замер 06.09.2026: 76 точек ровно ×10/÷10 к соседям у 55 объявлений
|
||||||
|
# (domklik 48, cian 21, yandex 6). Вид: 377 000 → 3 770 000 → 3 720 000 — цена
|
||||||
|
# возвращается к базе следующей же точкой.
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_drops_spike_up(caplog: pytest.LogCaptureFixture) -> None:
|
||||||
|
"""Точка ×10 с возвратом к базе выброшена; соседи по бокам целы."""
|
||||||
|
series = [19064000, 190150000, 19150000]
|
||||||
|
with caplog.at_level(logging.WARNING):
|
||||||
|
kept, dropped = drop_decimal_slips(series, lambda p: p, listing_id=406163)
|
||||||
|
assert kept == [19064000, 19150000]
|
||||||
|
assert dropped == 1
|
||||||
|
assert any("406163" in r.getMessage() and "19064000" in r.getMessage() for r in caplog.records)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_drops_spike_down() -> None:
|
||||||
|
"""Тот же критерий в другую сторону: ÷10 и возврат к базе (377 000 при базе 3.7 млн)."""
|
||||||
|
kept, dropped = drop_decimal_slips([3720000, 377000, 3770000], lambda p: p)
|
||||||
|
assert kept == [3720000, 3770000]
|
||||||
|
assert dropped == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_keeps_honest_doubling() -> None:
|
||||||
|
"""Рост ×2 — не сдвиг разряда, серия не тронута."""
|
||||||
|
series = [5000000, 10000000, 9900000]
|
||||||
|
assert drop_decimal_slips(series, lambda p: p) == (series, 0)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_keeps_tenfold_without_return_to_base() -> None:
|
||||||
|
"""×10 БЕЗ возврата к базе — возможное честное изменение цены, точку не трогаем."""
|
||||||
|
series = [1020000, 10200000, 10150000]
|
||||||
|
assert drop_decimal_slips(series, lambda p: p) == (series, 0)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_last_point_checked_against_current_price() -> None:
|
||||||
|
"""У последней точки следующей нет — свидетель это текущая цена объявления."""
|
||||||
|
kept, dropped = drop_decimal_slips([1020000, 10200000], lambda p: p, current_price=1020000)
|
||||||
|
assert kept == [1020000]
|
||||||
|
assert dropped == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_last_point_kept_without_current_price() -> None:
|
||||||
|
"""Второго свидетеля нет вообще → точка остаётся (могла быть честной сменой цены)."""
|
||||||
|
series = [1020000, 10200000]
|
||||||
|
assert drop_decimal_slips(series, lambda p: p, current_price=None) == (series, 0)
|
||||||
|
|
||||||
|
|
||||||
|
def test_save_detail_enrichment_drops_decimal_slip_before_insert() -> None:
|
||||||
|
"""Писатель не отправляет спайк в INSERT, а соседу пересчитывает diff_percent."""
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute.return_value.rowcount = 1
|
||||||
|
# SELECT price_rub FROM listings — текущая цена объявления (свидетель).
|
||||||
|
db.execute.return_value.fetchone.return_value = MagicMock(price_rub=3770000)
|
||||||
|
times = [datetime(2026, 4, d, tzinfo=UTC) for d in (1, 2, 3)]
|
||||||
|
e = DomClickDetailEnrichment(
|
||||||
|
item_id="2075729321",
|
||||||
|
source_url=_CARD_URL,
|
||||||
|
price_changes=[
|
||||||
|
{"change_time": times[0], "price_rub": 3720000, "diff_percent": None},
|
||||||
|
# Сдвиг разряда: ÷10 к предыдущей, следующая возвращается к базе.
|
||||||
|
{"change_time": times[1], "price_rub": 377000, "diff_percent": -89.87},
|
||||||
|
{"change_time": times[2], "price_rub": 3770000, "diff_percent": 900.0},
|
||||||
|
],
|
||||||
|
)
|
||||||
|
|
||||||
|
assert save_detail_enrichment(db, 999, e) is True
|
||||||
|
|
||||||
|
insert_calls = [c for c in db.execute.call_args_list if "offer_price_history" in str(c[0][0])]
|
||||||
|
assert [c[0][1]["price"] for c in insert_calls] == [3720000, 3770000]
|
||||||
|
# У точки-соседа diff_percent парсер считал от исчезнувшей базы 377 000 (+900%,
|
||||||
|
# такое validate_diff_percent отвергает). От настоящей базы 3 720 000 это +1.34%.
|
||||||
|
assert insert_calls[1][0][1]["diff"] == 1.34
|
||||||
|
|
||||||
|
|
||||||
# ── canon_sale_type (#2674) ───────────────────────────────────────────────────
|
# ── canon_sale_type (#2674) ───────────────────────────────────────────────────
|
||||||
# Кейсы — ФАКТИЧЕСКИЙ словарь прода на 2026-08-06:
|
# Кейсы — ФАКТИЧЕСКИЙ словарь прода на 2026-08-06:
|
||||||
# SELECT source, sale_type, count(*) FROM listings GROUP BY 1,2 →
|
# SELECT source, sale_type, count(*) FROM listings GROUP BY 1,2 →
|
||||||
|
|
|
||||||
|
|
@ -60,14 +60,19 @@ class _FakeNested:
|
||||||
|
|
||||||
|
|
||||||
class _FakeSession:
|
class _FakeSession:
|
||||||
"""Пишет все execute(stmt, params); rowcount=1 у каждого оператора."""
|
"""Пишет все execute(stmt, params); rowcount=1 у каждого оператора.
|
||||||
|
|
||||||
|
fetchone() → None: writer читает текущую цену листинга для гейта сдвигов разряда
|
||||||
|
(#3376), а этой фикстуре цена не нужна — «строки нет» гейт трактует как отсутствие
|
||||||
|
второго свидетеля и серию не трогает.
|
||||||
|
"""
|
||||||
|
|
||||||
def __init__(self) -> None:
|
def __init__(self) -> None:
|
||||||
self.calls: list[tuple[str, dict[str, Any]]] = []
|
self.calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
|
||||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||||
self.calls.append((str(stmt), params or {}))
|
self.calls.append((str(stmt), params or {}))
|
||||||
return SimpleNamespace(rowcount=1)
|
return SimpleNamespace(rowcount=1, fetchone=lambda: None)
|
||||||
|
|
||||||
def begin_nested(self) -> _FakeNested:
|
def begin_nested(self) -> _FakeNested:
|
||||||
return _FakeNested()
|
return _FakeNested()
|
||||||
|
|
|
||||||
|
|
@ -19,11 +19,25 @@ loader клал поле источника ``diff``, а там рубли. Пр
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
from collections.abc import Callable, Sequence
|
||||||
|
from typing import TypeVar
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
_DIFF_PERCENT_MAX_ABS = 100.0
|
_DIFF_PERCENT_MAX_ABS = 100.0
|
||||||
|
|
||||||
|
# «Десятичный сдвиг» — потерянный разряд в priceHistory источника (#3376).
|
||||||
|
# Окно вокруг ×10 широкое (9.5..10.5), потому что сдвиг разряда часто идёт вместе
|
||||||
|
# с настоящим мелким изменением цены: 377 000 → 3 720 000 это ×9.87, а не ровно ×10.
|
||||||
|
_SLIP_RATIO_MIN = 9.5
|
||||||
|
_SLIP_RATIO_MAX = 10.5
|
||||||
|
# Возврат к базе у СЛЕДУЮЩЕЙ точки — второй свидетель. Без него ×10 неотличим от
|
||||||
|
# честной смены цены (переезд объявления в другой сегмент, смена лота у продавца).
|
||||||
|
_BASE_RETURN_MIN = 0.9
|
||||||
|
_BASE_RETURN_MAX = 1.1
|
||||||
|
|
||||||
|
T = TypeVar("T")
|
||||||
|
|
||||||
|
|
||||||
def validate_diff_percent(
|
def validate_diff_percent(
|
||||||
value: float | int | None,
|
value: float | int | None,
|
||||||
|
|
@ -51,3 +65,82 @@ def validate_diff_percent(
|
||||||
)
|
)
|
||||||
return None
|
return None
|
||||||
return numeric
|
return numeric
|
||||||
|
|
||||||
|
|
||||||
|
def _positive_price(value: object) -> float | None:
|
||||||
|
"""Цена как положительное число; всё остальное (None, мусор, 0, bool) → None."""
|
||||||
|
if value is None or isinstance(value, bool):
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
numeric = float(value) # type: ignore[arg-type]
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
return None
|
||||||
|
return numeric if numeric > 0 else None
|
||||||
|
|
||||||
|
|
||||||
|
def _is_decimal_slip(prev: float, cur: float, nxt: float | None) -> bool:
|
||||||
|
"""Точка = сдвиг разряда: ×10 (или ÷10) к prev И возврат к базе у next."""
|
||||||
|
ratio = cur / prev
|
||||||
|
shifted = (_SLIP_RATIO_MIN <= ratio <= _SLIP_RATIO_MAX) or (
|
||||||
|
1 / _SLIP_RATIO_MAX <= ratio <= 1 / _SLIP_RATIO_MIN
|
||||||
|
)
|
||||||
|
if not shifted or nxt is None:
|
||||||
|
# nxt is None — второго свидетеля нет, и это НЕ повод выбрасывать: ×10 без
|
||||||
|
# возврата к базе бывает честным изменением цены.
|
||||||
|
return False
|
||||||
|
return _BASE_RETURN_MIN <= nxt / prev <= _BASE_RETURN_MAX
|
||||||
|
|
||||||
|
|
||||||
|
def drop_decimal_slips(
|
||||||
|
points: Sequence[T],
|
||||||
|
price_of: Callable[[T], object],
|
||||||
|
*,
|
||||||
|
current_price: float | int | None = None,
|
||||||
|
listing_id: int | str | None = None,
|
||||||
|
) -> tuple[list[T], int]:
|
||||||
|
"""Убрать из серии точки с потерянным разрядом (#3376).
|
||||||
|
|
||||||
|
Прод-замер 06.09.2026: 76 строк-выбросов у 55 объявлений во всех источниках
|
||||||
|
кроме avito — ``prev=377000 → 3770000 → next=3720000``. Это не рынок, а разряд,
|
||||||
|
потерянный в priceHistory источника: цена возвращается к базе следующей же
|
||||||
|
точкой. Читатели таблицы (медианный торг #3223, админка, /scrapers) получают
|
||||||
|
от такой точки выброс в сотни процентов.
|
||||||
|
|
||||||
|
Критерий требует ДВУХ свидетелей: скачок ×10 к предыдущей точке И возврат к
|
||||||
|
базе у следующей. Одного скачка мало — так выглядит и честная смена цены.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
points: точки одного объявления, упорядоченные по времени (возрастание).
|
||||||
|
price_of: как достать цену из точки (у писателей разная форма точки).
|
||||||
|
current_price: текущая цена объявления — второй свидетель для ПОСЛЕДНЕЙ
|
||||||
|
точки, у которой нет следующей. None → последнюю точку не трогаем.
|
||||||
|
listing_id: только для логов.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
(серия без сдвигов, число выброшенных точек).
|
||||||
|
"""
|
||||||
|
prices = [_positive_price(price_of(point)) for point in points]
|
||||||
|
kept: list[T] = []
|
||||||
|
dropped = 0
|
||||||
|
prev: float | None = None
|
||||||
|
for i, point in enumerate(points):
|
||||||
|
cur = prices[i]
|
||||||
|
# База — предыдущая ОСТАВЛЕННАЯ точка: у двух сдвигов подряд второй иначе
|
||||||
|
# считался бы от выброшенного соседа.
|
||||||
|
if prev is not None and cur is not None:
|
||||||
|
nxt = prices[i + 1] if i + 1 < len(prices) else _positive_price(current_price)
|
||||||
|
if _is_decimal_slip(prev, cur, nxt):
|
||||||
|
dropped += 1
|
||||||
|
logger.warning(
|
||||||
|
"offer_price_history: точка %s выброшена как сдвиг разряда "
|
||||||
|
"(listing_id=%s, prev=%s, next=%s)",
|
||||||
|
cur,
|
||||||
|
listing_id,
|
||||||
|
prev,
|
||||||
|
nxt,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
kept.append(point)
|
||||||
|
if cur is not None:
|
||||||
|
prev = cur
|
||||||
|
return kept, dropped
|
||||||
|
|
|
||||||
|
|
@ -25,7 +25,7 @@ from sqlalchemy.orm import Session
|
||||||
from scraper_kit.ceiling_height import plausible_ceiling_m
|
from scraper_kit.ceiling_height import plausible_ceiling_m
|
||||||
from scraper_kit.cian_exceptions import CianBlockedError
|
from scraper_kit.cian_exceptions import CianBlockedError
|
||||||
from scraper_kit.cian_state_parser import extract_all_states, extract_state
|
from scraper_kit.cian_state_parser import extract_all_states, extract_state
|
||||||
from scraper_kit.offer_price_history import validate_diff_percent
|
from scraper_kit.offer_price_history import drop_decimal_slips, validate_diff_percent
|
||||||
from scraper_kit.providers._base import build_curl_cffi_session
|
from scraper_kit.providers._base import build_curl_cffi_session
|
||||||
from scraper_kit.providers._proxy import curl_proxy_url
|
from scraper_kit.providers._proxy import curl_proxy_url
|
||||||
from scraper_kit.repair_state_normalizer import (
|
from scraper_kit.repair_state_normalizer import (
|
||||||
|
|
@ -467,10 +467,21 @@ def save_detail_enrichment(
|
||||||
exc,
|
exc,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Price changes — INSERT new rows, de-dup by listing_id + change_time
|
# Price changes — INSERT new rows, de-dup by listing_id + change_time.
|
||||||
for change in enrichment.price_changes:
|
# Сдвиги разряда (#3376) выбрасываем до вставки; гейту нужна серия по возрастанию
|
||||||
if not change.get("change_time") or not change.get("price_rub"):
|
# времени, а _extract_price_changes порядок источника не гарантирует. Второй
|
||||||
continue
|
# свидетель для последней точки — цена в listings (_snap_row, читалась выше).
|
||||||
|
_changes = sorted(
|
||||||
|
(c for c in enrichment.price_changes if c.get("change_time") and c.get("price_rub")),
|
||||||
|
key=lambda c: str(c["change_time"]),
|
||||||
|
)
|
||||||
|
_changes, decimal_slips_dropped = drop_decimal_slips(
|
||||||
|
_changes,
|
||||||
|
lambda c: c.get("price_rub"),
|
||||||
|
current_price=_snap_row.price_rub if _snap_row is not None else None,
|
||||||
|
listing_id=listing_id,
|
||||||
|
)
|
||||||
|
for change in _changes:
|
||||||
db.execute(
|
db.execute(
|
||||||
text("""
|
text("""
|
||||||
INSERT INTO offer_price_history
|
INSERT INTO offer_price_history
|
||||||
|
|
@ -553,11 +564,13 @@ def save_detail_enrichment(
|
||||||
|
|
||||||
db.commit()
|
db.commit()
|
||||||
logger.info(
|
logger.info(
|
||||||
"Cian detail enrichment saved for listing_id=%s (price_changes=%d, agent=%s, snapshot=%s)",
|
"Cian detail enrichment saved for listing_id=%s (price_changes=%d, agent=%s, "
|
||||||
|
"snapshot=%s, decimal_slips_dropped=%d)",
|
||||||
listing_id,
|
listing_id,
|
||||||
len(enrichment.price_changes),
|
len(_changes),
|
||||||
bool(enrichment.agent_profile),
|
bool(enrichment.agent_profile),
|
||||||
_snap_written,
|
_snap_written,
|
||||||
|
decimal_slips_dropped,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -62,7 +62,7 @@ from scraper_kit.domclick_exceptions import (
|
||||||
DomClickParseError,
|
DomClickParseError,
|
||||||
)
|
)
|
||||||
from scraper_kit.browser_fetcher import SidecarBanPageError
|
from scraper_kit.browser_fetcher import SidecarBanPageError
|
||||||
from scraper_kit.offer_price_history import validate_diff_percent
|
from scraper_kit.offer_price_history import drop_decimal_slips, validate_diff_percent
|
||||||
from scraper_kit.repair_state_normalizer import (
|
from scraper_kit.repair_state_normalizer import (
|
||||||
infer_repair_state_from_text,
|
infer_repair_state_from_text,
|
||||||
normalize_repair_state,
|
normalize_repair_state,
|
||||||
|
|
@ -792,9 +792,35 @@ def save_detail_enrichment(
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Сдвиги разряда (#3376) — выбрасываем ДО вставки. Серия уже отсортирована по
|
||||||
|
# change_time парсером. Текущая цена листинга — второй свидетель для последней
|
||||||
|
# точки; detail-backfill цену не меняет, так что строка listings актуальна.
|
||||||
|
current_price_row = db.execute(
|
||||||
|
text("SELECT price_rub FROM listings WHERE id = CAST(:lid AS bigint)"),
|
||||||
|
{"lid": listing_id},
|
||||||
|
).fetchone()
|
||||||
|
price_changes, decimal_slips_dropped = drop_decimal_slips(
|
||||||
|
e.price_changes,
|
||||||
|
lambda c: c.get("price_rub"),
|
||||||
|
current_price=current_price_row.price_rub if current_price_row else None,
|
||||||
|
listing_id=listing_id,
|
||||||
|
)
|
||||||
|
if decimal_slips_dropped:
|
||||||
|
# diff_percent парсер считал от соседа, которого больше нет (#3225: процент
|
||||||
|
# берётся из соседних цен, а не из поля источника). Пересчитываем по той же
|
||||||
|
# формуле, иначе у следующей точки останется процент от фантомной базы.
|
||||||
|
prev_price: int | None = None
|
||||||
|
for change in price_changes:
|
||||||
|
change["diff_percent"] = (
|
||||||
|
round((change["price_rub"] - prev_price) / prev_price * 100, 2)
|
||||||
|
if prev_price
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
prev_price = change["price_rub"]
|
||||||
|
|
||||||
# priceHistory — INSERT, де-дуп по (listing_id, change_time). SAVEPOINT per-row
|
# priceHistory — INSERT, де-дуп по (listing_id, change_time). SAVEPOINT per-row
|
||||||
# (backend.md): битый cast одной записи не валит весь UPDATE+commit.
|
# (backend.md): битый cast одной записи не валит весь UPDATE+commit.
|
||||||
for change in e.price_changes:
|
for change in price_changes:
|
||||||
ct = change.get("change_time")
|
ct = change.get("change_time")
|
||||||
price = change.get("price_rub")
|
price = change.get("price_rub")
|
||||||
if not ct or not price:
|
if not ct or not price:
|
||||||
|
|
@ -844,8 +870,10 @@ def save_detail_enrichment(
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
logger.info(
|
logger.info(
|
||||||
"domclick_detail save: listing id=%s enriched (price_changes=%d)",
|
"domclick_detail save: listing id=%s enriched "
|
||||||
|
"(price_changes=%d, decimal_slips_dropped=%d)",
|
||||||
listing_id,
|
listing_id,
|
||||||
len(e.price_changes),
|
len(price_changes),
|
||||||
|
decimal_slips_dropped,
|
||||||
)
|
)
|
||||||
return found
|
return found
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue