gendesign/backend/app/workers/tasks/supply_layers_refresh.py
bot-backend 6982255fb3
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
fix(workers): убрать max_retries, который ничего не делает (#2464)
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>
2026-08-20 17:17:38 +05:00

230 lines
11 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.

"""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()