All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 9s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Successful in 1m54s
CI / backend-tests (pull_request) Successful in 17m0s
11 тасок объявляли `max_retries=2`, но ретраи не реализовывали: ни `autoretry_for` в декораторе, ни вызова `self.retry()` в теле. Celery в таком виде параметр не применяет — при исключении таска падает с первой попытки. Читающий код видит «до 3 попыток», а их одна. Убран `max_retries` у: cbr_macro_sync, rosstat_macro_sync, developer_registry_refresh, location_refresh, mv_sales_tracker_refresh, refresh_analytics, refresh_layout_velocity, refresh_quarter_price_index, scrape_objective.sync_objective_group, supply_layers_refresh, scrape_kn.scrape_kn_region. Заодно убран `bind=True` там, где `self` не использовался вовсе; в `scrape_kn_region` он оставлен — `self.request.id` пишется в kn_scrape_log. Не тронуты и не должны быть: `resume_kn_run` (max_retries=12 + настоящий self.retry()), `nspd_sync`/`scrape_cadastre` (autoretry_for), `nspd_geo`/`objective_etl` (max_retries=0 — честное «ретраев нет»). Гейт `test_2464_retry_config_is_real.py` разбирает AST всех модулей `app/workers/tasks/` и требует: если декоратор объявляет ненулевой max_retries, в нём есть autoretry_for либо в теле функции есть self.retry(). Три таски из одиннадцати гейт нашёл сверх списка эпика. Проверка гейта: с фиксом зелено, при возврате `max_retries=2` в supply_layers_refresh — красно с указанием на эту таску. Плюс два контроля: гейт видит ≥20 тасок (не молчит из-за пустой выборки) и признаёт обе законные формы ретраев. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
230 lines
11 KiB
Python
230 lines
11 KiB
Python
"""Celery task: пересчёт 3-слойного склада предложения → ``supply_layers`` (м.125).
|
||
|
||
#950 (Site Finder v2 / GG-форсайт, EPIC 6 «3-layer склад предложения», ТЗ §9.3),
|
||
Step 5. Для каждого района Свердловской обл. вызывает детерминированный сервис
|
||
``supply_layers.compute_all_layers`` (L1 open / L2 hidden / L3 future) и UPSERT-ит
|
||
полученные ``SupplyLayerRow`` в таблицу ``supply_layers`` по логическому ключу
|
||
``uq_supply_layers_logical`` (layer, district_name, complex_id, obj_class,
|
||
dev_group_name, source, snapshot_date) — dev_group_name в ключе с м.128 (#970), иначе
|
||
dev_group'ы района схлопывались (~95.6% потерь L3).
|
||
|
||
Детерминированно, без LLM — сервис read-only считает агрегаты, этот worker только
|
||
пишет. ON CONFLICT делает повторные прогоны идемпотентными в рамках одного дня
|
||
(re-run того же ``snapshot_date`` обновляет мутабельные колонки + computed_at).
|
||
|
||
Контракт записи (КРИТИЧНО, флаг database-expert): на конфликте обновляем ВСЕ
|
||
мутабельные колонки (units/area/expected_online_date/confidence/method/computed_at),
|
||
иначе при ре-прогоне в таблице останутся stale-значения. Ключевые колонки (layer,
|
||
district_name, complex_id, obj_class, dev_group_name, source, snapshot_date) НЕ трогаем.
|
||
|
||
Mirror conventions (cbr_macro_sync.py / refresh_quarter_price_index.py):
|
||
• ``SessionLocal()`` + try/finally close, ``logger`` (не print).
|
||
• SAVEPOINT per-row (``with db.begin_nested():``, backend.md) — одна битая строка
|
||
(напр. CHECK-violation) не откатывает весь батч; commit один раз в конце.
|
||
• ``CAST(:x AS type)`` — НИКОГДА ``:x::type`` (psycopg v3).
|
||
• Сбой одного района (сервис уже degrade'ит до [], но страхуемся) не валит батч.
|
||
|
||
Coherent snapshot: ``run_date = date.today()`` вычисляется ОДИН раз в начале и
|
||
передаётся как ``snapshot_date`` КАЖДОЙ строке всех районов — весь батч делит один
|
||
снапшот (нет split на midnight-boundary; ON CONFLICT идемпотентен для same-day
|
||
re-run). Расписание — еженедельно, регистрируется в beat_schedule.py
|
||
(supply-layers-refresh-weekly).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from datetime import date
|
||
from typing import Any
|
||
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.db import SessionLocal
|
||
from app.services.site_finder.supply_layers import (
|
||
SupplyLayerRow,
|
||
compute_all_layers,
|
||
)
|
||
from app.workers.celery_app import celery_app
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Регион данных (ЕКБ / Свердловская обл.). domrf_kn_objects.district_name заполнен
|
||
# spatial-join'ом только для region_cd=66 (м.57) — зеркало supply_layers._EKB_REGION_CODE.
|
||
_EKB_REGION_CODE: int = 66
|
||
|
||
# Множество районов для прогона = UNION distinct non-NULL районов из ОБОИХ источников
|
||
# (objective_lots.district + domrf_kn_objects.district_name WHERE region_cd=66). Не
|
||
# хардкодим 8 имён — робастно к новым районам. CAST(:x AS type) (psycopg v3).
|
||
_DISTRICTS_SQL = text(
|
||
"""
|
||
SELECT district_name FROM (
|
||
SELECT DISTINCT ol.district AS district_name
|
||
FROM objective_lots ol
|
||
WHERE ol.district IS NOT NULL
|
||
UNION
|
||
SELECT DISTINCT o.district_name AS district_name
|
||
FROM domrf_kn_objects o
|
||
WHERE o.district_name IS NOT NULL
|
||
AND o.region_cd = CAST(:region_cd AS integer)
|
||
) d
|
||
ORDER BY district_name
|
||
"""
|
||
)
|
||
|
||
# UPSERT в supply_layers (м.125 + м.128). Конфликт по uq_supply_layers_logical (NULLS
|
||
# NOT DISTINCT) → ON CONFLICT по точному column-list ключа. DO UPDATE обновляет ВСЕ
|
||
# мутабельные колонки + computed_at=now() (database-expert: иначе stale-данные).
|
||
# Ключевые колонки (layer/district_name/complex_id/obj_class/dev_group_name/source/
|
||
# snapshot_date) НЕ обновляются — это и есть ключ. dev_group_name добавлен в ключ м.128
|
||
# (#970): без него все dev_group'ы района схлопывались в один NULL-key (complex_id всегда
|
||
# NULL для L3) → UPSERT перезатирал их (~95.6% потерь L3). CAST(:x AS type) везде
|
||
# (psycopg v3, никогда ::).
|
||
_UPSERT_SQL = text(
|
||
"""
|
||
INSERT INTO supply_layers (
|
||
layer, district_name, complex_id, obj_class, dev_group_name,
|
||
units_estimate, area_estimate, expected_online_date,
|
||
source, confidence, method, snapshot_date
|
||
) VALUES (
|
||
CAST(:layer AS smallint),
|
||
CAST(:district_name AS text),
|
||
CAST(:complex_id AS bigint),
|
||
CAST(:obj_class AS text),
|
||
CAST(:dev_group_name AS text),
|
||
CAST(:units_estimate AS integer),
|
||
CAST(:area_estimate AS numeric),
|
||
CAST(:expected_online_date AS date),
|
||
CAST(:source AS text),
|
||
CAST(:confidence AS text),
|
||
CAST(:method AS text),
|
||
CAST(:snapshot_date AS date)
|
||
)
|
||
ON CONFLICT (
|
||
layer, district_name, complex_id, obj_class, dev_group_name, source, snapshot_date
|
||
)
|
||
DO UPDATE SET
|
||
units_estimate = EXCLUDED.units_estimate,
|
||
area_estimate = EXCLUDED.area_estimate,
|
||
expected_online_date = EXCLUDED.expected_online_date,
|
||
confidence = EXCLUDED.confidence,
|
||
method = EXCLUDED.method,
|
||
computed_at = now()
|
||
"""
|
||
)
|
||
|
||
|
||
def _list_districts(db: Session) -> list[str]:
|
||
"""Distinct non-NULL районы из objective_lots + domrf_kn_objects (region_cd=66).
|
||
|
||
UNION обоих источников — робастно к новым районам (не хардкодим 8 имён). На сбое
|
||
запроса логируем и возвращаем [] (батч завершится no-op, не упадёт).
|
||
"""
|
||
try:
|
||
rows = db.execute(_DISTRICTS_SQL, {"region_cd": _EKB_REGION_CODE}).mappings().all()
|
||
except Exception:
|
||
logger.exception("supply_layers_refresh: districts enumeration query failed")
|
||
return []
|
||
return [r["district_name"] for r in rows if r["district_name"]]
|
||
|
||
|
||
def _upsert_rows(db: Session, rows: list[SupplyLayerRow]) -> tuple[int, int]:
|
||
"""UPSERT строк слоёв в supply_layers. SAVEPOINT per-row (backend.md): одна битая
|
||
строка (напр. CHECK-violation) не откатывает весь батч.
|
||
|
||
Returns: (upserted, skipped) — upserted = успешные INSERT/UPDATE, skipped = строки,
|
||
чей UPSERT упал (залогированы warning'ом). Commit делает caller (один на батч).
|
||
"""
|
||
upserted = 0
|
||
skipped = 0
|
||
for row in rows:
|
||
params = row.as_dict()
|
||
try:
|
||
with db.begin_nested():
|
||
db.execute(_UPSERT_SQL, params)
|
||
upserted += 1
|
||
except Exception as e:
|
||
skipped += 1
|
||
logger.warning(
|
||
"supply_layers upsert failed (layer=%s district=%s source=%s): %s",
|
||
row.layer,
|
||
row.district_name,
|
||
row.source,
|
||
e,
|
||
)
|
||
return upserted, skipped
|
||
|
||
|
||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло
|
||
# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и
|
||
# autoretry_for задан не был — конфигурация не имела эффекта. Где ретраи нужны, они
|
||
# задаются явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry()
|
||
# (scrape_kn.resume_kn_run). Где их сознательно нет — пишется max_retries=0 с
|
||
# пояснением (nspd_geo, objective_etl.import_anton_objective).
|
||
@celery_app.task(
|
||
name="tasks.supply_layers_refresh.supply_layers_refresh",
|
||
)
|
||
def supply_layers_refresh() -> dict[str, Any]:
|
||
"""Пересчитать 3-слойный склад предложения по всем районам и UPSERT в supply_layers.
|
||
|
||
Coherent snapshot: ``run_date = date.today()`` фиксируется ОДИН раз и передаётся
|
||
как ``snapshot_date`` каждой строке всех районов (избегает midnight-split, делает
|
||
same-day re-run идемпотентным). Цикл по районам: сбой/пустота сервиса одного района
|
||
не валит остальные (сервис уже degrade'ит до [], плюс per-district guard). Commit
|
||
один раз в конце (outer tx); per-row SAVEPOINT изолирует битые строки.
|
||
|
||
Returns:
|
||
Счётчики ``{"snapshot_date", "districts", "rows", "upserted", "skipped"}``.
|
||
|
||
Raises:
|
||
Пробрасывает только инфраструктурные сбои (commit / нефатальные внутри уже
|
||
обработаны) — чтобы фейл был виден в Celery / GlitchTip, а не молча проглочен.
|
||
"""
|
||
run_date = date.today()
|
||
db = SessionLocal()
|
||
try:
|
||
districts = _list_districts(db)
|
||
logger.info(
|
||
"supply_layers_refresh: %d districts to process (snapshot_date=%s)",
|
||
len(districts),
|
||
run_date.isoformat(),
|
||
)
|
||
|
||
all_rows: list[SupplyLayerRow] = []
|
||
for district in districts:
|
||
try:
|
||
rows = compute_all_layers(db, district=district, snapshot_date=run_date)
|
||
except Exception:
|
||
# Сервис уже graceful (→ []), но страхуемся: сбой одного района не
|
||
# должен ронять весь батч. Логируем и продолжаем со следующим.
|
||
logger.exception(
|
||
"supply_layers_refresh: compute_all_layers failed for district=%s",
|
||
district,
|
||
)
|
||
continue
|
||
all_rows.extend(rows)
|
||
|
||
upserted, skipped = _upsert_rows(db, all_rows)
|
||
db.commit()
|
||
|
||
logger.info(
|
||
"supply_layers_refresh: done snapshot_date=%s districts=%d rows=%d "
|
||
"upserted=%d skipped=%d",
|
||
run_date.isoformat(),
|
||
len(districts),
|
||
len(all_rows),
|
||
upserted,
|
||
skipped,
|
||
)
|
||
return {
|
||
"snapshot_date": run_date.isoformat(),
|
||
"districts": len(districts),
|
||
"rows": len(all_rows),
|
||
"upserted": upserted,
|
||
"skipped": skipped,
|
||
}
|
||
except Exception as e:
|
||
logger.exception("supply_layers_refresh failed: %s", e)
|
||
raise
|
||
finally:
|
||
db.close()
|