fix(workers): убрать max_retries, который ничего не делает (#2464) #2977
12 changed files with 201 additions and 31 deletions
|
|
@ -110,13 +110,22 @@ def _upsert_inflation(db: Session, rows: list[tuple[date, Decimal]]) -> int:
|
|||
return upserted
|
||||
|
||||
|
||||
# Ретраев здесь НЕТ намеренно, и параметров, обещающих их, тоже быть не должно
|
||||
# (#2464). Раньше стояло `bind=True, max_retries=2` — но self не использовался,
|
||||
# self.retry() не вызывался и autoretry_for задан не был, поэтому конфигурация
|
||||
# ретраев не имела ни малейшего эффекта: таска падала окончательно с первой ошибки,
|
||||
# а параметр обещал до двух повторов. Соседи, где ретраи действительно нужны,
|
||||
# задают их явно: autoretry_for в nspd_sync и scrape_cadastre, self.retry() в
|
||||
# scrape_kn.
|
||||
#
|
||||
# Отсутствие ретраев — это и есть задуманное поведение, оно описано в докстринге
|
||||
# ниже: «первая возникшая ошибка пробрасывается в конце (surfaces в
|
||||
# Celery/GlitchTip), не глотается». Ряды тянутся по расписанию, следующий тик
|
||||
# повторит попытку; молча ретраить внутри тика значило бы прятать отказ источника.
|
||||
@celery_app.task(
|
||||
bind=True,
|
||||
name="tasks.cbr_macro_sync.cbr_macro_sync",
|
||||
max_retries=2,
|
||||
)
|
||||
def cbr_macro_sync(
|
||||
self: Any,
|
||||
from_date: str | None = None,
|
||||
to_date: str | None = None,
|
||||
) -> dict[str, Any]:
|
||||
|
|
|
|||
|
|
@ -28,12 +28,15 @@ from app.workers.celery_app import celery_app
|
|||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло
|
||||
# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и
|
||||
# autoretry_for задан не был — конфигурация не имела эффекта. Соседи, где ретраи
|
||||
# нужны, задают их явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry()
|
||||
# (scrape_kn).
|
||||
@celery_app.task(
|
||||
bind=True,
|
||||
name="tasks.developer_registry_refresh.refresh_developer_registry",
|
||||
max_retries=2,
|
||||
)
|
||||
def refresh_developer_registry(self: Any) -> dict[str, Any]:
|
||||
def refresh_developer_registry() -> dict[str, Any]:
|
||||
"""REFRESH MATERIALIZED VIEW CONCURRENTLY developer_registry.
|
||||
|
||||
Лёгкая задача (реестр ~1024 застройщика). CONCURRENTLY — non-blocking для
|
||||
|
|
|
|||
|
|
@ -30,12 +30,15 @@ from app.workers.celery_app import celery_app
|
|||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло
|
||||
# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и
|
||||
# autoretry_for задан не был — конфигурация не имела эффекта. Соседи, где ретраи
|
||||
# нужны, задают их явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry()
|
||||
# (scrape_kn).
|
||||
@celery_app.task(
|
||||
bind=True,
|
||||
name="tasks.location_refresh.location_refresh",
|
||||
max_retries=2,
|
||||
)
|
||||
def location_refresh(self: Any, region: str | None = None) -> dict[str, Any]:
|
||||
def location_refresh(region: str | None = None) -> dict[str, Any]:
|
||||
"""Пересчитать + upsert-нуть district-level индексы по всем районам в `location`.
|
||||
|
||||
Идемпотентно (ON CONFLICT по district_name). Graceful: сбойный район
|
||||
|
|
|
|||
|
|
@ -26,12 +26,15 @@ from app.workers.celery_app import celery_app
|
|||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло
|
||||
# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и
|
||||
# autoretry_for задан не был — конфигурация не имела эффекта. Соседи, где ретраи
|
||||
# нужны, задают их явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry()
|
||||
# (scrape_kn).
|
||||
@celery_app.task(
|
||||
bind=True,
|
||||
name="tasks.mv_sales_tracker_refresh.refresh_sales_tracker_mvs",
|
||||
max_retries=2,
|
||||
)
|
||||
def refresh_sales_tracker_mvs_task(self: Any) -> dict[str, Any]:
|
||||
def refresh_sales_tracker_mvs_task() -> dict[str, Any]:
|
||||
"""REFRESH both sales-tracker MVs (#61).
|
||||
|
||||
Both MVs are refreshed CONCURRENTLY (non-blocking, require their UNIQUE
|
||||
|
|
|
|||
|
|
@ -17,13 +17,16 @@ from app.workers.celery_app import celery_app
|
|||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#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(
|
||||
bind=True,
|
||||
name="tasks.refresh_analytics.refresh_ekb_districts_medians",
|
||||
max_retries=2,
|
||||
)
|
||||
def refresh_ekb_districts_medians(
|
||||
self: Any,
|
||||
window_months: int = 24,
|
||||
min_deals: int = 50,
|
||||
) -> dict[str, Any]:
|
||||
|
|
|
|||
|
|
@ -22,12 +22,16 @@ from app.workers.celery_app import celery_app
|
|||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#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(
|
||||
bind=True,
|
||||
name="tasks.refresh_layout_velocity.refresh_layout_velocity",
|
||||
max_retries=2,
|
||||
)
|
||||
def refresh_layout_velocity_task(self: Any) -> dict[str, Any]:
|
||||
def refresh_layout_velocity_task() -> dict[str, Any]:
|
||||
"""REFRESH MATERIALIZED VIEW mv_layout_velocity (best_layouts, #113 / #1666).
|
||||
|
||||
MV рефрешится CONCURRENTLY (non-blocking, требует unique-индекс
|
||||
|
|
|
|||
|
|
@ -20,12 +20,16 @@ from app.workers.celery_app import celery_app
|
|||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#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(
|
||||
bind=True,
|
||||
name="tasks.refresh_quarter_price_index.refresh_quarter_price_index_chain",
|
||||
max_retries=2,
|
||||
)
|
||||
def refresh_quarter_price_index_chain(self: Any) -> dict[str, Any]:
|
||||
def refresh_quarter_price_index_chain() -> dict[str, Any]:
|
||||
"""Refresh mv_quarter_price_per_m2 then mv_quarter_price_index in sequence.
|
||||
|
||||
Both MVs are refreshed CONCURRENTLY (non-blocking). Falls back to
|
||||
|
|
|
|||
|
|
@ -231,12 +231,22 @@ def _upsert_emiss_rows(db: Session, rows: list[EmissRow]) -> int:
|
|||
return upserted
|
||||
|
||||
|
||||
# Ретраев здесь НЕТ намеренно, и параметров, обещающих их, тоже быть не должно
|
||||
# (#2464). Раньше стояло `bind=True, max_retries=2` — но self не использовался,
|
||||
# self.retry() не вызывался и autoretry_for задан не был, поэтому конфигурация
|
||||
# ретраев не имела ни малейшего эффекта: таска падала окончательно с первой ошибки,
|
||||
# а параметр обещал до двух повторов. Соседи, где ретраи действительно нужны,
|
||||
# задают их явно: autoretry_for в nspd_sync и scrape_cadastre, self.retry() в
|
||||
# scrape_kn.
|
||||
#
|
||||
# Отсутствие ретраев — это и есть задуманное поведение, оно описано в докстринге
|
||||
# ниже: «первая возникшая ошибка пробрасывается в конце (surfaces в
|
||||
# Celery/GlitchTip), не глотается». Ряды тянутся по расписанию, следующий тик
|
||||
# повторит попытку; молча ретраить внутри тика значило бы прятать отказ источника.
|
||||
@celery_app.task(
|
||||
bind=True,
|
||||
name="tasks.rosstat_macro_sync.rosstat_macro_sync",
|
||||
max_retries=2,
|
||||
)
|
||||
def rosstat_macro_sync(self: Any) -> dict[str, Any]:
|
||||
def rosstat_macro_sync() -> dict[str, Any]:
|
||||
"""Загрузить ряды Росстата (open-data + ЕМИСС + xlsx-СМР) и апсертить в macro_indicator.
|
||||
|
||||
Источники выполняются НЕЗАВИСИМО (per-source try/except): сбой одного источника
|
||||
|
|
|
|||
|
|
@ -157,7 +157,10 @@ def _region_lock(region_code: int, developers: list[str] | None) -> Iterator[boo
|
|||
logger.warning("release lock %s failed: %s", key, e)
|
||||
|
||||
|
||||
@celery_app.task(bind=True, name="tasks.scrape_kn.scrape_kn_region", max_retries=2)
|
||||
# bind=True здесь настоящий: self.request.id пишется в kn_scrape_log. А вот
|
||||
# max_retries=2 был инертен — self.retry() в этой таске не вызывается и
|
||||
# autoretry_for не задан (#2464). self.retry() ниже принадлежит resume_kn_run.
|
||||
@celery_app.task(bind=True, name="tasks.scrape_kn.scrape_kn_region")
|
||||
def scrape_kn_region(
|
||||
self: Any,
|
||||
region_code: int,
|
||||
|
|
|
|||
|
|
@ -146,13 +146,16 @@ def _save_raw(
|
|||
return int(row)
|
||||
|
||||
|
||||
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#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(
|
||||
bind=True,
|
||||
name="tasks.scrape_objective.sync_objective_group",
|
||||
max_retries=2,
|
||||
)
|
||||
def sync_objective_group(
|
||||
self: Any,
|
||||
group_name: str | None = None,
|
||||
triggered_by: str = "beat",
|
||||
use_ddu: bool = True,
|
||||
|
|
|
|||
|
|
@ -155,12 +155,16 @@ def _upsert_rows(db: Session, rows: list[SupplyLayerRow]) -> tuple[int, int]:
|
|||
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(
|
||||
bind=True,
|
||||
name="tasks.supply_layers_refresh.supply_layers_refresh",
|
||||
max_retries=2,
|
||||
)
|
||||
def supply_layers_refresh(self: Any) -> dict[str, Any]:
|
||||
def supply_layers_refresh() -> dict[str, Any]:
|
||||
"""Пересчитать 3-слойный склад предложения по всем районам и UPSERT в supply_layers.
|
||||
|
||||
Coherent snapshot: ``run_date = date.today()`` фиксируется ОДИН раз и передаётся
|
||||
|
|
|
|||
121
backend/tests/workers/test_2464_retry_config_is_real.py
Normal file
121
backend/tests/workers/test_2464_retry_config_is_real.py
Normal file
|
|
@ -0,0 +1,121 @@
|
|||
"""Таска, объявившая max_retries, обязана ретраи РЕАЛИЗОВАТЬ (#2464).
|
||||
|
||||
`cbr_macro_sync` и `rosstat_macro_sync` были объявлены как `bind=True, max_retries=2`, но
|
||||
`self` не использовался, `self.retry()` не вызывался и `autoretry_for` задан не был. То
|
||||
есть конфигурация ретраев не имела ни малейшего эффекта: таска падала окончательно с первой
|
||||
ошибки, а параметр обещал до двух повторов. Ни красного, ни ошибки — параметр просто
|
||||
декорация, и отличить её от работающей настройки можно только чтением тела.
|
||||
|
||||
Гейт разбирает КАЖДЫЙ декоратор `@celery_app.task(...)` через ast: если в нём есть
|
||||
`max_retries`, то либо в том же декораторе должен стоять `autoretry_for`, либо в теле
|
||||
функции — вызов `self.retry(...)`.
|
||||
|
||||
Оговорка про силу проверки: гейт СТРУКТУРНЫЙ, он читает исходники. Поведенческим его
|
||||
сделать нельзя — инертная конфигурация по определению ничего не меняет в поведении, и
|
||||
поймать её можно только по несоответствию объявления и кода.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import ast
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
_TASKS_DIR = Path(__file__).resolve().parents[2] / "app" / "workers"
|
||||
|
||||
|
||||
def _task_decorators(tree: ast.AST):
|
||||
"""(имя функции, декоратор celery-таски, узел функции) для каждой таски модуля."""
|
||||
for node in ast.walk(tree):
|
||||
if not isinstance(node, ast.FunctionDef | ast.AsyncFunctionDef):
|
||||
continue
|
||||
for dec in node.decorator_list:
|
||||
if not isinstance(dec, ast.Call):
|
||||
continue
|
||||
target = dec.func
|
||||
name = getattr(target, "attr", None) or getattr(target, "id", None)
|
||||
if name == "task":
|
||||
yield node.name, dec, node
|
||||
|
||||
|
||||
def _calls_self_retry(fn: ast.AST) -> bool:
|
||||
for node in ast.walk(fn):
|
||||
if isinstance(node, ast.Call):
|
||||
f = node.func
|
||||
if isinstance(f, ast.Attribute) and f.attr == "retry":
|
||||
owner = getattr(f.value, "id", None)
|
||||
if owner == "self":
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def _offenders() -> list[str]:
|
||||
bad: list[str] = []
|
||||
for path in sorted(_TASKS_DIR.rglob("*.py")):
|
||||
tree = ast.parse(path.read_text())
|
||||
for fn_name, dec, fn_node in _task_decorators(tree):
|
||||
kwargs = {kw.arg: kw.value for kw in dec.keywords if kw.arg}
|
||||
if "max_retries" not in kwargs:
|
||||
continue
|
||||
# max_retries=0 — ЧЕСТНОЕ объявление «ретраев нет», а не пустое обещание.
|
||||
# Такие места в репозитории есть и снабжены пояснением (nspd_geo:
|
||||
# «resume через worker_ready, не через retry»; objective_etl: «при сбое
|
||||
# лучше человек посмотрит»). Флагуем только ненулевое обещание.
|
||||
v = kwargs["max_retries"]
|
||||
if isinstance(v, ast.Constant) and v.value == 0:
|
||||
continue
|
||||
if "autoretry_for" in kwargs or _calls_self_retry(fn_node):
|
||||
continue
|
||||
bad.append(f"{path.relative_to(_TASKS_DIR.parent.parent)}::{fn_name}")
|
||||
return bad
|
||||
|
||||
|
||||
def test_declared_retries_are_actually_implemented() -> None:
|
||||
"""max_retries без autoretry_for и без self.retry() — обещание без исполнения."""
|
||||
bad = _offenders()
|
||||
assert not bad, (
|
||||
"таски объявляют max_retries, но ретраи не реализуют — параметр не имеет эффекта:\n "
|
||||
+ "\n ".join(bad)
|
||||
+ "\nЛибо задай autoretry_for / вызови self.retry(), либо убери max_retries."
|
||||
)
|
||||
|
||||
|
||||
def test_gate_sees_the_tasks_at_all() -> None:
|
||||
"""Контроль на сам гейт: он должен что-то находить.
|
||||
|
||||
Если разбор перестанет узнавать декораторы (переименуют celery_app, сменят
|
||||
обёртку), проверка выше станет тавтологически зелёной.
|
||||
"""
|
||||
found = 0
|
||||
for path in sorted(_TASKS_DIR.rglob("*.py")):
|
||||
found += sum(1 for _ in _task_decorators(ast.parse(path.read_text())))
|
||||
assert found >= 20, f"гейт нашёл всего {found} celery-тасок — разбор сломался"
|
||||
|
||||
|
||||
def test_gate_recognises_both_valid_forms() -> None:
|
||||
"""Контроль: обе законные формы ретраев распознаются, а не только одна.
|
||||
|
||||
Иначе гейт краснел бы на исправных тасках и его бы отключили.
|
||||
"""
|
||||
src_autoretry = (
|
||||
"@celery_app.task(bind=True, max_retries=2, autoretry_for=(ValueError,))\n"
|
||||
"def t(self):\n return 1\n"
|
||||
)
|
||||
src_self_retry = (
|
||||
"@celery_app.task(bind=True, max_retries=2)\n"
|
||||
"def t(self):\n raise self.retry(countdown=1)\n"
|
||||
)
|
||||
src_inert = "@celery_app.task(bind=True, max_retries=2)\ndef t(self):\n return 1\n"
|
||||
|
||||
def _is_ok(src: str) -> bool:
|
||||
tree = ast.parse(src)
|
||||
for _n, dec, fn in _task_decorators(tree):
|
||||
kwargs = {kw.arg for kw in dec.keywords if kw.arg}
|
||||
return "autoretry_for" in kwargs or _calls_self_retry(fn)
|
||||
raise AssertionError("декоратор не распознан")
|
||||
|
||||
assert _is_ok(src_autoretry), "форма autoretry_for не распознана"
|
||||
assert _is_ok(src_self_retry), "форма self.retry() не распознана"
|
||||
assert not _is_ok(src_inert), "инертная форма ошибочно признана исправной"
|
||||
Loading…
Add table
Reference in a new issue