gendesign/tradein-mvp/backend/app/core/db.py
bot-backend a8503d3a58
All checks were successful
CI Trade-In / backend-tests (pull_request) Successful in 5m9s
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 12s
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
tradein: потолок на запрос и на ожидание блокировки для движка БД (#3463)
На боевой БД `statement_timeout`, `lock_timeout` и
`idle_in_transaction_session_timeout` равны 0, а у движка `app/core/db.py` не
было `connect_args` вовсе. После #3444/#3449 шаги БД на пути `/estimate` идут
через обёртку, которая при отмене по бюджету ДОЖИДАЕТСЯ своего потока (иначе он
остаётся сиротой в общей `Session`) — ожидание верное, но его верхняя граница
равна длительности самого запроса, а у запроса границы не было. Один
`ACCESS EXCLUSIVE` на таблице → четыре повисших запроса → `_ESTIMATE_CONCURRENCY`
исчерпан → `/estimate` отдаёт 429 всем остальным.

Потолок ставится на КОННЕКТЕ (libpq `options`), а не в обёртке: таймаут в
обёртке вернул бы ровно ту сироту, ради которой писался #3449.

statement_timeout = 30 с: выше самого длинного ОБЪЯВЛЕННОГО бюджета `/estimate`
(20 с, `estimate_avito_imv_timeout_s`) в 1.5 раза и в 7 раз выше самого долгого
ЗАМЕРЕННОГО запроса через этот движок (4.27 с, `pg_stat_statements` на проде за
16 суток), но конечен. lock_timeout = 5 с: та же величина, что у миграций
проекта, и больше `deadlock_timeout` (1 с на проде).

`idle_in_transaction_session_timeout` намеренно не трогаем: тем же движком живёт
планировщик, а его свипы держат транзакцию открытой всё время внешнего HTTP
(замер: живая сессия `idle in transaction` 29 с).

Задачи планировщика проверены, а не предположены: у `listing_source_snapshot`
свой `SET LOCAL statement_timeout = 900000`, и тест доказывает, что `SET LOCAL`
ПЕРЕКРЫВАЕТ сессионный потолок и не течёт за свою транзакцию. Самая долгая
чисто-БД задача по `scrape_runs` за 14 суток укладывается в 9.7 с целиком;
единственный запрос длиннее 20 с на всей БД (`REFRESH MATERIALIZED VIEW
CONCURRENTLY`, 30.85 с) идёт мимо движка — по своему сырому psycopg-соединению.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-12 19:46:57 +05:00

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

import asyncio
import logging
from collections.abc import Callable, Generator
from typing import Any
from sqlalchemy import create_engine
from sqlalchemy.orm import DeclarativeBase, Session, sessionmaker
from app.core.config import settings
logger = logging.getLogger(__name__)
# #3463. Потолок ОДНОГО statement'а, секунды·1000. Ставится на КОННЕКТЕ (libpq
# `options`), а не в питоновской обёртке: обёртка (`run_db_thread` ниже) при
# отмене обязана ДОЖДАТЬСЯ потока, иначе поток остаётся сиротой в общей
# `Session` — ровно то, ради чего писался #3449. Значит верхняя граница ожидания
# = длительность самого запроса, и задать её может только сервер.
#
# 30 с выбраны так, чтобы потолок НИКОГДА не стал биндящим ограничением для
# честной работы, но остался конечным:
# * самый длинный ОБЪЯВЛЕННЫЙ бюджет на `/estimate` — 20 с (`estimate_avito_imv_timeout_s`,
# config.py:852); дальше 12 с геокод, 8 с Yandex/Cian/house_meta. 30 с = 1.5× от максимума;
# * замер на боевой БД 12.09 (`pg_stat_statements`, накоплено с 27.08): САМЫЙ долгий
# запрос через этот движок — 4.27 с; всего 2 запроса длиннее 5 с за 16 суток, и оба
# мимо движка (`REFRESH MATERIALIZED VIEW CONCURRENTLY` 30.85 с — своё сырое
# psycopg-соединение в tasks/refresh_search_matview.py; `COPY _stg` 19.67 с — psql
# из scripts/local-avito-msk/collect.py);
# * самая долгая ЧИСТО-БД задача планировщика по `scrape_runs` за 14 суток —
# listing_source_snapshot, 9.7 с ЦЕЛИКОМ (и у неё сверх того свой
# `SET LOCAL statement_timeout = 900000`, который перекрывает это значение —
# гейт tests/test_3463_db_timeouts.py).
_STATEMENT_TIMEOUT_MS = 30_000
# Ожидание БЛОКИРОВКИ — заведомо меньше: ждать лок дольше секунд смысла нет, лучше
# деградировать. 5 с — та же величина, что у миграций проекта
# (`SET LOCAL lock_timeout = '5s'` в data/sql/250,251,260,272,277…), снизу ограничена
# deadlock_timeout (на проде 1 с — сверено 12.09). Именно этот потолок закрывает
# сценарий #3463: под `ACCESS EXCLUSIVE` на `geocode_cache` запрос ЖДЁТ лок, а не
# считает, — statement_timeout тут только страховка от «считает вечно».
_LOCK_TIMEOUT_MS = 5_000
# idle_in_transaction_session_timeout НАМЕРЕННО не трогаем: тот же движок обслуживает
# и планировщик, а его свипы держат транзакцию открытой всё время внешнего HTTP
# (замер на проде 12.09: живая сессия `idle in transaction` 29 с). Сессионный потолок
# на простой в транзакции убивал бы рабочий сбор, а не зависший запрос.
DB_CONNECT_ARGS = {
"options": f"-c statement_timeout={_STATEMENT_TIMEOUT_MS} -c lock_timeout={_LOCK_TIMEOUT_MS}"
}
engine = create_engine(
settings.database_url,
pool_pre_ping=True,
future=True,
# #3463. Накрывает ВСЕ три сервиса образа (backend / scraper / tgbot — один и тот
# же `app.core.db`, см. docker-compose.prod.yml) и обе стороны: продуктовый путь
# `/estimate` и задачи планировщика. Миграции идут мимо (psql из
# .forgejo/workflows/deploy-tradein.yml, не этот движок) — их DDL под своим
# `SET LOCAL lock_timeout` и потолком не ограничен.
connect_args=DB_CONNECT_ARGS,
# #3194: SQLAlchemy печатает ВСЕ bind-параметры в тексте StatementError —
# через них в GlitchTip уезжали ключ шифрования кук и сами куки
# (pgp_sym_encrypt(:cookies_json, :key)). Флаг на УРОВНЕ ДВИЖКА кроет все
# сайты вызова разом, включая будущие.
# НЕ закрывает: текст ошибки самого драйвера (Postgres DETAIL со значением)
# и сырые psycopg-подключения мимо движков — это отдельный класс.
hide_parameters=True,
# #3408 п.1. Потолок пула обязан быть НЕ МЕНЬШЕ суммы потолков одновременности,
# которые сам же процесс и объявляет: 4 оценки (`api/v1/trade_in`
# `_ESTIMATE_CONCURRENCY`) + 4 подсказки (`api/public/mera` `_SUGGEST_CONCURRENCY`)
# + 8 фоновых догрузок (`services/estimator` `_MAX_DEFERRED_REFRESH_TASKS`) = 16.
# Дефолт SQLAlchemy 5+10=15 меньше этой суммы, то есть исчерпать пул можно
# штатной работой, не абузом. Гейт — tests/test_3408_pool_ceiling.py.
#
# Инвариант «пул >= суммы объявленных потолков» — это ПОЛ, а не гарантия:
# сумма считает по одному коннекту на объявленную единицу работы, а коннект
# держит и всё остальное — любая ручка с `Depends(get_db)` тоже. Глобального
# обработчика `sqlalchemy.exc.TimeoutError` в app/main.py нет, поэтому на
# исчерпанном пуле соседние ручки отдают 500 (после ожидания `pool_timeout`),
# а не медленный 200.
#
# pool_size оставлен дефолтным (5): это ПОСТОЯННО открытые коннекты, а в покое
# прод держит 5-6 (замер 11.09). Растёт только overflow — коннекты пика,
# которые пул закрывает сам. Потолок процесса: 5 + 15 = 20; воркер один
# (docker-compose.prod.yml, uvicorn без --workers), Postgres max_connections=100.
max_overflow=15,
# #3408 п.2. Дефолтные 30 с ожидания коннекта длиннее ЛЮБОГО бюджета внешнего
# источника в эстиматоре: 8 с (Yandex / Cian / house_meta), 12 с (geocode),
# 20 с (Avito IMV — `estimate_avito_imv_timeout_s`, config.py:852). Сам чекаут
# прервать `asyncio.wait_for` не может: он занимает поток `asyncio.to_thread`
# целиком, а пул потоков конечен (min(32, cpu+4)) — исчерпанный пул коннектов
# так превращается в исчерпанный пул потоков. 5 с короче самого КОРОТКОГО
# бюджета: занятый пул деградирует ОДИН источник, а не весь запрос.
#
# Отдельный коммит в конце ветки намеренно (ревью PR #3444): это единственная
# правка, которая меняет режим отказа с «медленно» на «быстро с ошибкой», и
# едет она во ВСЕ сервисы образа — backend, scraper, tgbot
# (tradein-mvp/docker-compose.prod.yml). За 29 ч логов исчерпания пула не было
# ни разу, то есть новое значение на проде пока не на чем проверить.
# ТРИГГЕР ОТКАТА на 30 с: любое `QueuePool limit ... timed out` в логах
# бэкенда ЛИБО рост failed+zombie в `scrape_runs` после деплоя.
pool_timeout=5,
)
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine, expire_on_commit=False)
class Base(DeclarativeBase):
pass
def get_db() -> Generator[Session, None, None]:
db = SessionLocal()
try:
yield db
finally:
db.close()
async def run_db_thread[T](fn: Callable[..., T], *args: Any, **kwargs: Any) -> T:
"""`asyncio.to_thread(fn, ...)`, который при отмене ДОЖИДАЕТСЯ своего потока (#3449).
Для синхронной работы по сессии, которая шагу НЕ принадлежит — она общая со
всем остальным запросом (`Depends(get_db)`). Поток отменить нельзя: у
`asyncio.to_thread` отменяется только ожидание со стороны loop'а. Корутина
умирает по бюджету источника (`estimator._with_budget` = `asyncio.wait_for`,
геокодер — 12 с), а поток продолжает работать с ТОЙ ЖЕ `Session`, пока
вызывающий уже идёт дальше по коду — следующий источник,
`_fetch_anchor_comps`, `_persist_estimate_and_commit`. Два потока в одной
`Session` дают «another operation is in progress» / `InvalidRequestError` на
СЛЕДУЮЩЕМ шаге: у источников такую ошибку глушит `except` вокруг вызова, у
персиста оценки не глушит ничего — 500 и потерянная оценка клиента.
Поэтому отмена пробрасывается ПОСЛЕ того, как поток отпустил сессию. Цена —
бюджет источника переезжает на длину ОДНОГО шага БД (чекаут ≤ `pool_timeout`
плюс сам запрос), а не на длину фетча, ради которой бюджет и заведён.
Только защита от сироты: транзакцию шаг НЕ завершает (в середине геокодинга
commit зафиксировал бы частичное состояние оценки). Кому нужен ещё и возврат
коннекта в пул перед внешним HTTP — `estimator._db_step`, он поверх этого.
Гейт — tests/test_3449_geocoder_cancel_orphan.py.
"""
step = asyncio.ensure_future(asyncio.to_thread(fn, *args, **kwargs))
try:
return await asyncio.shield(step)
except asyncio.CancelledError:
await asyncio.wait([step])
if not step.cancelled() and step.exception() is not None:
# Результата уже никто не ждёт: без явного чтения asyncio напечатает
# «Task exception was never retrieved» вообще без контекста.
logger.warning("шаг БД упал уже после отмены: %s", step.exception())
raise