Когда сервис получает SIGTERM, он должен остановиться аккуратно: дать запросам завершиться, сбросить буферы, закрыть соединения с базой. Если закрыть пул соединений слишком рано — фоновая задача, которая ещё не закончила commit, получит InterfaceError: connection already closed. Слишком поздно — соединения утекут и база будет держать их открытыми без причины.
Разберём, как всё устроено правильно.
engine.dispose() вызывается последним
Пул соединений — это общий ресурс, которым пользуются и HTTP-обработчики, и фоновые задачи. Закрыть пул нужно только тогда, когда все, кто им пользуется, уже завершили работу.
dispose закрывает соединения в пуле и не ждёт выданные; к этому моменту задачи и потребитель должны вернуть свои соединения.
FastAPI управляет жизненным циклом через lifespan — асинхронный контекстный менеджер. Shutdown-секция выполняется в том порядке, в котором написана:
from contextlib import asynccontextmanager
from fastapi import FastAPI
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker
engine = create_async_engine(
"postgresql+asyncpg://...",
pool_size=10,
max_overflow=5,
pool_pre_ping=True,
)
session_factory = async_sessionmaker(engine, expire_on_commit=False)
@asynccontextmanager
async def lifespan(app: FastAPI):
# startup
yield
# shutdown — порядок важен
await _stop_background_tasks() # 1. дожать фоновые задачи
await _stop_kafka() # 2. завершить producer/consumer
await engine.dispose() # 3. закрыть пул — последним
app = FastAPI(lifespan=lifespan)
engine.dispose() закрывает все соединения, которые лежат в пуле, и ничего не ждёт: соединение, которое кто-то ещё держит, остаётся у владельца и закроется, когда он его вернёт. Поэтому к этому моменту всё должно быть возвращено — фоновые задачи завершены на шаге 1, иначе их запросы упадут на закрытом движке.
Не добавляй отдельный signal.signal(SIGTERM, ...) для engine.dispose() — uvicorn уже вызывает lifespan-shutdown при получении SIGTERM, дублирование только запутает порядок.
Что происходит с активными транзакциями при SIGTERM
Разные виды кода завершают транзакции по-разному.
HTTP-обработчик
Когда HTTP-обработчик находится внутри транзакции в момент получения SIGTERM, uvicorn даёт ему время завершиться (опция --timeout-graceful-shutdown; по умолчанию ожидание не ограничено, поэтому её задают явно). Сессия, открытая через async with session_factory(), на выходе из блока закрывается: явный commit() фиксирует работу, а всё незафиксированное откатывается — и при исключении, и при отмене.
async def get_session() -> AsyncIterator[AsyncSession]:
async with session_factory() as session:
yield session
@router.post("/orders")
async def create_order(
body: CreateOrderRequest,
session: AsyncSession = Depends(get_session),
) -> OrderResponse:
order = Order(customer_id=body.customer_id, total=body.total)
session.add(order)
await session.commit()
return OrderResponse.model_validate(order)
Если таймаут истёк — uvicorn отменяет задачу запроса, CancelledError выходит из блока, и __aexit__ закрывает сессию с откатом.
Фоновая задача
Фоновые задачи завершаются через флаг готовности. Цикл должен проверять этот флаг, чтобы мягко остановиться после текущей итерации:
async def process_outbox():
while app_state.is_ready:
async with session_factory() as session:
async with session.begin():
rows = await session.execute(
select(OutboxEvent)
.where(OutboxEvent.status == "pending")
.limit(20)
.with_for_update(skip_locked=True)
)
batch = rows.scalars().all()
for event in batch:
await publish_event(event)
event.status = "published"
await asyncio.sleep(1)
При SIGTERM app_state.is_ready переходит в False, задача завершает текущую итерацию, транзакция делает commit, цикл выходит. Не используй while True — без флага задача не остановится.
Принудительная отмена с CancelledError
Если задача не успела завершиться сама, её отменяют через cancel():
async def _stop_background_tasks():
for task in _background_tasks:
task.cancel()
try:
await asyncio.wait_for(task, timeout=25.0)
except (asyncio.CancelledError, asyncio.TimeoutError):
pass
Внутри самой задачи при CancelledError транзакция откатится автоматически через __aexit__ контекстного менеджера — это нормально:
async def sync_product_catalog():
async with session_factory() as session:
async with session.begin():
try:
products = await fetch_external_catalog()
session.add_all(products)
except asyncio.CancelledError:
# транзакция откатится через __aexit__, это нормально
raise
Consumer Kafka
Если потребитель Kafka обрабатывает сообщения в транзакции, останов работает так: задачу потребителя отменяют и дожидаются, транзакция текущего сообщения коммитится, затем задача коммитит offset, и только после этого вызывают consumer.stop() — с ручным commit он сам ничего не фиксирует. Если SIGTERM застал посередине — транзакция откатывается, offset не коммитится, при следующем старте сообщение будет обработано повторно. Защиту от дублирования обеспечивает таблица обработанных событий:
async def consume_order_events():
async for msg in consumer:
event = json.loads(msg.value)
async with session_factory() as session:
async with session.begin():
inserted = await session.execute(
insert(ProcessedEvent).values(event_id=event["event_id"])
.on_conflict_do_nothing()
)
if inserted.rowcount == 1:
order = await session.get(Order, event["order_id"])
order.status = event["status"]
await consumer.commit()
Alembic запускается только при старте
Частое заблуждение: «нужно почистить схему при выходе» или «запустить alembic downgrade в shutdown». Это не паттерн.
Alembic применяет миграции при старте приложения и после этого просто закрывает свои соединения. При завершении Alembic ничего не делает — это правильно:
@asynccontextmanager
async def lifespan(app: FastAPI):
# startup: применить миграции
await run_migrations() # alembic upgrade head
yield
# shutdown: только закрытие ресурсов, никакого DDL
await engine.dispose()
Не добавляй никакой DDL в секцию shutdown.
Долгие транзакции мешают остановке и соседям
Транзакция, которая держит блокировку, при остановке становится проблемой дважды: она задерживает завершение пода и блокирует соседей, пока не закончится. Опаснее всего сочетание «долгая транзакция плюс блокировка строк»:
async def process_batch() -> None:
async with session_factory() as session, session.begin():
orders = await repo.lock_pending_for_update(session, 1000) # SELECT ... FOR UPDATE
for order in orders:
await external.notify(order) # сетевой вызов внутри транзакции
order.mark_notified()
Здесь тысяча строк заблокирована на всё время обхода, а внутри ещё и сетевой вызов, длительность которого вам не принадлежит. При остановке lifespan ждёт эту задачу, съедая бюджет; если не дождётся, отмена оборвёт транзакцию на середине.
Что делают, чтобы этого не было: пачки по сто строк и транзакция на пачку, а не на весь прогон; никаких внешних вызовов внутри транзакции, уведомление уходит после фиксации; проверка флага остановки между пачками, чтобы цикл остановился сам, зафиксировав сделанное; FOR UPDATE SKIP LOCKED вместо FOR UPDATE, чтобы остановка одного пода не тормозила остальных. Проверить, есть ли у вас такие транзакции, можно одним запросом под нагрузкой: транзакции длительностью больше нескольких секунд видны в pg_stat_activity по xact_start.
Если база недоступна в момент остановки
Под останавливается, а база в это время недоступна: упала, потеряна сеть, кончились соединения. Любой шаг остановки, который пытается получить соединение (отметить задачу, записать последнее состояние), уходит в ожидание на весь таймаут пула, и один такой шаг съедает бюджет остальных.
Правило: в завершении к базе не обращаются, если без этого можно обойтись, а если обращение нужно, у него свой короткий таймаут, и отказ не мешает остановке продолжаться:
async def mark_stopped(node_id: str) -> None:
try:
async with asyncio.timeout(2):
async with engine.begin() as conn:
await conn.execute(text("UPDATE nodes SET stopped_at = now() WHERE id = :id"), {"id": node_id})
except Exception:
log.warning("не удалось отметить остановку узла, продолжаем", exc_info=True)
Проверяется на стенде за минуту: остановите базу и сразу за ней под, посмотрите, за какое время завершился процесс. Если вместо «остановился за три секунды» получилось принудительное завершение через шестьдесят, в цепочке завершения есть шаг, который ждёт соединения.
Частые ошибки
engine.dispose() в начале shutdown. Если вызвать dispose() до отмены фоновых задач — задачи потеряют доступ к базе в середине транзакции. Dispose всегда идёт последним.
# так не надо: пул закрыт, а фоновые задачи ещё в транзакциях
await engine.dispose()
await _stop_background_tasks()
signal.signal(SIGTERM, ...) для dispose в теле модуля. Это создаёт параллельный канал завершения в обход lifespan. Uvicorn уже обрабатывает SIGTERM и вызывает lifespan — дополнительный обработчик не нужен.
while True в фоновом цикле. Без проверки флага готовности задача не остановится по команде. Используй while app_state.is_ready.
alembic downgrade в shutdown. Это не паттерн. Shutdown — это закрытие ресурсов, не откат схемы.
Логирование закрытия пула как ошибки. Сообщение asyncpg pool closed — нормальное событие завершения, не ошибка. Используй уровень INFO.
Глубже: две версии кода против одной базы: совместимость N-1расширенное
Во время rolling update старая и новая версии приложения несколько минут работают с одной базой, и откат через rollout undo возвращает старую версию к базе, которую уже поменяла новая. Отсюда правило N-1: каждая миграция совместима с предыдущей версией кода, а downgrade в Alembic пишут, но на него не рассчитывают.
Переименовать колонку одним шагом нельзя: старый под упадёт на первом же SELECT. Делают расширение и сжатие в три выката. Первый: op.add_column новой колонки, допускающей NULL, код пишет в обе и читает старую. Второй: заполнение истории пачками по тысяче строк в фоновой задаче, не в миграции, код читает новую. Третий: NOT NULL через CHECK ... NOT VALID и VALIDATE, затем удаление старой колонки. Индексы создают с CONCURRENTLY в миграции с autocommit-блоком, иначе таблица блокируется на запись на всё время построения.
Где запускать alembic upgrade head: не в lifespan. При трёх репликах три пода стартуют одновременно и трое пробуют мигрировать; Alembic распределённой блокировки не ставит, и исход зависит от удачи. Миграции идут отдельным шагом до выката: initContainer, который дождётся завершения, или Job в конвейере, после которого катится Deployment. Тогда поды новой версии стартуют уже на новой схеме, а поды старой её переживают по правилу N-1.
Коротко
engine.dispose()вызывается в lifespan-shutdown последним — после фоновых задач и Kafka.- HTTP-транзакции завершаются через uvicorn graceful (
--timeout-graceful-shutdown). - Фоновые задачи останавливаются через флаг готовности (
app_state.is_ready), а неwhile True. CancelledErrorв задаче — транзакция откатывается автоматически через__aexit__.- Alembic работает только при старте, в shutdown никакого DDL.
async with session_factory()на выходе только закрывает сессию; фиксирует работуcommit()или блокsession.begin(), всё остальное откатывается.- Старая и новая версии несколько минут делят одну базу: миграции совместимы с предыдущим кодом (расширение, заполнение пачками, сжатие), индексы
CONCURRENTLY, аalembic upgrade headидёт отдельным шагом до выката, не вlifespan.
Что почитать дальше
- HTTP drain — uvicorn graceful, preStop sleep, долгие эндпоинты.
- Фоновые задачи и outbox — asyncio-задачи, APScheduler, CancelledError.
- Kafka shutdown — aiokafka consumer и producer, ручной commit.
- Бюджеты и наблюдаемость — раскладка 60s-бюджета, метрики завершения.