gendesign/backend/app/services/site_finder/velocity.py
lekss361 b7ce4a589c fix(sf-10): velocity LEFT JOIN + velocity_data_available flag
Replace INNER JOIN on objective_complex_mapping with LEFT JOIN approach:
competitors without mapping now return velocity=0 with
velocity_data_available=False instead of being silently dropped.

UI: VelocityBlock shows 'нет данных velocity' badge when flag is False
and hides the gauge (meaningless zero). TS type updated (additive optional
field, backward compat).

Audit (EKB, region_cd=66): mapped=129, unmapped=1387, total=1516.

Closes #271 item #10
2026-05-17 15:51:43 +03:00

491 lines
22 KiB
Python
Raw 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.

"""Velocity-score — темп продаж конкурентов вокруг участка.
Per #34 D2: утилизация objective_corpus_room_month (еженедельно обновляемые данные).
Ранее использовался domrf_kn_sale_graph (последнее обновление 2026-01, устарел).
Главный demand-сигнал «продастся ли» — среднемесячный объём продаж
конкурирующих ЖК в радиусе radius_km от участка, нормированный к
ЕКБ-медиане по данным Objective.
Foundation: domrf_kn_objects (lat/lon, comm_name, obj_class, region_cd),
objective_complex_mapping (domrf_obj_id ↔ objective_complex_name),
objective_corpus_room_month (project_name, deals_total_vol_m2,
deals_total_count, report_month).
Linkage: domrf_kn_objects.obj_id
→ objective_complex_mapping.domrf_obj_id
→ objective_complex_mapping.objective_complex_name
→ objective_corpus_room_month.project_name
"""
from __future__ import annotations
import logging
from dataclasses import dataclass
from typing import Any, Literal
from sqlalchemy import text
from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
# Fallback если в БД нет данных за окно months_window.
# Эмпирика по ЕКБ: ~4 500 м²/мес на один ЖК (apartments, 2024-2025).
_EKB_MEDIAN_FALLBACK_SQM_PER_MONTH: float = 4500.0
@dataclass(frozen=True)
class VelocityResult:
"""Результат расчёта velocity-score для участка."""
competitors_count: int
monthly_velocity_sqm: float # avg м²/мес по конкурентам в радиусе
ekb_median_sqm: float # benchmark ЕКБ для нормализации
velocity_score: float # 0..1 — отношение к benchmark
confidence: Literal["high", "medium", "low"]
months_observed: int # фактически использованных месяцев
period_start: str # YYYY-MM
period_end: str # YYYY-MM
sample_competitors: list[dict[str, Any]] # top-5 для UI
by_room_bucket: dict[str, dict[str, Any]] # агрегат по комнатности
# True если ≥1 конкурент имеет маппинг в objective_complex_mapping;
# False → конкуренты найдены, но данных Objective нет — velocity = 0.
velocity_data_available: bool = True
def as_dict(self) -> dict[str, Any]:
return {
"competitors_count": self.competitors_count,
"monthly_velocity_sqm": round(self.monthly_velocity_sqm, 1),
"ekb_median_sqm": round(self.ekb_median_sqm, 1),
"velocity_score": round(self.velocity_score, 3),
"confidence": self.confidence,
"months_observed": self.months_observed,
"period": {"start": self.period_start, "end": self.period_end},
"sample_competitors": self.sample_competitors,
"by_room_bucket": self.by_room_bucket,
"velocity_data_available": self.velocity_data_available,
}
def compute_velocity(
db: Session,
parcel_geom_wkt: str,
radius_km: float = 3.0,
obj_class: str | None = None,
months_window: int = 6,
) -> VelocityResult | None:
"""Вычислить velocity-score для участка.
Алгоритм:
1. Найти ЖК-конкуренты в радиусе radius_km (через lat/lon ST_DWithin).
2. Взять objective_corpus_room_month за последние months_window месяцев
через objective_complex_mapping (domrf_obj_id → project_name).
3. Посчитать суммарный объём deals_total_vol_m2.
4. Нормировать на ЕКБ-медиану → score 0..1.
Возвращает None если parcel_geom_wkt невалиден или конкурентов нет.
"""
# ── Step 1: конкуренты по lat/lon в радиусе ──────────────────────────────
# DISTINCT ON (obj_id) ORDER BY snapshot_date DESC — latest snapshot only.
# obj_class в domrf_kn_objects заполнен слабо (много NULL); фильтруем
# только если явно передан.
class_filter = "AND o.obj_class = :obj_class" if obj_class else ""
# SAVEPOINT per query: failure rollbacks ТОЛЬКО savepoint, не outer tx.
# db.rollback() здесь НЕЛЬЗЯ — он orphan'ит outer SessionTransaction
# (см. PR #155 bot review — SQLAlchemy 2.0 begin_nested context cleanup).
try:
with db.begin_nested():
comp_rows = (
db.execute(
text(
f"""
WITH latest_obj AS (
SELECT DISTINCT ON (obj_id)
obj_id,
comm_name,
dev_name,
obj_class,
latitude,
longitude,
district_name
FROM domrf_kn_objects
WHERE latitude IS NOT NULL
AND longitude IS NOT NULL
AND region_cd = 66
{class_filter}
ORDER BY obj_id, snapshot_date DESC NULLS LAST
)
SELECT
o.obj_id,
o.comm_name,
o.dev_name,
o.obj_class,
o.district_name,
ST_Distance(
ST_SetSRID(ST_MakePoint(o.longitude, o.latitude), 4326)::geography,
ST_Centroid(ST_GeomFromText(:parcel_wkt, 4326))::geography
) AS distance_m
FROM latest_obj o
WHERE ST_DWithin(
ST_SetSRID(ST_MakePoint(o.longitude, o.latitude), 4326)::geography,
ST_Centroid(ST_GeomFromText(:parcel_wkt, 4326))::geography,
:radius_m
)
ORDER BY distance_m ASC
LIMIT 200
"""
),
{
"parcel_wkt": parcel_geom_wkt,
"radius_m": radius_km * 1000.0,
"obj_class": obj_class,
},
)
.mappings()
.all()
)
except Exception:
logger.exception("velocity: competitor query failed for wkt=%s", parcel_geom_wkt[:80])
# SAVEPOINT auto-rollbacks через __exit__ context manager.
# Outer tx остаётся clean — caller продолжает работать без cascade.
return None
if not comp_rows:
return None
obj_ids: list[int] = [int(r["obj_id"]) for r in comp_rows]
competitor_meta: dict[int, dict[str, Any]] = {
int(r["obj_id"]): {
"name": r["comm_name"],
"dev_name": r["dev_name"],
"obj_class": r["obj_class"],
"district_name": r["district_name"],
"distance_m": round(float(r["distance_m"]), 0),
}
for r in comp_rows
}
# ── Step 2: objective_corpus_room_month за последние N месяцев ───────────
# Linkage: domrf_obj_id → objective_complex_mapping → project_name →
# objective_corpus_room_month.
# deals_total_vol_m2 = DDU + DKP м² за месяц (primary signal).
# deals_total_count > 0 — фильтрует месяцы без сделок.
# LEFT JOIN с objective_complex_mapping: конкуренты без маппинга не
# выпадают — они включаются в список но с total_sqm=NULL (→ 0.0).
# GROUP BY domrf_obj_id чтобы сохранить совместимость с caller (obj_id key).
try:
with db.begin_nested():
sales_rows = (
db.execute(
text(
"""
WITH all_competitors AS (
SELECT unnest(CAST(:obj_ids AS int[])) AS obj_id
),
mapped AS (
SELECT cm.domrf_obj_id AS obj_id,
cm.objective_complex_name
FROM objective_complex_mapping cm
WHERE cm.domrf_obj_id = ANY(:obj_ids)
)
SELECT
ac.obj_id,
SUM(COALESCE(crm.deals_total_vol_m2,
crm.deals_total_count * 45.0)) AS total_sqm,
COUNT(DISTINCT crm.report_month) AS months_with_data,
MIN(crm.report_month) AS period_start,
MAX(crm.report_month) AS period_end,
CASE WHEN m.obj_id IS NOT NULL THEN TRUE
ELSE FALSE END AS has_mapping
FROM all_competitors ac
LEFT JOIN mapped m ON m.obj_id = ac.obj_id
LEFT JOIN objective_corpus_room_month crm
ON crm.project_name = m.objective_complex_name
AND crm.report_month >= (
CURRENT_DATE - CAST(:window_interval AS interval))
AND crm.deals_total_count > 0
GROUP BY ac.obj_id, m.obj_id
"""
),
{
"obj_ids": obj_ids,
"window_interval": f"{months_window} months",
},
)
.mappings()
.all()
)
except Exception:
logger.exception("velocity: objective sales query failed for obj_ids=%s", obj_ids[:5])
# SAVEPOINT auto-rollback'нут — outer tx clean
return None
if not sales_rows:
return None
# Проверяем: есть ли хотя бы один конкурент с маппингом (has_mapping=True).
# Если нет — возвращаем velocity=0 с явным флагом velocity_data_available=False,
# вместо того чтобы отбросить всех конкурентов (старый INNER JOIN поведение).
has_any_mapping = any(bool(r["has_mapping"]) for r in sales_rows)
if not has_any_mapping:
logger.info(
"velocity: %d competitors found but none mapped in objective_complex_mapping;"
" returning velocity=0 with data_available=False",
len(obj_ids),
)
ekb_median = (
_get_ekb_median(db, months_window=months_window) or _EKB_MEDIAN_FALLBACK_SQM_PER_MONTH
)
n_comps = len(comp_rows)
sample = [
{
"obj_id": oid,
**competitor_meta[oid],
"total_sqm_period": 0.0,
"by_room_bucket": {},
}
for oid in obj_ids[:5]
if oid in competitor_meta
]
return VelocityResult(
competitors_count=n_comps,
monthly_velocity_sqm=0.0,
ekb_median_sqm=ekb_median,
velocity_score=0.0,
confidence="low",
months_observed=0,
period_start="",
period_end="",
sample_competitors=sample,
by_room_bucket={},
velocity_data_available=False,
)
# ── Step 2b: разбивка по комнатности (room_bucket) ───────────────────────
# Тот же маппинг domrf_obj_id → project_name. Агрегируем по room_bucket
# для отображения структуры спроса в UI.
try:
with db.begin_nested():
bucket_rows = (
db.execute(
text(
"""
WITH mapped AS (
SELECT cm.domrf_obj_id AS obj_id,
cm.objective_complex_name
FROM objective_complex_mapping cm
WHERE cm.domrf_obj_id = ANY(:obj_ids)
)
SELECT
m.obj_id,
crm.room_bucket,
SUM(crm.deals_total_count) AS units_sold,
SUM(COALESCE(crm.deals_total_vol_m2,
crm.deals_total_count * 45.0)) AS sqm_sold
FROM objective_corpus_room_month crm
JOIN mapped m ON m.objective_complex_name = crm.project_name
WHERE crm.report_month >= (CURRENT_DATE - CAST(:window_interval AS interval))
AND crm.deals_total_count > 0
GROUP BY m.obj_id, crm.room_bucket
"""
),
{
"obj_ids": obj_ids,
"window_interval": f"{months_window} months",
},
)
.mappings()
.all()
)
except Exception:
logger.warning("velocity: bucket breakdown query failed, continuing without it")
bucket_rows = []
# Агрегируем по room_bucket поверх всех конкурентов.
by_bucket_agg: dict[str, dict[str, Any]] = {}
per_comp_buckets: dict[int, dict[str, int]] = {}
for row in bucket_rows:
bucket = str(row["room_bucket"])
oid = int(row["obj_id"])
units = int(row["units_sold"] or 0)
sqm = float(row["sqm_sold"] or 0.0)
if bucket not in by_bucket_agg:
by_bucket_agg[bucket] = {"units": 0, "sqm": 0.0, "complexes": set()}
by_bucket_agg[bucket]["units"] += units
by_bucket_agg[bucket]["sqm"] += sqm
by_bucket_agg[bucket]["complexes"].add(oid)
if oid not in per_comp_buckets:
per_comp_buckets[oid] = {}
per_comp_buckets[oid][bucket] = units
by_room_bucket: dict[str, dict[str, Any]] = {
bucket: {
"units": data["units"],
"sqm": round(data["sqm"], 0),
"complexes_count": len(data["complexes"]),
}
for bucket, data in by_bucket_agg.items()
}
# Считаем только по строкам с маппингом — unmapped строки дают total_sqm=NULL.
mapped_sales_rows = [r for r in sales_rows if bool(r["has_mapping"])]
total_sqm = sum(float(r["total_sqm"] or 0.0) for r in mapped_sales_rows)
months_observed = max((int(r["months_with_data"] or 0) for r in mapped_sales_rows), default=0)
period_start_dates = [r["period_start"] for r in mapped_sales_rows if r["period_start"]]
period_end_dates = [r["period_end"] for r in mapped_sales_rows if r["period_end"]]
period_start = min(period_start_dates).strftime("%Y-%m") if period_start_dates else ""
period_end = max(period_end_dates).strftime("%Y-%m") if period_end_dates else ""
# ── Step 3: ЕКБ-медиана ──────────────────────────────────────────────────
ekb_median = (
_get_ekb_median(db, months_window=months_window) or _EKB_MEDIAN_FALLBACK_SQM_PER_MONTH
)
n_comps = len(comp_rows)
# Если mapped-конкурентов нет данных — partial coverage → velocity=0.
if months_observed == 0 or total_sqm <= 0:
logger.info(
"velocity: %d competitors found, %d mapped, but no sales data in window;"
" returning velocity=0 with data_available=False",
len(obj_ids),
len(mapped_sales_rows),
)
sample_partial = sorted(
[
{
"obj_id": oid,
**competitor_meta[oid],
"total_sqm_period": 0.0,
"by_room_bucket": {},
}
for oid in obj_ids
if oid in competitor_meta
],
key=lambda x: x["total_sqm_period"], # type: ignore[arg-type]
reverse=True,
)[:5]
return VelocityResult(
competitors_count=n_comps,
monthly_velocity_sqm=0.0,
ekb_median_sqm=ekb_median,
velocity_score=0.0,
confidence="low",
months_observed=0,
period_start="",
period_end="",
sample_competitors=sample_partial,
by_room_bucket={},
velocity_data_available=False,
)
# Среднемесячный объём в расчёте: суммарный по всем конкурентам / месяцев.
# Чем больше конкурентов с данными — тем весомее результат.
monthly_velocity = total_sqm / months_observed
# ── Step 4: нормализация → score 0..1 ────────────────────────────────────
# Логика: сравниваем суммарный velocity радиуса с «нормой» одного ЖК.
# Если в радиусе продаётся N × ekb_median → рынок горячий.
# Нормируем: score = min(1.0, total_velocity / (n_competitors × ekb_median × 2))
# Cap 2×median = «насыщен». Итоговый score 0..1.
# n_with_sales — только mapped конкуренты (у unmapped данных нет).
n_with_sales = len(mapped_sales_rows)
denominator = n_with_sales * ekb_median * 2.0 if n_with_sales > 0 else ekb_median * 2.0
velocity_score = min(1.0, max(0.0, monthly_velocity / denominator))
# ── Step 5: confidence ───────────────────────────────────────────────────
mapped_conf: Literal["high", "medium", "low"]
if n_comps >= 10 and months_observed >= 5:
mapped_conf = "high"
elif n_comps >= 5 and months_observed >= 3:
mapped_conf = "medium"
else:
mapped_conf = "low"
# ── Step 6: top-5 конкурентов по объёму продаж ───────────────────────────
sales_by_id: dict[int, float] = {
int(r["obj_id"]): float(r["total_sqm"] or 0.0) for r in mapped_sales_rows
}
sample = sorted(
[
{
"obj_id": oid,
**competitor_meta[oid],
"total_sqm_period": round(sales_by_id.get(oid, 0.0), 0),
"by_room_bucket": per_comp_buckets.get(oid, {}),
}
for oid in obj_ids
if oid in competitor_meta
],
key=lambda x: x["total_sqm_period"],
reverse=True,
)[:5]
return VelocityResult(
competitors_count=n_comps,
monthly_velocity_sqm=monthly_velocity,
ekb_median_sqm=ekb_median,
velocity_score=velocity_score,
confidence=mapped_conf,
months_observed=months_observed,
period_start=period_start,
period_end=period_end,
sample_competitors=sample,
by_room_bucket=by_room_bucket,
velocity_data_available=True,
)
def _get_ekb_median(db: Session, months_window: int = 6) -> float | None:
"""ЕКБ-wide медиана monthly velocity (м²/мес) per ЖК — benchmark.
Источник: objective_corpus_room_month (актуальные данные, обновляется еженедельно).
Ранее использовался domrf_kn_sale_graph с фильтром region_cd=66.
objective_corpus_room_month не имеет region_cd — данные Objective'а
покрывают primarily ЕКБ, что для baseline допустимо.
Учитываются только ЖК с ≥3 месяцами данных за окно (стабильный сигнал).
Fallback к _EKB_MEDIAN_FALLBACK_SQM_PER_MONTH если нет данных в БД.
"""
try:
with db.begin_nested():
row = (
db.execute(
text(
"""
WITH per_project AS (
SELECT
project_name,
SUM(COALESCE(deals_total_vol_m2,
deals_total_count * 45.0)) AS total_sqm,
COUNT(DISTINCT report_month) AS months_data
FROM objective_corpus_room_month
WHERE report_month >= (CURRENT_DATE - CAST(:window_interval AS interval))
AND deals_total_count > 0
GROUP BY project_name
HAVING COUNT(DISTINCT report_month) >= 3
),
per_project_velocity AS (
SELECT total_sqm / months_data AS velocity
FROM per_project
)
SELECT PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY velocity) AS median
FROM per_project_velocity
"""
),
{"window_interval": f"{months_window} months"},
)
.mappings()
.first()
)
except Exception:
logger.warning("velocity: ekb_median query failed, using fallback")
# SAVEPOINT auto-rollback'нут
return None
if row and row["median"] is not None:
return float(row["median"])
return None