fix(tradein): гейт против сдвига разряда (×10) в offer_price_history у всех писателей + миграция 286 (#3376) #3383
9 changed files with 1119 additions and 16 deletions
|
|
@ -140,6 +140,16 @@ def record_yandex_price_history(db: Session, lots: list[ScrapedLot]) -> int:
|
||||||
{"listing_id": listing_id},
|
{"listing_id": listing_id},
|
||||||
).fetchone()
|
).fetchone()
|
||||||
|
|
||||||
|
# ПОТОЛОК: гейта сдвига разряда (#3376) здесь НЕТ, и он бы тут не
|
||||||
|
# сработал. drop_decimal_slips требует двух свидетелей — скачка ×10
|
||||||
|
# к предыдущей точке и подтверждения у следующей. На этом пути серия
|
||||||
|
# максимум из двух точек (последняя лежащая в БД + текущая), у
|
||||||
|
# последней точки свидетель — текущая цена лота, а она и ЕСТЬ эта
|
||||||
|
# точка: свидетель совпадает с подозреваемым, отношение всегда 1.0.
|
||||||
|
# То есть проводка была бы декорацией: ветка, которая по построению
|
||||||
|
# не может выбросить ни одной точки. Отлов ×10 у yandex требует
|
||||||
|
# другого механизма — сравнения со СЛЕДУЮЩИМ наблюдением, а значит
|
||||||
|
# DELETE уже вставленной строки. Отдельная задача: #3385.
|
||||||
if latest is None:
|
if latest is None:
|
||||||
# История пуста — seed (+ опционально previous-точка).
|
# История пуста — seed (+ опционально previous-точка).
|
||||||
prev = lot.price_previous_rub
|
prev = lot.price_previous_rub
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,318 @@
|
||||||
|
-- 286_offer_price_history_decimal_slips.sql
|
||||||
|
-- Удалить из offer_price_history точки с потерянным разрядом и починить diff_percent (#3376).
|
||||||
|
--
|
||||||
|
-- ЧТО НЕ ТАК
|
||||||
|
-- В priceHistory источников встречаются точки ровно ×10 / ÷10 к соседям:
|
||||||
|
-- 377 000 → 3 770 000 → 3 720 000; 19 064 000 → 190 150 000 → 19 150 000. Цена
|
||||||
|
-- возвращается к базе следующей же точкой — это потерянный разряд у источника,
|
||||||
|
-- а не рынок. Второй вид того же дефекта — в НАЧАЛЕ серии: 330 000 → 3 300 000
|
||||||
|
-- (текущая цена объявления 3 300 000), 420 000 → 4 200 000 → 4 500 000. Дефектная
|
||||||
|
-- точка первая, слева у неё базы нет, поэтому её ловит зеркальное правило.
|
||||||
|
--
|
||||||
|
-- Читатели колонки — медианный торг лендинга (#3223), админка, /scrapers,
|
||||||
|
-- публичный API. Оценщик (services/estimator.py) историю не читает, на оценку
|
||||||
|
-- эти точки не влияли.
|
||||||
|
--
|
||||||
|
-- Код починен в том же PR: гейт drop_decimal_slips в scraper_kit/offer_price_history.py
|
||||||
|
-- подключён к писателям НА ЖИВОМ ТРАКТЕ (domclick/detail.py, cian/detail.py). Эта
|
||||||
|
-- миграция отрабатывает задним числом по уже собранным строкам. Разовый
|
||||||
|
-- scripts/local-cian/playwright_history.py гейта не знает, но переиграть миграцию не
|
||||||
|
-- может: он выбирает только объявления вообще БЕЗ offer_price_history.
|
||||||
|
--
|
||||||
|
-- ВЫБОРКА КАНДИДАТОВ ЗЕРКАЛИТ ГЕЙТ 1:1, А НЕ «ТОТ ЖЕ КРИТЕРИЙ»
|
||||||
|
-- Первая редакция этого файла брала предыдущую точку через lag(), то есть СЫРУЮ
|
||||||
|
-- соседнюю строку — включая те, что удаляются этим же проходом. Гейт в коде берёт
|
||||||
|
-- предыдущую ОСТАВЛЕННУЮ. Расхождение видно на двух сериях (обе — от ревьюера):
|
||||||
|
-- 1M → 10M → 1M → 10M (текущая цена 1M): гейт оставляет [1M, 1M], lag-версия
|
||||||
|
-- оставляла [1M] — честная точка удалена;
|
||||||
|
-- 1M → 10M → 1.05M → 9.9M (текущая цена 1M): lag-версия удаляла на первом проходе
|
||||||
|
-- ещё и 1.05M (её сырая база — уже удалённая 10M), а на втором доедала 9.9M,
|
||||||
|
-- то есть НЕ была идемпотентной.
|
||||||
|
-- Поэтому кандидаты выбираются PL/pgSQL-циклом, который повторяет гейт пошагово:
|
||||||
|
-- строки по (change_time, id), переменная base — последняя ОСТАВЛЕННАЯ цена,
|
||||||
|
-- свидетель — следующая СЫРАЯ точка либо, у последней, listings.price_rub.
|
||||||
|
-- Правило первой точки — вторым, УЖЕ ПО ОСТАВШИМСЯ строкам (в гейте оно тоже
|
||||||
|
-- применяется к kept-серии: свидетели первой точки сами могут оказаться сдвигами).
|
||||||
|
--
|
||||||
|
-- КРИТЕРИЙ — ДВА СВИДЕТЕЛЯ, И У ПЕРВОЙ ТОЧКИ ТОЖЕ ДВА
|
||||||
|
-- Внутри серии: 1) скачок к предыдущей ОСТАВЛЕННОЙ точке, price/base ∈ [9.5, 10.5]
|
||||||
|
-- в любую сторону; 2) возврат к базе у следующей точки, next/base ∈ [0.9, 1.1].
|
||||||
|
-- Первая точка: 1) она ×10/÷10 ко второй оставшейся; 2) вторая подтверждена ТРЕТЬЕЙ
|
||||||
|
-- оставшейся, third/second ∈ [0.9, 1.1]. Третьей нет → кандидата нет.
|
||||||
|
-- Окно шире чистой десятки, потому что сдвиг разряда часто идёт вместе с настоящим
|
||||||
|
-- мелким изменением цены (377 000 → 3 720 000 это ×9.87). Одного скачка НЕ хватает:
|
||||||
|
-- ×10 без подтверждения — это возможное честное изменение цены (400 000 → 4 000 000 →
|
||||||
|
-- 8 000 000 — разгон, а не разряд), и такие точки остаются. Нет второго свидетеля
|
||||||
|
-- вообще (последняя точка серии, а listings.price_rub пуст) → точку не трогаем.
|
||||||
|
--
|
||||||
|
-- ПОЧЕМУ У ПЕРВОЙ ТОЧКИ НЕТ ОТКАТА НА listings.price_rub. Он был и снят сознательно:
|
||||||
|
-- у 17 из 21 кандидата первой точки (yandex 6, domklik 9, cian 2) listings.price_rub
|
||||||
|
-- ТОЧНО равна второй точке — у yandex это буквально одна переменная lot.price_rub,
|
||||||
|
-- записанная в двух местах. «Два свидетеля» там вырождаются в одного, а DELETE
|
||||||
|
-- необратим. Эти 17 строк остаются жить; их процент, если он не процент, занулит
|
||||||
|
-- финальный шаг |diff_percent| > 100 → NULL. Для yandex отлов переехал в #3385 —
|
||||||
|
-- сравнение со СЛЕДУЮЩИМ наблюдением. У ПОСЛЕДНЕЙ точки цена объявления свидетелем
|
||||||
|
-- остаётся: она там сравнивается с базой СЛЕВА, то есть с другим наблюдением.
|
||||||
|
--
|
||||||
|
-- СКОЛЬКО ЖДЁМ УДАЛЕНИЙ
|
||||||
|
-- Замер на проде 06.09.2026 ровно под ЭТИМ критерием: 34 строки внутри серий
|
||||||
|
-- (cian 12, domklik 22) + 4 первых точки по третьей точке (domklik 4) = 38.
|
||||||
|
-- По источникам: cian 12, domklik 26, yandex 0 — yandex целиком отпал вместе с
|
||||||
|
-- откатом на listings.price_rub (см. выше), его случай уехал в #3385.
|
||||||
|
-- Порог остановки 200 — примерно пятикратный запас: поймали больше — критерий
|
||||||
|
-- ловит не то, и миграция обязана упасть, а не молча вычистить историю.
|
||||||
|
--
|
||||||
|
-- ТОЛЬКО СТРОКИ ЗАГРУЗЧИКА: 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, а не разбор чужого массива; их
|
||||||
|
-- не удаляем и в цепочку не берём. Поэтому единственная avito-строка из замера
|
||||||
|
-- остаётся: она триггерная, и если цена в карточке действительно менялась ×10, то
|
||||||
|
-- это наблюдение, а не дефект разбора.
|
||||||
|
--
|
||||||
|
-- ОТСЮДА ЖЕ ТРЕТИЙ СВИДЕТЕЛЬ, КОТОРЫЙ СПАСАЕТ ОТ УДАЛЕНИЯ (L-4): если ровно такая
|
||||||
|
-- же цена есть у ТРИГГЕРНОЙ строки того же объявления — значит listings.price_rub
|
||||||
|
-- когда-то реально равнялась этому значению, и подозреваемая точка загрузчика
|
||||||
|
-- подтверждена независимым писателем. Такую не удаляем и оставляем в цепочке базой.
|
||||||
|
-- ЧЕСТНО ПРО ЕГО ВЕС: на прод-данных 06.09.2026 он не спас НИ ОДНОЙ строки (0 из 21
|
||||||
|
-- кандидата правила первой точки). Механизм проверен только на синтетике — держим
|
||||||
|
-- как страховку на будущих прогонах, а не как замеренную защиту.
|
||||||
|
--
|
||||||
|
-- ПЕРЕСЧЁТ diff_percent У СЛЕДУЮЩЕЙ СТРОКИ
|
||||||
|
-- У строки, шедшей за удалённой, предыдущая цена сменилась — процент, посчитанный
|
||||||
|
-- от фантомной базы, надо пересчитать по формуле 285: (price − prev)/prev*100,
|
||||||
|
-- prev отсутствует/нулевой → NULL (процента не существует, ноль читался бы как
|
||||||
|
-- «цена не менялась»). Цепочки подряд удалённых строк это покрывает: next_id
|
||||||
|
-- последней удалённой в цепочке и есть первая уцелевшая.
|
||||||
|
--
|
||||||
|
-- ФИНАЛЬНЫЙ ШАГ: |diff_percent| > 100 → NULL ПО ВСЕМ ИСТОЧНИКАМ
|
||||||
|
-- Ровно то, что делает validate_diff_percent на записи с #3225: такая величина —
|
||||||
|
-- уже не процент (у domklik туда клали рубли). 285 пересчитала только domklik и
|
||||||
|
-- только из соседних цен, поэтому у cian/yandex/avito-строк загрузчика мусор мог
|
||||||
|
-- остаться. Честные скачки ×2…×6 при этом СОХРАНЯЮТ строку и цену — обнуляется
|
||||||
|
-- только процент, посчитанный не от той величины.
|
||||||
|
--
|
||||||
|
-- ИДЕМПОТЕНТНОСТЬ
|
||||||
|
-- Свойство доказано НА ГЕЙТЕ, а не на этом файле: property-тест
|
||||||
|
-- test_drop_decimal_slips_is_idempotent перебирает серии длины 2-6 по алфавиту
|
||||||
|
-- {1M, 10M, 100M, 1.05M, 9.9M} × current_price ∈ {—, 1M, 10M} и требует
|
||||||
|
-- gate(gate(s)) == gate(s). Раз выборка кандидатов здесь повторяет гейт пошагово,
|
||||||
|
-- свойство переносится: второй прогон даёт 0 удалений. Пересчёт diff_percent
|
||||||
|
-- детерминирован и ограничен IS DISTINCT FROM, финальный NULL-шаг — тоже (после
|
||||||
|
-- него |x| > 100 не остаётся). Новых объектов схемы нет, временные таблицы уходят
|
||||||
|
-- по ON COMMIT DROP.
|
||||||
|
-- Прогон обеих фаз ДВАЖДЫ подряд с ROLLBACK вместо COMMIT — tradein-mvp/scripts/sql/
|
||||||
|
-- 286_dryrun.sql; второй проход обязан дать deleted=0 recomputed=0 nulled=0.
|
||||||
|
BEGIN;
|
||||||
|
-- Конвенция проекта (#2752): DELETE/UPDATE берут блокировки на строках
|
||||||
|
-- offer_price_history, а CREATE TEMP TABLE — DDL; без lock_timeout деплой встанет
|
||||||
|
-- в очередь за чужой сессией и утащит за собой запросы приложения.
|
||||||
|
SET LOCAL lock_timeout = '5s';
|
||||||
|
|
||||||
|
CREATE TEMP TABLE oph_decimal_slips (
|
||||||
|
id bigint PRIMARY KEY,
|
||||||
|
listing_id bigint NOT NULL,
|
||||||
|
next_id bigint, -- строка, шедшая следом: у неё сменится база diff_percent
|
||||||
|
source text NOT NULL
|
||||||
|
) ON COMMIT DROP;
|
||||||
|
|
||||||
|
-- ФАЗА 1 — правило внутри серии, пошагово как в гейте.
|
||||||
|
DO $$
|
||||||
|
DECLARE
|
||||||
|
r record;
|
||||||
|
cur_listing bigint := NULL;
|
||||||
|
base numeric := NULL; -- последняя ОСТАВЛЕННАЯ цена этого объявления
|
||||||
|
witness numeric;
|
||||||
|
confirmed boolean;
|
||||||
|
BEGIN
|
||||||
|
FOR r IN
|
||||||
|
SELECT oph.id,
|
||||||
|
oph.listing_id,
|
||||||
|
oph.source,
|
||||||
|
oph.price_rub,
|
||||||
|
lead(oph.price_rub) OVER w AS next_price,
|
||||||
|
lead(oph.id) OVER w AS next_id,
|
||||||
|
l.price_rub AS listing_price
|
||||||
|
FROM offer_price_history oph
|
||||||
|
LEFT JOIN listings l ON l.id = oph.listing_id
|
||||||
|
WHERE oph.change_time <> oph.recorded_at
|
||||||
|
-- Объявления с одной строкой загрузчика кандидатов дать не могут ни по
|
||||||
|
-- одному из правил (нет ни базы слева, ни второй точки справа).
|
||||||
|
AND oph.listing_id IN (
|
||||||
|
SELECT listing_id
|
||||||
|
FROM offer_price_history
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
GROUP BY listing_id
|
||||||
|
HAVING count(*) >= 2
|
||||||
|
)
|
||||||
|
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
|
||||||
|
ORDER BY oph.listing_id, oph.change_time, oph.id
|
||||||
|
LOOP
|
||||||
|
IF cur_listing IS DISTINCT FROM r.listing_id THEN
|
||||||
|
cur_listing := r.listing_id;
|
||||||
|
base := NULL; -- новое объявление, базы слева нет
|
||||||
|
END IF;
|
||||||
|
|
||||||
|
-- Свидетель справа: следующая точка серии, а у последней — текущая цена
|
||||||
|
-- объявления. Нет ни того, ни другого → witness IS NULL, и сравнение ниже
|
||||||
|
-- даёт NULL, то есть «не сдвиг»: без второго свидетеля точку не трогаем.
|
||||||
|
witness := COALESCE(r.next_price, r.listing_price);
|
||||||
|
|
||||||
|
IF base IS NOT NULL
|
||||||
|
AND base > 0 AND r.price_rub > 0 AND witness > 0
|
||||||
|
AND (r.price_rub / base BETWEEN 9.5 AND 10.5
|
||||||
|
OR base / r.price_rub BETWEEN 9.5 AND 10.5)
|
||||||
|
AND witness / base BETWEEN 0.9 AND 1.1
|
||||||
|
THEN
|
||||||
|
-- L-4: цену подтверждает триггерная строка того же объявления → это
|
||||||
|
-- наблюдавшаяся listings.price_rub, а не дефект разбора.
|
||||||
|
SELECT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM offer_price_history t
|
||||||
|
WHERE t.listing_id = r.listing_id
|
||||||
|
AND t.change_time = t.recorded_at
|
||||||
|
AND t.price_rub = r.price_rub
|
||||||
|
)
|
||||||
|
INTO confirmed;
|
||||||
|
|
||||||
|
IF NOT confirmed THEN
|
||||||
|
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
|
||||||
|
VALUES (r.id, r.listing_id, r.next_id, r.source);
|
||||||
|
CONTINUE; -- база НЕ двигается: этой точки в серии больше нет
|
||||||
|
END IF;
|
||||||
|
END IF;
|
||||||
|
|
||||||
|
IF r.price_rub > 0 THEN
|
||||||
|
base := r.price_rub;
|
||||||
|
END IF;
|
||||||
|
END LOOP;
|
||||||
|
END $$;
|
||||||
|
|
||||||
|
-- ФАЗА 2 — правило первой точки, по ОСТАВШИМСЯ строкам (кандидаты фазы 1 исключены).
|
||||||
|
-- Набор множественный, состояния не требует — обычный INSERT ... SELECT.
|
||||||
|
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
|
||||||
|
WITH kept AS (
|
||||||
|
SELECT oph.id,
|
||||||
|
oph.listing_id,
|
||||||
|
oph.source,
|
||||||
|
oph.price_rub,
|
||||||
|
row_number() OVER w AS rn,
|
||||||
|
lead(oph.id) OVER w AS second_id,
|
||||||
|
lead(oph.price_rub) OVER w AS second_price,
|
||||||
|
lead(oph.price_rub, 2) OVER w AS third_price
|
||||||
|
FROM offer_price_history oph
|
||||||
|
WHERE oph.change_time <> oph.recorded_at
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM oph_decimal_slips s WHERE s.id = oph.id)
|
||||||
|
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
|
||||||
|
)
|
||||||
|
SELECT k.id, k.listing_id, k.second_id, k.source
|
||||||
|
FROM kept k
|
||||||
|
WHERE k.rn = 1
|
||||||
|
AND k.price_rub > 0
|
||||||
|
AND k.second_price > 0
|
||||||
|
AND (k.second_price / k.price_rub BETWEEN 9.5 AND 10.5
|
||||||
|
OR k.price_rub / k.second_price BETWEEN 9.5 AND 10.5)
|
||||||
|
-- Свидетель — только ТРЕТЬЯ точка истории. listings.price_rub здесь НЕ
|
||||||
|
-- подставляется: у 17 из 21 кандидата она в точности равна второй точке,
|
||||||
|
-- то есть это то же наблюдение, а не второй свидетель (см. шапку).
|
||||||
|
AND k.third_price > 0
|
||||||
|
AND k.third_price / k.second_price BETWEEN 0.9 AND 1.1
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM offer_price_history t
|
||||||
|
WHERE t.listing_id = k.listing_id
|
||||||
|
AND t.change_time = t.recorded_at
|
||||||
|
AND t.price_rub = k.price_rub
|
||||||
|
);
|
||||||
|
|
||||||
|
DO $$
|
||||||
|
DECLARE
|
||||||
|
to_delete bigint;
|
||||||
|
per_source text;
|
||||||
|
deleted_rows bigint;
|
||||||
|
updated_rows bigint;
|
||||||
|
nulled_rows bigint;
|
||||||
|
nulled_src text;
|
||||||
|
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 = 38 '
|
||||||
|
'(34 внутри серий + 4 первых точки; cian 12, domklik 26, yandex 0) — '
|
||||||
|
'критерий ловит не то, миграция остановлена', 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;
|
||||||
|
|
||||||
|
-- Финальный шаг: величина не того рода → NULL. Правило validate_diff_percent,
|
||||||
|
-- применённое ко ВСЕМ источникам и только к строкам загрузчика.
|
||||||
|
SELECT string_agg(source || '=' || cnt, ', ' ORDER BY source)
|
||||||
|
INTO nulled_src
|
||||||
|
FROM (
|
||||||
|
SELECT source, count(*) AS cnt
|
||||||
|
FROM offer_price_history
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
AND abs(diff_percent) > 100
|
||||||
|
GROUP BY source
|
||||||
|
) s;
|
||||||
|
|
||||||
|
UPDATE offer_price_history
|
||||||
|
SET diff_percent = NULL
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
AND abs(diff_percent) > 100;
|
||||||
|
GET DIAGNOSTICS nulled_rows = ROW_COUNT;
|
||||||
|
|
||||||
|
RAISE NOTICE 'offer_price_history: |diff_percent| > 100 обнулено = % (%)',
|
||||||
|
nulled_rows, COALESCE(nulled_src, 'ни одной');
|
||||||
|
END $$;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
|
|
@ -32,6 +32,7 @@ importers) — `scraper_kit.domclick_exceptions` остаётся единств
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import itertools
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
from datetime import UTC, datetime, timedelta
|
from datetime import UTC, datetime, timedelta
|
||||||
|
|
@ -41,7 +42,11 @@ 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,
|
||||||
|
recompute_diff_percent,
|
||||||
|
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 +736,254 @@ 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 под ЭТИМ критерием: 38 точек (cian 12, domklik 26,
|
||||||
|
# yandex 0). Вид: 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 БЕЗ возврата к базе — возможное честное изменение цены, точку не трогаем.
|
||||||
|
|
||||||
|
Хвост 1 020 000 → 10 200 000 → 10 150 000 намеренно стоит НЕ первым: ровно та
|
||||||
|
же тройка в начале серии — это уже сдвиг по правилу первой точки (см.
|
||||||
|
test_drop_decimal_slips_drops_first_point_witnessed_by_third), и правила здесь
|
||||||
|
не спорят, а смотрят на разные концы серии. Внутри серии у 10 200 000 есть
|
||||||
|
база слева, и она НЕ подтверждена следующей точкой (9.95 ≠ 1) — точка честная.
|
||||||
|
"""
|
||||||
|
series = [1000000, 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)
|
||||||
|
|
||||||
|
|
||||||
|
# ── правило ПЕРВОЙ точки (#3376, прод-разбор 06.09.2026) ─────────────────────
|
||||||
|
# Серии, у которых дефектная точка первая: 420 000 → 4 200 000 → 4 500 000. Слева
|
||||||
|
# базы нет, основное правило такую точку не видит. Свидетелей по-прежнему двое, и
|
||||||
|
# оба — из самой истории: вторая точка и ТРЕТЬЯ. Кандидатов первой точки на проде
|
||||||
|
# было 21, свидетеля из истории имеют 4 (все domklik).
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_drops_first_point_witnessed_by_third() -> None:
|
||||||
|
"""420 000 → 4 200 000 → 4 500 000: вторую точку подтверждает третья."""
|
||||||
|
kept, dropped = drop_decimal_slips([420000, 4200000, 4500000], lambda p: p)
|
||||||
|
assert (kept, dropped) == ([4200000, 4500000], 1)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_keeps_first_point_when_series_runs_away() -> None:
|
||||||
|
"""Первая ×10 ко второй, но третья ушла ЕЩЁ дальше — это разгон цены, не разряд."""
|
||||||
|
series = [400000, 4000000, 8000000]
|
||||||
|
assert drop_decimal_slips(series, lambda p: p) == (series, 0)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_keeps_first_point_without_third_point() -> None:
|
||||||
|
"""330 000 → 3 300 000: третьей точки нет — свидетеля нет, ДАЖЕ с ценой лота.
|
||||||
|
|
||||||
|
Цена объявления в правиле первой точки не участвует сознательно: у 17 из 21
|
||||||
|
прод-кандидата (yandex 6, domklik 9, cian 2) listings.price_rub в точности
|
||||||
|
равнялась второй точке — у yandex это буквально одна переменная lot.price_rub,
|
||||||
|
записанная в двух местах. Такой «второй свидетель» — то же наблюдение, а
|
||||||
|
удаление точки необратимо. Верни COALESCE на current_price в правило первой
|
||||||
|
точки — первый assert покраснеет.
|
||||||
|
"""
|
||||||
|
series = [330000, 3300000]
|
||||||
|
assert drop_decimal_slips(series, lambda p: p, current_price=3300000) == (series, 0)
|
||||||
|
assert drop_decimal_slips(series, lambda p: p, current_price=None) == (series, 0)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_first_point_decided_on_kept_series() -> None:
|
||||||
|
"""Свидетели первой точки — ОСТАВШИЕСЯ соседи, а не сырые.
|
||||||
|
|
||||||
|
1M → 10M → 100M → 10M: 100M выбрасывает основное правило, и только после этого
|
||||||
|
видно, что первая точка ÷10 к оставшейся серии [10M, 10M]. Считай правило по
|
||||||
|
сырым соседям — первый проход оставил бы 1M, а второй выбросил, то есть гейт
|
||||||
|
перестал бы быть идемпотентным (перебор ниже ловит 1116 таких прогонов, и эта
|
||||||
|
серия — первый из них).
|
||||||
|
"""
|
||||||
|
kept, dropped = drop_decimal_slips(
|
||||||
|
[1_000_000, 10_000_000, 100_000_000, 10_000_000], lambda p: p
|
||||||
|
)
|
||||||
|
assert (kept, dropped) == ([10_000_000, 10_000_000], 2)
|
||||||
|
|
||||||
|
|
||||||
|
# ── база = предыдущая ОСТАВЛЕННАЯ точка (фальсификаторы ревьюера) ────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_alternating_series_keeps_honest_points() -> None:
|
||||||
|
"""1M → 10M → 1M → 10M при цене 1M: обе честные точки на месте.
|
||||||
|
|
||||||
|
Первая редакция миграции 286 брала базой СЫРУЮ предыдущую строку (lag) и на
|
||||||
|
этой серии оставляла [1M] вместо [1M, 1M] — удаляла честную точку.
|
||||||
|
"""
|
||||||
|
kept, dropped = drop_decimal_slips(
|
||||||
|
[1_000_000, 10_000_000, 1_000_000, 10_000_000], lambda p: p, current_price=1_000_000
|
||||||
|
)
|
||||||
|
assert (kept, dropped) == ([1_000_000, 1_000_000], 2)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_base_is_previous_kept_not_raw() -> None:
|
||||||
|
"""1M → 10M → 1.05M → 9.9M при цене 1M: после выброса 10M база для 1.05M — это 1M.
|
||||||
|
|
||||||
|
По сырой базе 1.05M тоже стала бы кандидатом (10M/1.05M = 9.52), а следующий
|
||||||
|
прогон доел бы 9.9M — ровно та неидемпотентность, которую нашёл ревьюер.
|
||||||
|
"""
|
||||||
|
kept, dropped = drop_decimal_slips(
|
||||||
|
[1_000_000, 10_000_000, 1_050_000, 9_900_000], lambda p: p, current_price=1_000_000
|
||||||
|
)
|
||||||
|
assert (kept, dropped) == ([1_000_000, 1_050_000, 9_900_000], 1)
|
||||||
|
|
||||||
|
|
||||||
|
# ── свойство: gate(gate(s)) == gate(s) ───────────────────────────────────────
|
||||||
|
|
||||||
|
_SLIP_ALPHABET = (1_000_000, 10_000_000, 100_000_000, 1_050_000, 9_900_000)
|
||||||
|
_SLIP_CURRENTS = (None, 1_000_000, 10_000_000)
|
||||||
|
|
||||||
|
|
||||||
|
def test_drop_decimal_slips_is_idempotent(caplog: pytest.LogCaptureFixture) -> None:
|
||||||
|
"""Повторный прогон гейта не трогает ничего — на всех сериях длины 2-6.
|
||||||
|
|
||||||
|
Свойство несущее, а не декоративное: миграция 286 зеркалит эту функцию
|
||||||
|
пошагово, а доказать идемпотентность на самом SQL-файле нечем — прод-прогон
|
||||||
|
один. Раз выборка кандидатов повторяет гейт, свойство переносится на неё.
|
||||||
|
|
||||||
|
ФАЛЬСИФИЦИРУЕМОСТЬ (перемерено 06.09.2026 на этом же переборе — 19 525 серий
|
||||||
|
× 3 цены = 58 575 прогонов):
|
||||||
|
• база = сырая предыдущая точка вместо оставленной → 136 неидемпотентных
|
||||||
|
прогонов, первый же — (1M, 10M, 1.05M, 9.9M) при цене 1M: [1M, 9.9M] → [1M];
|
||||||
|
• решение по первой точке на сырых соседях вместо оставшихся → 1116 прогонов,
|
||||||
|
первый — (1M, 10M, 100M, 10M) без текущей цены: [1M, 10M, 10M] → [10M, 10M].
|
||||||
|
"""
|
||||||
|
# Гейт логирует каждый выброс; на таком переборе это десятки тысяч записей.
|
||||||
|
caplog.set_level(logging.CRITICAL, logger="scraper_kit.offer_price_history")
|
||||||
|
for length in range(2, 7):
|
||||||
|
for series in itertools.product(_SLIP_ALPHABET, repeat=length):
|
||||||
|
for current in _SLIP_CURRENTS:
|
||||||
|
once, _ = drop_decimal_slips(list(series), lambda p: p, current_price=current)
|
||||||
|
twice, dropped_again = drop_decimal_slips(once, lambda p: p, current_price=current)
|
||||||
|
assert (twice, dropped_again) == (once, 0), (
|
||||||
|
f"серия {series} при текущей цене {current}: {once} → {twice}"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ── recompute_diff_percent — общий пересчёт после выброса ────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_recompute_diff_percent_from_surviving_neighbours() -> None:
|
||||||
|
"""Процент считается от оставшейся базы; у первой точки базы нет → NULL."""
|
||||||
|
changes: list[dict[str, object]] = [
|
||||||
|
{"price_rub": 3720000, "diff_percent": -89.87},
|
||||||
|
{"price_rub": 3770000, "diff_percent": 900.0},
|
||||||
|
]
|
||||||
|
recompute_diff_percent(changes)
|
||||||
|
assert [c["diff_percent"] for c in changes] == [None, 1.34]
|
||||||
|
|
||||||
|
|
||||||
|
def test_recompute_diff_percent_survives_none_price() -> None:
|
||||||
|
"""price_rub=None не роняет пересчёт: ручной ingest кладёт price_changes сырьём.
|
||||||
|
|
||||||
|
scripts/ingest_domclick_jsonl.py:108 берёт rec["price_changes"] из JSONL без
|
||||||
|
валидации — арифметика по None здесь дала бы TypeError на всю запись.
|
||||||
|
"""
|
||||||
|
changes: list[dict[str, object]] = [
|
||||||
|
{"price_rub": 5_000_000, "diff_percent": 12.0},
|
||||||
|
{"price_rub": None, "diff_percent": 900.0},
|
||||||
|
{"price_rub": 5_100_000, "diff_percent": -90.0},
|
||||||
|
]
|
||||||
|
recompute_diff_percent(changes)
|
||||||
|
assert [c["diff_percent"] for c in changes] == [None, None, 2.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
|
||||||
|
|
||||||
|
|
||||||
|
def test_save_detail_enrichment_witness_comes_only_from_listing_price() -> None:
|
||||||
|
"""Серия из ДВУХ точек: второго свидетеля даёт исключительно цена листинга.
|
||||||
|
|
||||||
|
Здесь чтение цены не декорация, а единственный вход гейта: у последней точки
|
||||||
|
следующей нет. Сломай его (RETURNING убран / fetchone → None) — свидетеля не
|
||||||
|
станет, спайк 10 200 000 уедет в offer_price_history, и тест покраснеет на
|
||||||
|
втором INSERT'е. Проверено подменой fetchone → None.
|
||||||
|
"""
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute.return_value.rowcount = 1
|
||||||
|
db.execute.return_value.fetchone.return_value = MagicMock(price_rub=1020000)
|
||||||
|
times = [datetime(2026, 4, d, tzinfo=UTC) for d in (1, 2)]
|
||||||
|
e = DomClickDetailEnrichment(
|
||||||
|
item_id="x",
|
||||||
|
source_url=_CARD_URL,
|
||||||
|
price_changes=[
|
||||||
|
{"change_time": times[0], "price_rub": 1020000, "diff_percent": None},
|
||||||
|
{"change_time": times[1], "price_rub": 10200000, "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] == [1020000]
|
||||||
|
# Цену отдаёт тот же UPDATE — отдельного SELECT по PK больше нет.
|
||||||
|
assert "RETURNING price_rub" in str(db.execute.call_args_list[0][0][0])
|
||||||
|
assert not [c for c in db.execute.call_args_list if "SELECT price_rub" in str(c[0][0])]
|
||||||
|
|
||||||
|
|
||||||
# ── 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,20 @@ class _FakeNested:
|
||||||
|
|
||||||
|
|
||||||
class _FakeSession:
|
class _FakeSession:
|
||||||
"""Пишет все execute(stmt, params); rowcount=1 у каждого оператора."""
|
"""Пишет все execute(stmt, params); rowcount=1 у каждого оператора.
|
||||||
|
|
||||||
|
fetchone() → None: UPDATE listings отдаёт текущую цену через RETURNING — она нужна
|
||||||
|
гейту сдвигов разряда (#3376) как второй свидетель. Этой фикстуре цена не нужна,
|
||||||
|
«строки нет» гейт трактует как отсутствие свидетеля и серию не трогает; факт
|
||||||
|
«листинг найден» здесь по-прежнему приходит из rowcount, а не из RETURNING.
|
||||||
|
"""
|
||||||
|
|
||||||
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()
|
||||||
|
|
|
||||||
|
|
@ -308,6 +308,34 @@ def test_save_detail_enrichment_oph_on_conflict_constraint():
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_save_detail_enrichment_recomputes_diff_after_decimal_slip():
|
||||||
|
"""cian: после выброса сдвига разряда соседу пересчитывается diff_percent (#3376).
|
||||||
|
|
||||||
|
377 000 → 3 770 000 → 377 000 при цене листинга 377 000: средняя точка ×10 с
|
||||||
|
возвратом к базе — выбрасывается. У последней точки парсер посчитал −90% от
|
||||||
|
исчезнувшей базы 3 770 000; от настоящей базы 377 000 это 0%. Без пересчёта в
|
||||||
|
колонку уехал бы процент от цены, которой в истории больше нет (у domclick это
|
||||||
|
уже чинилось, cian отставал).
|
||||||
|
"""
|
||||||
|
db = _mock_db_detail(price_rub=377_000)
|
||||||
|
enrichment = DetailEnrichment(
|
||||||
|
price_changes=[
|
||||||
|
{"change_time": "2026-04-01T00:00:00Z", "price_rub": 377_000, "diff_percent": None},
|
||||||
|
{"change_time": "2026-04-05T00:00:00Z", "price_rub": 3_770_000, "diff_percent": 900.0},
|
||||||
|
{"change_time": "2026-04-10T00:00:00Z", "price_rub": 377_000, "diff_percent": -90.0},
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
save_detail_enrichment(db, 23, enrichment)
|
||||||
|
|
||||||
|
oph = [
|
||||||
|
(params["price"], params["diff"])
|
||||||
|
for sql, params, _ in db._execute_log
|
||||||
|
if "offer_price_history" in sql
|
||||||
|
]
|
||||||
|
assert oph == [(377_000, None), (377_000, 0.0)]
|
||||||
|
|
||||||
|
|
||||||
def test_save_detail_enrichment_skips_price_change_without_change_time():
|
def test_save_detail_enrichment_skips_price_change_without_change_time():
|
||||||
"""price_changes без change_time пропускаются, не вызывают INSERT."""
|
"""price_changes без change_time пропускаются, не вызывают INSERT."""
|
||||||
db = _mock_db_detail(price_rub=5_000_000)
|
db = _mock_db_detail(price_rub=5_000_000)
|
||||||
|
|
|
||||||
|
|
@ -14,16 +14,38 @@ loader клал поле источника ``diff``, а там рубли. Пр
|
||||||
Цена решения: настоящий рост цены больше чем в 2 раза (+100%) тоже уйдёт в NULL.
|
Цена решения: настоящий рост цены больше чем в 2 раза (+100%) тоже уйдёт в NULL.
|
||||||
Это осознанный размен — колонка аналитическая, а молчаливый мусор в ней дороже
|
Это осознанный размен — колонка аналитическая, а молчаливый мусор в ней дороже
|
||||||
пропущенного выброса.
|
пропущенного выброса.
|
||||||
|
|
||||||
|
Здесь же живёт второй гейт той же границы записи — drop_decimal_slips (#3376):
|
||||||
|
он выбрасывает из СЕРИИ точки с потерянным разрядом источника (×10/÷10 с
|
||||||
|
подтверждением у соседей), а recompute_diff_percent чинит процент у соседа
|
||||||
|
выброшенной точки. Оба вызываются писателями на живом тракте (domclick/detail.py,
|
||||||
|
cian/detail.py) и зеркалятся миграцией 286. Загейчены именно они: разовый
|
||||||
|
scripts/local-cian/playwright_history.py пишет мимо гейта, но переиграть миграцию
|
||||||
|
не может — он выбирает только объявления вообще БЕЗ offer_price_history.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
from collections.abc import Callable, MutableMapping, 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 +73,165 @@ 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 _shifted_by_decade(a: float, b: float) -> bool:
|
||||||
|
"""b отличается от a ровно на разряд — ×10 или ÷10 (окно 9.5..10.5)."""
|
||||||
|
return (_SLIP_RATIO_MIN <= b / a <= _SLIP_RATIO_MAX) or (
|
||||||
|
_SLIP_RATIO_MIN <= a / b <= _SLIP_RATIO_MAX
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _is_decimal_slip(prev: float, cur: float, nxt: float | None) -> bool:
|
||||||
|
"""Точка = сдвиг разряда: ×10 (или ÷10) к prev И возврат к базе у next."""
|
||||||
|
if not _shifted_by_decade(prev, cur) or nxt is None:
|
||||||
|
# nxt is None — второго свидетеля нет, и это НЕ повод выбрасывать: ×10 без
|
||||||
|
# возврата к базе бывает честным изменением цены.
|
||||||
|
return False
|
||||||
|
return _BASE_RETURN_MIN <= nxt / prev <= _BASE_RETURN_MAX
|
||||||
|
|
||||||
|
|
||||||
|
def _is_first_point_slip(first: float, second: float, witness: float | None) -> bool:
|
||||||
|
"""Первая точка серии = сдвиг разряда: у неё нет prev, свидетели справа.
|
||||||
|
|
||||||
|
Прод-разбор 06.09.2026: серии вида ``420 000 → 4 200 000 → 4 500 000``.
|
||||||
|
Дефектная точка ПЕРВАЯ, основное правило её не видит: базы слева нет.
|
||||||
|
Свидетелей по-прежнему двое — ×10 ко ВТОРОЙ точке и подтверждение второй точки
|
||||||
|
ТРЕТЬЕЙ. Одного скачка мало: ``400 000 → 4 000 000 → 8 000 000`` это разгон
|
||||||
|
цены, а не потерянный разряд.
|
||||||
|
|
||||||
|
Свидетель здесь только третья точка истории; текущей цены объявления в этом
|
||||||
|
правиле нет СОЗНАТЕЛЬНО. У 17 из 21 прод-кандидата (yandex 6, domklik 9,
|
||||||
|
cian 2) ``listings.price_rub`` точно равнялась второй точке — у yandex это
|
||||||
|
буквально одна переменная ``lot.price_rub``, записанная в двух местах. Такой
|
||||||
|
«второй свидетель» — то же самое наблюдение, и два свидетеля вырождаются в
|
||||||
|
одного, а выброс точки необратим. У ПОСЛЕДНЕЙ точки (_is_decimal_slip) цена
|
||||||
|
объявления свидетелем остаётся: там она сравнивается с базой СЛЕВА, то есть с
|
||||||
|
другим наблюдением.
|
||||||
|
"""
|
||||||
|
if not _shifted_by_decade(first, second) or witness is None:
|
||||||
|
return False
|
||||||
|
return _BASE_RETURN_MIN <= witness / second <= _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).
|
||||||
|
|
||||||
|
Вид дефекта: ``prev=377000 → 3770000 → next=3720000``. Это не рынок, а разряд,
|
||||||
|
потерянный в priceHistory источника: цена возвращается к базе следующей же
|
||||||
|
точкой. Читатели таблицы (медианный торг #3223, админка, /scrapers) получают
|
||||||
|
от такой точки выброс в сотни процентов. Прод-замер 06.09.2026 под ЭТИМ
|
||||||
|
критерием: 38 строк (cian 12, domklik 26, yandex 0).
|
||||||
|
|
||||||
|
Критерий требует ДВУХ свидетелей: скачок ×10 к предыдущей точке И возврат к
|
||||||
|
базе у следующей. Одного скачка мало — так выглядит и честная смена цены.
|
||||||
|
У ПЕРВОЙ точки предыдущей нет, для неё работает зеркальное правило
|
||||||
|
(_is_first_point_slip) — оба свидетеля справа, и оба из самой истории.
|
||||||
|
|
||||||
|
ИДЕМПОТЕНТНОСТЬ — свойство, а не пожелание: ``gate(gate(s)) == gate(s)``.
|
||||||
|
Его держит property-тест перебором серий длины 2-6 (см.
|
||||||
|
test_drop_decimal_slips_is_idempotent), и от него же зависит миграция 286,
|
||||||
|
которая эту функцию зеркалит: миграция обязана быть безопасной при повторном
|
||||||
|
прогоне. Два места, где свойство легко потерять:
|
||||||
|
• база — предыдущая ОСТАВЛЕННАЯ точка, а не предыдущая сырая. Иначе на серии
|
||||||
|
1M → 10M → 1.05M → 9.9M первый проход выбросит и 1.05M (её база — уже
|
||||||
|
выброшенная 10M), а второй доест 9.9M.
|
||||||
|
• решение по ПЕРВОЙ точке принимается по ОСТАВШЕЙСЯ серии, а не по сырой:
|
||||||
|
свидетели первой точки (вторая и третья точки) сами могут оказаться
|
||||||
|
сдвигами. На сырой серии 1M → 10M → 100M → 10M первый проход оставил бы
|
||||||
|
первую точку (третья сырая точка 100M второй не подтверждает), а второй —
|
||||||
|
выбросил, потому что 100M к тому моменту уже удалена.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
points: точки одного объявления, упорядоченные по времени (возрастание).
|
||||||
|
price_of: как достать цену из точки (у писателей разная форма точки).
|
||||||
|
current_price: текущая цена объявления — второй свидетель ТОЛЬКО для
|
||||||
|
ПОСЛЕДНЕЙ точки, у которой нет следующей. None → последнюю точку не
|
||||||
|
трогаем. В правиле первой точки не участвует (_is_first_point_slip).
|
||||||
|
listing_id: только для логов.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
(серия без сдвигов, число выброшенных точек).
|
||||||
|
"""
|
||||||
|
prices = [_positive_price(price_of(point)) for point in points]
|
||||||
|
now_price = _positive_price(current_price)
|
||||||
|
kept: list[T] = []
|
||||||
|
kept_prices: list[float | None] = []
|
||||||
|
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 now_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)
|
||||||
|
kept_prices.append(cur)
|
||||||
|
if cur is not None:
|
||||||
|
prev = cur
|
||||||
|
|
||||||
|
# Первая точка — последней: её свидетели это уже ОСТАВШИЕСЯ вторая и третья
|
||||||
|
# точки (см. про идемпотентность выше). Третьей нет → свидетеля нет: current_price
|
||||||
|
# здесь не подставляем, она слишком часто копия второй точки (см.
|
||||||
|
# _is_first_point_slip).
|
||||||
|
if len(kept_prices) >= 2 and kept_prices[0] is not None and kept_prices[1] is not None:
|
||||||
|
witness = kept_prices[2] if len(kept_prices) >= 3 else None
|
||||||
|
if _is_first_point_slip(kept_prices[0], kept_prices[1], witness):
|
||||||
|
logger.warning(
|
||||||
|
"offer_price_history: первая точка %s выброшена как сдвиг разряда "
|
||||||
|
"(listing_id=%s, next=%s, witness=%s)",
|
||||||
|
kept_prices[0],
|
||||||
|
listing_id,
|
||||||
|
kept_prices[1],
|
||||||
|
witness,
|
||||||
|
)
|
||||||
|
return kept[1:], dropped + 1
|
||||||
|
return kept, dropped
|
||||||
|
|
||||||
|
|
||||||
|
def recompute_diff_percent(changes: Sequence[MutableMapping[str, object]]) -> None:
|
||||||
|
"""Пересчитать diff_percent из соседних цен ПО МЕСТУ — после выброса точек.
|
||||||
|
|
||||||
|
Формула та же, что у починенного парсера (#3225) и у миграции 285:
|
||||||
|
(price − prev)/prev*100, prev нет → NULL. Ноль здесь читался бы как «цена не
|
||||||
|
менялась», то есть как измерение, которого не было.
|
||||||
|
|
||||||
|
Точки с непригодной ценой (None из ручного ingest, 0, мусор) базой не
|
||||||
|
становятся и получают NULL: писатели их всё равно не вставляют.
|
||||||
|
"""
|
||||||
|
prev: float | None = None
|
||||||
|
for change in changes:
|
||||||
|
price = _positive_price(change.get("price_rub"))
|
||||||
|
if price is None:
|
||||||
|
change["diff_percent"] = None
|
||||||
|
continue
|
||||||
|
change["diff_percent"] = round((price - prev) / prev * 100, 2) if prev else None
|
||||||
|
prev = price
|
||||||
|
|
|
||||||
|
|
@ -25,7 +25,11 @@ 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,
|
||||||
|
recompute_diff_percent,
|
||||||
|
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 +471,30 @@ 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, читалась выше).
|
||||||
|
# ponytail: сортировка лексикографическая по str(change_time). На наблюдаемом
|
||||||
|
# формате cian (`2026-06-11T08:23:11.950626Z`, всегда одна зона) она совпадает с
|
||||||
|
# хронологией; на смешанных форматах/зонах разойдётся с миграцией 286, которая
|
||||||
|
# сортирует timestamptz. Апгрейд — парсить в datetime, если формат поедет.
|
||||||
|
_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,
|
||||||
|
)
|
||||||
|
if decimal_slips_dropped:
|
||||||
|
# То же, что у domclick: diff_percent соседа считался от базы, которой
|
||||||
|
# больше нет. Без пересчёта у него остаётся процент от фантомной цены
|
||||||
|
# (377 000 → 3 770 000 → 377 000 оставлял бы у последней точки −90%).
|
||||||
|
recompute_diff_percent(_changes)
|
||||||
|
for change in _changes:
|
||||||
db.execute(
|
db.execute(
|
||||||
text("""
|
text("""
|
||||||
INSERT INTO offer_price_history
|
INSERT INTO offer_price_history
|
||||||
|
|
@ -553,11 +577,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,11 @@ 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,
|
||||||
|
recompute_diff_percent,
|
||||||
|
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,
|
||||||
|
|
@ -774,6 +778,7 @@ def save_detail_enrichment(
|
||||||
raw_payload = COALESCE(raw_payload, '{}'::jsonb) || CAST(:raw_extra AS jsonb),
|
raw_payload = COALESCE(raw_payload, '{}'::jsonb) || CAST(:raw_extra AS jsonb),
|
||||||
detail_enriched_at = NOW()
|
detail_enriched_at = NOW()
|
||||||
WHERE id = CAST(:lid AS bigint)
|
WHERE id = CAST(:lid AS bigint)
|
||||||
|
RETURNING price_rub
|
||||||
"""),
|
"""),
|
||||||
{
|
{
|
||||||
"lid": listing_id,
|
"lid": listing_id,
|
||||||
|
|
@ -791,10 +796,32 @@ def save_detail_enrichment(
|
||||||
"raw_extra": raw_extra_json,
|
"raw_extra": raw_extra_json,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
# rowcount читаем ДО fetchone(): у RETURNING-результата курсор после выборки
|
||||||
|
# уже мягко закрыт, и порядок «сначала счётчик, потом строка» снимает вопрос,
|
||||||
|
# доживёт ли rowcount до этого момента у конкретного драйвера.
|
||||||
|
found = result.rowcount > 0
|
||||||
|
# Текущая цена листинга — второй свидетель для последней точки priceHistory
|
||||||
|
# (#3376). Берём RETURNING'ом у UPDATE выше, а не отдельным SELECT по PK:
|
||||||
|
# это та же строка, только что обновлённая, лишний round-trip не нужен.
|
||||||
|
listing_row = result.fetchone()
|
||||||
|
|
||||||
|
# Сдвиги разряда (#3376) — выбрасываем ДО вставки. Серия уже отсортирована по
|
||||||
|
# change_time парсером.
|
||||||
|
price_changes, decimal_slips_dropped = drop_decimal_slips(
|
||||||
|
e.price_changes,
|
||||||
|
lambda c: c.get("price_rub"),
|
||||||
|
current_price=listing_row.price_rub if listing_row else None,
|
||||||
|
listing_id=listing_id,
|
||||||
|
)
|
||||||
|
if decimal_slips_dropped:
|
||||||
|
# diff_percent парсер считал от соседа, которого больше нет (#3225: процент
|
||||||
|
# берётся из соседних цен, а не из поля источника). Пересчитываем по той же
|
||||||
|
# формуле, иначе у следующей точки останется процент от фантомной базы.
|
||||||
|
recompute_diff_percent(price_changes)
|
||||||
|
|
||||||
# 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:
|
||||||
|
|
@ -831,11 +858,10 @@ def save_detail_enrichment(
|
||||||
)
|
)
|
||||||
|
|
||||||
# Дом-поля — только если листинг нашёлся: без строки в listings нет и house_id_fk.
|
# Дом-поля — только если листинг нашёлся: без строки в listings нет и house_id_fk.
|
||||||
if result.rowcount > 0:
|
if found:
|
||||||
_fill_house_params_from_detail(db, listing_id, e)
|
_fill_house_params_from_detail(db, listing_id, e)
|
||||||
|
|
||||||
db.commit()
|
db.commit()
|
||||||
found = result.rowcount > 0
|
|
||||||
if not found:
|
if not found:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"domclick_detail save: listing not found id=%s (item_id=%s)",
|
"domclick_detail save: listing not found id=%s (item_id=%s)",
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
250
tradein-mvp/scripts/sql/286_dryrun.sql
Normal file
250
tradein-mvp/scripts/sql/286_dryrun.sql
Normal file
|
|
@ -0,0 +1,250 @@
|
||||||
|
-- 286_dryrun.sql — прогон миграции 286 на проде БЕЗ записи (#3376).
|
||||||
|
--
|
||||||
|
-- ЗАЧЕМ. Миграция 286 удаляет строки из offer_price_history. Прежде чем этому
|
||||||
|
-- уехать на деплой, надо увидеть на БОЕВЫХ данных: сколько строк она заберёт по
|
||||||
|
-- источникам, сколько соседей пересчитает, сколько процентов обнулит — и что
|
||||||
|
-- ВТОРОЙ прогон подряд не делает уже ничего (идемпотентность).
|
||||||
|
--
|
||||||
|
-- КАК ЗАПУСКАТЬ (ничего не остаётся, в конце ROLLBACK):
|
||||||
|
-- ssh poincare
|
||||||
|
-- docker exec -i tradein-postgres psql -U <user> -d tradein -v ON_ERROR_STOP=1 \
|
||||||
|
-- -f - < 286_dryrun.sql
|
||||||
|
-- (или скопировать файл в контейнер и psql -f). Смотреть строки NOTICE.
|
||||||
|
--
|
||||||
|
-- ЧЕМ ОТЛИЧАЕТСЯ ОТ 286 — ровно тремя вещами, всё остальное скопировано дословно:
|
||||||
|
-- 1. вся работа обёрнута в FOR pass_no IN 1..2 — два прохода подряд в ОДНОЙ
|
||||||
|
-- транзакции; второй обязан дать deleted=0 recomputed=0 nulled=0;
|
||||||
|
-- 2. порог 200 даёт WARNING, а не EXCEPTION: на dry-run важнее увидеть число и
|
||||||
|
-- добежать до второго прохода, чем оборвать транзакцию;
|
||||||
|
-- 3. ROLLBACK вместо COMMIT.
|
||||||
|
-- ФАЙЛ РАЗОВЫЙ и живёт вне data/sql (автоприменение его не подхватывает). Правишь
|
||||||
|
-- 286 — правь и здесь, иначе сверять будет нечего.
|
||||||
|
BEGIN;
|
||||||
|
SET LOCAL lock_timeout = '5s';
|
||||||
|
|
||||||
|
-- ── ПРЕДЗАМЕР ОДНИМ SELECT (для сверки с числами в шапке 286) ────────────────
|
||||||
|
-- Правило внутри серии пошаговое (база = предыдущая ОСТАВЛЕННАЯ точка) и одним
|
||||||
|
-- SELECT не выражается — его считает DO-блок ниже. А эти два выражаются.
|
||||||
|
|
||||||
|
-- (а) кандидаты правила ПЕРВОЙ точки по СЫРЫМ строкам, до удалений фазы 1.
|
||||||
|
-- Ожидание из шапки 286: 4, все domklik (yandex 0 — их свидетелем была бы
|
||||||
|
-- listings.price_rub, а откат на неё снят).
|
||||||
|
SELECT 'предзамер: первая точка (сырые строки)' AS metric, k.source, count(*)
|
||||||
|
FROM (
|
||||||
|
SELECT oph.source,
|
||||||
|
oph.price_rub,
|
||||||
|
row_number() OVER w AS rn,
|
||||||
|
lead(oph.price_rub) OVER w AS second_price,
|
||||||
|
lead(oph.price_rub, 2) OVER w AS third_price
|
||||||
|
FROM offer_price_history oph
|
||||||
|
WHERE oph.change_time <> oph.recorded_at
|
||||||
|
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
|
||||||
|
) k
|
||||||
|
WHERE k.rn = 1
|
||||||
|
AND k.price_rub > 0
|
||||||
|
AND k.second_price > 0
|
||||||
|
AND (k.second_price / k.price_rub BETWEEN 9.5 AND 10.5
|
||||||
|
OR k.price_rub / k.second_price BETWEEN 9.5 AND 10.5)
|
||||||
|
AND k.third_price > 0
|
||||||
|
AND k.third_price / k.second_price BETWEEN 0.9 AND 1.1
|
||||||
|
GROUP BY k.source
|
||||||
|
ORDER BY k.source;
|
||||||
|
|
||||||
|
-- (б) сколько строк заберёт финальный шаг |diff_percent| > 100 → NULL.
|
||||||
|
SELECT 'предзамер: |diff_percent| > 100' AS metric, source, count(*)
|
||||||
|
FROM offer_price_history
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
AND abs(diff_percent) > 100
|
||||||
|
GROUP BY source
|
||||||
|
ORDER BY source;
|
||||||
|
|
||||||
|
-- ── ДВА ПРОХОДА ЛОГИКИ 286 ───────────────────────────────────────────────────
|
||||||
|
CREATE TEMP TABLE oph_decimal_slips (
|
||||||
|
id bigint PRIMARY KEY,
|
||||||
|
listing_id bigint NOT NULL,
|
||||||
|
next_id bigint,
|
||||||
|
source text NOT NULL
|
||||||
|
) ON COMMIT DROP;
|
||||||
|
|
||||||
|
DO $$
|
||||||
|
DECLARE
|
||||||
|
pass_no int;
|
||||||
|
r record;
|
||||||
|
cur_listing bigint;
|
||||||
|
base numeric;
|
||||||
|
witness numeric;
|
||||||
|
confirmed boolean;
|
||||||
|
to_delete bigint;
|
||||||
|
per_source text;
|
||||||
|
deleted_rows bigint;
|
||||||
|
updated_rows bigint;
|
||||||
|
nulled_rows bigint;
|
||||||
|
nulled_src text;
|
||||||
|
BEGIN
|
||||||
|
FOR pass_no IN 1..2 LOOP
|
||||||
|
DELETE FROM oph_decimal_slips;
|
||||||
|
cur_listing := NULL;
|
||||||
|
base := NULL;
|
||||||
|
|
||||||
|
-- ФАЗА 1 — правило внутри серии, пошагово как в гейте.
|
||||||
|
FOR r IN
|
||||||
|
SELECT oph.id,
|
||||||
|
oph.listing_id,
|
||||||
|
oph.source,
|
||||||
|
oph.price_rub,
|
||||||
|
lead(oph.price_rub) OVER w AS next_price,
|
||||||
|
lead(oph.id) OVER w AS next_id,
|
||||||
|
l.price_rub AS listing_price
|
||||||
|
FROM offer_price_history oph
|
||||||
|
LEFT JOIN listings l ON l.id = oph.listing_id
|
||||||
|
WHERE oph.change_time <> oph.recorded_at
|
||||||
|
AND oph.listing_id IN (
|
||||||
|
SELECT listing_id
|
||||||
|
FROM offer_price_history
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
GROUP BY listing_id
|
||||||
|
HAVING count(*) >= 2
|
||||||
|
)
|
||||||
|
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
|
||||||
|
ORDER BY oph.listing_id, oph.change_time, oph.id
|
||||||
|
LOOP
|
||||||
|
IF cur_listing IS DISTINCT FROM r.listing_id THEN
|
||||||
|
cur_listing := r.listing_id;
|
||||||
|
base := NULL;
|
||||||
|
END IF;
|
||||||
|
|
||||||
|
witness := COALESCE(r.next_price, r.listing_price);
|
||||||
|
|
||||||
|
IF base IS NOT NULL
|
||||||
|
AND base > 0 AND r.price_rub > 0 AND witness > 0
|
||||||
|
AND (r.price_rub / base BETWEEN 9.5 AND 10.5
|
||||||
|
OR base / r.price_rub BETWEEN 9.5 AND 10.5)
|
||||||
|
AND witness / base BETWEEN 0.9 AND 1.1
|
||||||
|
THEN
|
||||||
|
SELECT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM offer_price_history t
|
||||||
|
WHERE t.listing_id = r.listing_id
|
||||||
|
AND t.change_time = t.recorded_at
|
||||||
|
AND t.price_rub = r.price_rub
|
||||||
|
)
|
||||||
|
INTO confirmed;
|
||||||
|
|
||||||
|
IF NOT confirmed THEN
|
||||||
|
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
|
||||||
|
VALUES (r.id, r.listing_id, r.next_id, r.source);
|
||||||
|
CONTINUE;
|
||||||
|
END IF;
|
||||||
|
END IF;
|
||||||
|
|
||||||
|
IF r.price_rub > 0 THEN
|
||||||
|
base := r.price_rub;
|
||||||
|
END IF;
|
||||||
|
END LOOP;
|
||||||
|
|
||||||
|
-- ФАЗА 2 — правило первой точки, по оставшимся строкам.
|
||||||
|
INSERT INTO oph_decimal_slips (id, listing_id, next_id, source)
|
||||||
|
WITH kept AS (
|
||||||
|
SELECT oph.id,
|
||||||
|
oph.listing_id,
|
||||||
|
oph.source,
|
||||||
|
oph.price_rub,
|
||||||
|
row_number() OVER w AS rn,
|
||||||
|
lead(oph.id) OVER w AS second_id,
|
||||||
|
lead(oph.price_rub) OVER w AS second_price,
|
||||||
|
lead(oph.price_rub, 2) OVER w AS third_price
|
||||||
|
FROM offer_price_history oph
|
||||||
|
WHERE oph.change_time <> oph.recorded_at
|
||||||
|
AND NOT EXISTS (SELECT 1 FROM oph_decimal_slips s WHERE s.id = oph.id)
|
||||||
|
WINDOW w AS (PARTITION BY oph.listing_id ORDER BY oph.change_time, oph.id)
|
||||||
|
)
|
||||||
|
SELECT k.id, k.listing_id, k.second_id, k.source
|
||||||
|
FROM kept k
|
||||||
|
WHERE k.rn = 1
|
||||||
|
AND k.price_rub > 0
|
||||||
|
AND k.second_price > 0
|
||||||
|
AND (k.second_price / k.price_rub BETWEEN 9.5 AND 10.5
|
||||||
|
OR k.price_rub / k.second_price BETWEEN 9.5 AND 10.5)
|
||||||
|
-- Свидетель — только третья точка истории (см. шапку 286).
|
||||||
|
AND k.third_price > 0
|
||||||
|
AND k.third_price / k.second_price BETWEEN 0.9 AND 1.1
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM offer_price_history t
|
||||||
|
WHERE t.listing_id = k.listing_id
|
||||||
|
AND t.change_time = t.recorded_at
|
||||||
|
AND t.price_rub = k.price_rub
|
||||||
|
);
|
||||||
|
|
||||||
|
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;
|
||||||
|
|
||||||
|
-- Отличие 2 от 286: здесь WARNING, а не EXCEPTION.
|
||||||
|
IF to_delete > 200 THEN
|
||||||
|
RAISE WARNING 'pass%: кандидатов % > порога 200 — на проде 286 здесь УПАЛА БЫ',
|
||||||
|
pass_no, 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;
|
||||||
|
|
||||||
|
WITH neighbours AS (
|
||||||
|
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 = rc.new_diff
|
||||||
|
FROM recomputed rc
|
||||||
|
WHERE oph.id = rc.id
|
||||||
|
AND oph.diff_percent IS DISTINCT FROM rc.new_diff;
|
||||||
|
GET DIAGNOSTICS updated_rows = ROW_COUNT;
|
||||||
|
|
||||||
|
SELECT string_agg(source || '=' || cnt, ', ' ORDER BY source)
|
||||||
|
INTO nulled_src
|
||||||
|
FROM (
|
||||||
|
SELECT source, count(*) AS cnt
|
||||||
|
FROM offer_price_history
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
AND abs(diff_percent) > 100
|
||||||
|
GROUP BY source
|
||||||
|
) s;
|
||||||
|
|
||||||
|
UPDATE offer_price_history
|
||||||
|
SET diff_percent = NULL
|
||||||
|
WHERE change_time <> recorded_at
|
||||||
|
AND abs(diff_percent) > 100;
|
||||||
|
GET DIAGNOSTICS nulled_rows = ROW_COUNT;
|
||||||
|
|
||||||
|
RAISE NOTICE 'pass%: deleted=% (%) recomputed=% nulled=% (%)',
|
||||||
|
pass_no,
|
||||||
|
deleted_rows, COALESCE(per_source, 'ни одного'),
|
||||||
|
updated_rows,
|
||||||
|
nulled_rows, COALESCE(nulled_src, 'ни одной');
|
||||||
|
END LOOP;
|
||||||
|
END $$;
|
||||||
|
|
||||||
|
ROLLBACK;
|
||||||
Loading…
Add table
Reference in a new issue