gendesign/tradein-mvp/backend/app/tasks/landing_stats.py
bot-backend e9a2fff0b3
All checks were successful
CI Trade-In / changes (pull_request) Successful in 13s
CI / changes (pull_request) Successful in 16s
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 5m32s
test(mera/b2c): гейт на сам SQL доли снижений + чистка протухших метрик
Ревью: гейт охранял не то место. Подмена знаменателя красила три теста, но
дефект «84.8% вместо 48.1%» живёт в SQL — во включении однострочных записей
истории (у domklik одна запись = «цену не менял») в знаменатель. Ревьюер
вернул дефект условием n_rows >= 2 в CTE moved, и все 25 тестов остались
зелёными: текстовые пины держали только span_days и max_abs_pct.

Новый пин держит обе половины: однострочные попадают в moved веткой CASE со
значением 0, и нигде в запросе нет фильтра по числу записей истории (ни в
WHERE, ни HAVING). Живой прогон на подготовленных строках не заведён
намеренно: DATABASE_URL в тестовой джобе — заглушка, Postgres там нет, и тест
по образцу test_purge_expired_trade_in_data.py молча скипался бы, то есть не
гейтил бы ничего. Фальсифицировано руками — с n_rows >= 2 тест красный и
называет причину.

Второе: метрика, у которой пропал вход, больше не доживает в таблице со
старым computed_at (ручка отдавала её неотличимо от свежей). Строки вне
сегодняшнего набора удаляются в той же транзакции. На ПУСТОМ наборе чистка
не ходит: разом отвалившиеся все входы — признак поломки прогона, а не пяти
одновременных «данных больше нет». Оба поведения покрыты тестами, оба
проверены на сломанном коде.
2026-08-29 19:10:04 +05:00

