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
Ревью: гейт охранял не то место. Подмена знаменателя красила три теста, но дефект «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 (ручка отдавала её неотличимо от свежей). Строки вне сегодняшнего набора удаляются в той же транзакции. На ПУСТОМ наборе чистка не ходит: разом отвалившиеся все входы — признак поломки прогона, а не пяти одновременных «данных больше нет». Оба поведения покрыты тестами, оба проверены на сломанном коде.
397 lines
21 KiB
Python
397 lines
21 KiB
Python
"""Пересчёт витринных метрик публичного лэндинга МЕРЫ (таблица 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
|