397 lines
21 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Пересчёт витринных метрик публичного лэндинга МЕРЫ (таблица landing_stats).
ЗАЧЕМ
-----
Числа на лэндинге (frontend/src/app/mera-public/marketing-v3.ts) были литералами
— то есть придуманными. Публичная страница, которая продаёт «расчёт по данным»,
не может показывать цифры, которых в данных нет: это ровно та подмена, против
которой продукт и позиционируется. Здесь каждая витринная величина считается
запросом к проду, и вместе с ней пишется размер выборки.
ГЛАВНОЕ ПРАВИЛО: НЕТ ВХОДА — НЕТ СТРОКИ
---------------------------------------
Ни одна метрика не пишется с подставленным значением. Если выборка пуста
(нет оценок, нет истории цен, нет сделок) — строка в landing_stats просто не
появляется, ручка её не отдаёт, фронт не рисует блок. Ноль здесь читался бы как
измеренный ноль («ни одно объявление не снижало цену»), а это враньё другого
рода, чем отсутствие данных. Правило действует и на ВТОРОМ прогоне: пропавшая
метрика удаляется из таблицы (см. refresh_landing_stats), иначе она осталась бы
на витрине со старым computed_at и читалась бы как измеренная сегодня.
ЧЕГО ЗДЕСЬ НЕТ И НЕ БУДЕТ
-------------------------
«Точность прогноза» и «срок продажи» — величин с такими именами в базе нет.
Точность считает бэктест (своя задача, свои допущения), а срок продажи требует
пары «объявление снято → сделка», которой у нас нет: снятие объявления не
означает продажу. `listing_age_median_days` НЕ является сроком продажи и назван
экспозицией активного объявления — см. note метрики.
ПОЧЕМУ ТОЛЬКО DOMKLIK В ЦЕНОВЫХ МЕТРИКАХ
----------------------------------------
`offer_price_history` наполняется триггером, и наполняется по-разному:
у avito/yandex стартовая цена в историю НЕ пишется (первая строка появляется
только при изменении, то есть «снизил» и «не снижал» неразличимы), а yandex
вдобавок сеет синтетическую пару со сдвигом в сутки. Считать долю снижений по
такой смеси — значит получить число, у которого нет смысла. Domklik пишет старт,
поэтому только он.
Задача синхронная (только SELECT'ы + UPSERT), запускается kit-scheduler'ом через
product_handlers._job_landing_stats в run_in_executor — по образцу
deals_freshness_monitor / listing_source_snapshot.
"""
from __future__ import annotations
import logging
from decimal import Decimal
from typing import Any
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.services import scrape_runs as runs_mod
logger = logging.getLogger(__name__)
__all__ = ["EKB", "collect_landing_metrics", "refresh_landing_stats"]
# Город витрины. Лэндинг сегодня продаёт Екатеринбург, и метрики обязаны быть
# про него же: медиана по всей области смешала бы рынки с разной динамикой.
EKB = "Екатеринбург"
# Порог наблюдения для ценовых метрик. За две недели объявление успевает получить
# первую правку цены; более короткие живут слишком мало, чтобы «не снижал» было
# наблюдением, а не «не успел».
_PRICE_SPAN_DAYS = 14
# Отсечка аномалий: изменение больше 30% за наблюдение — это, как правило, смена
# объекта под тем же id (перевыставили другую квартиру) или опечатка в цене,
# а не торг. Медиану такие хвосты не двигают, но долю снижений — двигают.
_PRICE_MAX_ABS_PCT = 30
# ── Оценки ──────────────────────────────────────────────────────────────────
# Период считаем по фактическим краям created_at, а не «с даты запуска»: витрина
# обещает «за N дней работы», и N должен быть измеренным.
_ESTIMATES_SQL = text("""
SELECT count(*) AS total,
EXTRACT(EPOCH FROM (max(created_at) - min(created_at)))
/ 86400.0 AS period_days
FROM trade_in_estimates
""")
# n_analogs > 0: оценка без аналогов — это отказ расчёта, а не «ноль аналогов»;
# включив её, мы бы занизили медиану наблюдениями, где измерять было нечего.
_ANALOGS_SQL = text("""
SELECT count(*) AS n,
percentile_cont(0.5) WITHIN GROUP (ORDER BY n_analogs) AS median
FROM trade_in_estimates
WHERE n_analogs > 0
""")
# Возраст АКТИВНОГО объявления = экспозиция на сегодня, а не срок продажи:
# знаменатель — те, кто ещё висит, поэтому величина по построению занижена
# относительно «сколько в итоге продавалось». Это ограничение уезжает в note.
_LISTING_AGE_SQL = text("""
SELECT count(*) AS n,
percentile_cont(0.5) WITHIN GROUP (
ORDER BY (CURRENT_DATE - listing_date)
) AS median
FROM listings
WHERE is_active
AND city = CAST(:city AS text)
AND listing_date IS NOT NULL
AND listing_date <= CURRENT_DATE
""")
# ── Динамика цены объявлений ────────────────────────────────────────────────
#
# Знаменатель — объявления, которые МОЖНО было наблюдать: от первой записи в
# истории до последнего показа прошло >= 14 дней. Сюда попадают и те, у кого
# запись одна (domklik пишет старт → одна запись означает «цену не менял»); без
# них доля снижений считалась бы только по менявшим и давала 85% вместо 48%.
#
# Скорость снижения нормируем на 30 дней по интервалу МЕЖДУ КРАЙНИМИ ПРАВКАМИ,
# а не по всему наблюдению: цена не менялась после последней правки, и растягивая
# знаменатель на «висит до сих пор», мы измеряли бы терпение продавца, а не торг.
_PRICE_MOVES_SQL = text("""
WITH hist AS (
SELECT listing_id,
min(change_time) AS first_change,
max(change_time) AS last_change,
count(*) AS n_rows
FROM offer_price_history
WHERE source = 'domklik'
GROUP BY listing_id
),
observed AS (
SELECT h.listing_id,
h.n_rows,
EXTRACT(EPOCH FROM (h.last_change - h.first_change)) / 86400.0 AS change_days
FROM hist h
JOIN listings l ON l.id = h.listing_id
WHERE GREATEST(h.last_change, COALESCE(l.last_seen_at, h.last_change)) - h.first_change
>= make_interval(days => CAST(:span_days AS integer))
),
priced AS (
SELECT o.listing_id,
o.n_rows,
o.change_days,
(SELECT p.price_rub FROM offer_price_history p
WHERE p.listing_id = o.listing_id AND p.source = 'domklik'
ORDER BY p.change_time ASC, p.id ASC LIMIT 1) AS price_first,
(SELECT p.price_rub FROM offer_price_history p
WHERE p.listing_id = o.listing_id AND p.source = 'domklik'
ORDER BY p.change_time DESC, p.id DESC LIMIT 1) AS price_last
FROM observed o
),
moved AS (
SELECT listing_id,
change_days,
CASE WHEN n_rows >= 2
THEN (price_last - price_first) / price_first * 100.0
ELSE 0
END AS pct
FROM priced
WHERE price_first IS NOT NULL AND price_first > 0
)
SELECT count(*) AS n,
count(*) FILTER (WHERE pct < 0) AS n_cut,
percentile_cont(0.5) WITHIN GROUP (
ORDER BY pct * 30.0 / NULLIF(change_days, 0)
) FILTER (WHERE pct < 0) AS median_pct_per_month
FROM moved
WHERE abs(pct) <= CAST(:max_abs_pct AS numeric)
""")
# 12 месяцев от сегодня. deal_date у Росреестра — лейбл начала квартала, поэтому
# окно накрывает 4-5 кварталов и число «за год» тут приблизительно по построению;
# это сказано в note, а не спрятано.
_DEALS_SQL = text("""
SELECT count(*) AS n
FROM deals
WHERE city = CAST(:city AS text)
AND deal_date >= (CURRENT_DATE - INTERVAL '12 months')
""")
_UPSERT_SQL = text("""
INSERT INTO landing_stats (metric, value_num, value_text, sample_n, note, computed_at)
VALUES (
CAST(:metric AS text),
CAST(:value_num AS numeric),
CAST(:value_text AS text),
CAST(:sample_n AS integer),
CAST(:note AS text),
now()
)
ON CONFLICT (metric) DO UPDATE SET
value_num = EXCLUDED.value_num,
value_text = EXCLUDED.value_text,
sample_n = EXCLUDED.sample_n,
note = EXCLUDED.note,
computed_at = EXCLUDED.computed_at
""")
# Строки метрик, которых в СЕГОДНЯШНЕМ наборе нет, удаляются. Метрика исчезает
# из набора ровно тогда, когда у неё пропал вход (см. «нет входа — нет строки»),
# и оставленная строка продолжала бы отдаваться ручкой как обычная — со старым
# computed_at, который витрина не обязана читать. Удалённая метрика — блок,
# которого на странице нет; протухшая — блок с враньём.
_PRUNE_SQL = text("""
DELETE FROM landing_stats
WHERE metric <> ALL(CAST(:kept AS text[]))
""")
def _num(value: Any) -> float | None:
"""Привести значение агрегата к float; None остаётся None.
percentile_cont возвращает Decimal/float в зависимости от типа входа — в
numeric-колонку и в JSON поедет одинаково только после явного приведения.
"""
if value is None:
return None
if isinstance(value, Decimal):
return float(value)
return float(value)
def collect_landing_metrics(db: Session) -> list[dict[str, Any]]:
"""Посчитать метрики витрины. Метрика без данных в список НЕ попадает.
Отделено от записи, чтобы тест мог проверить сами ЗНАЧЕНИЯ на подготовленной
базе, не разбирая по дороге счётчики прогона.
"""
metrics: list[dict[str, Any]] = []
row = db.execute(_ESTIMATES_SQL).first()
total = int(row.total) if row is not None and row.total else 0
if total > 0:
metrics.append(
{
"metric": "estimates_total",
"value_num": float(total),
"value_text": None,
"sample_n": total,
"note": "Расчётов сделано в системе (все города, весь срок работы)",
}
)
period = _num(row.period_days)
# Один-единственный расчёт даёт период 0 дней — это не измерение, а
# артефакт единственной точки; такую строку не пишем.
if period is not None and total > 1:
metrics.append(
{
"metric": "estimates_period_days",
"value_num": round(period, 1),
"value_text": None,
"sample_n": total,
"note": "Дней между первым и последним расчётом",
}
)
row = db.execute(_ANALOGS_SQL).first()
if row is not None and row.n and _num(row.median) is not None:
metrics.append(
{
"metric": "analogs_median",
"value_num": round(_num(row.median) or 0.0, 1),
"value_text": None,
"sample_n": int(row.n),
"note": "Медиана числа аналогов на расчёт (только расчёты, где аналоги нашлись)",
}
)
row = db.execute(_LISTING_AGE_SQL, {"city": EKB}).first()
if row is not None and row.n and _num(row.median) is not None:
metrics.append(
{
"metric": "listing_age_median_days",
"value_num": round(_num(row.median) or 0.0, 1),
"value_text": None,
"sample_n": int(row.n),
"note": (
"Медианная ЭКСПОЗИЦИЯ активного объявления в Екатеринбурге "
"(сколько дней висит на сегодня). Это НЕ срок продажи: "
"считается по тем, кто ещё продаётся, и снятие объявления "
"не означает сделку"
),
}
)
row = db.execute(
_PRICE_MOVES_SQL,
{"span_days": _PRICE_SPAN_DAYS, "max_abs_pct": _PRICE_MAX_ABS_PCT},
).first()
if row is not None and row.n:
n = int(row.n)
base_note = (
f"Только Домклик (единственный источник, где триггер пишет стартовую цену), "
f"наблюдение от {_PRICE_SPAN_DAYS} дней, изменения свыше "
f"{_PRICE_MAX_ABS_PCT}% отброшены как смена объекта"
)
metrics.append(
{
"metric": "price_cut_share_pct",
"value_num": round(int(row.n_cut) * 100.0 / n, 1),
"value_text": None,
"sample_n": n,
"note": f"Доля объявлений, снижавших цену. {base_note}",
}
)
median_move = _num(row.median_pct_per_month)
if median_move is not None:
metrics.append(
{
"metric": "price_cut_median_pct_per_month",
"value_num": round(median_move, 2),
"value_text": None,
# Выборка ЗДЕСЬ — только снижавшие: медиана считается по ним,
# и подставить сюда общий n значило бы приписать величине
# выборку, по которой её не считали.
"sample_n": int(row.n_cut),
"note": (
f"Медианное изменение цены за 30 дней среди снижавших "
f"(отрицательное). {base_note}"
),
}
)
row = db.execute(_DEALS_SQL, {"city": EKB}).first()
if row is not None and row.n:
metrics.append(
{
"metric": "deals_total_12m",
"value_num": float(row.n),
"value_text": None,
"sample_n": int(row.n),
"note": (
"Сделок Росреестра по Екатеринбургу за последние 12 месяцев. "
"Дата сделки — лейбл начала квартала, поэтому окно накрывает "
"целые кварталы, а не ровно год"
),
}
)
return metrics
def refresh_landing_stats(
db: Session,
run_id: int,
params: dict[str, Any] | None = None,
) -> dict[str, int]:
"""Пересчитать landing_stats и финализировать прогон.
Sync (вызывается scheduler-триггером в executor, как check_deals_freshness).
`params` не используется — принимается ради единой сигнатуры обработчиков.
Метрика, у которой пропал вход, СНИМАЕТСЯ с витрины, а не доживает со старым
computed_at: строки, которых нет в сегодняшнем наборе, удаляются в той же
транзакции. Иначе «нет входа — нет строки» действует только на первом
прогоне, а дальше отсутствие данных выглядит как данные — ручка отдаёт такую
строку неотличимо от свежей, и отличить её можно только сравнив computed_at с
соседями, чего фронт не делает.
Пустой результат — НЕ ошибка прогона: на свежей базе метрик может не быть ни
одной, и падать в failed из-за этого значит завести шумный алерт там, где
система работает штатно. Но и чистка в этом случае НЕ выполняется: разом
отвалившиеся все входы — это признак поломки самого прогона (пустая/недоступная
база), а не пяти одновременных «данных больше нет», и стирать по такому
признаку всю витрину нельзя. Чистка ходит только с непустым набором, где
пропажу конкретной метрики видно на фоне посчитавшихся соседей.
"""
del params
counters: dict[str, int] = {"metrics_written": 0, "metrics_removed": 0}
try:
runs_mod.update_heartbeat(db, run_id, counters)
metrics = collect_landing_metrics(db)
for row in metrics:
db.execute(_UPSERT_SQL, row)
if metrics:
removed = db.execute(_PRUNE_SQL, {"kept": [m["metric"] for m in metrics]})
counters["metrics_removed"] = int(removed.rowcount or 0)
db.commit()
counters["metrics_written"] = len(metrics)
if not metrics:
logger.warning(
"landing_stats run_id=%d: ни одной метрики не посчиталось — "
"витрина покажет прошлый срез (или пусто, если его не было)",
run_id,
)
runs_mod.mark_done(db, run_id, counters)
logger.info(
"refresh_landing_stats run_id=%d done: %d метрик (%s)",
run_id,
len(metrics),
", ".join(m["metric"] for m in metrics) or "",
)
return counters
except Exception as exc:
logger.exception("refresh_landing_stats run_id=%d failed", run_id)
try:
db.rollback()
except Exception:
pass
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters)
raise