Как писать асинхронные интеграции в Python: пул соединений, таймауты и graceful shutdown
Разберём типовые ошибки при работе с async-вызовами (утечки задач, неправильные таймауты, конкуренция за ресурсы) и соберём практический шаблон сервиса. Плюс — как корректно завершать фоновые корутины на shutdown, чтобы не терять запросы и не ронять обраб
Содержание
Как писать асинхронные интеграции в Python: пул соединений, таймауты и graceful shutdown
Асинхронные интеграции в Python — это не «магия без блокировок», а дисциплина: вы создаёте модель выполнения, где каждая корутина, таймаут и соединение должны иметь понятный жизненный цикл. На практике чаще всего ломаются не «async/await» как таковые, а инфраструктура вокруг: пул соединений, отмена задач, корректная остановка сервиса и обработка ошибок в условиях нагрузки и частичных отказов.
В этой статье разберём типовые ошибки (утечки задач, неправильные таймауты, конкуренция за ресурсы) и соберём практический шаблон асинхронного сервиса, который умеет:
- корректно управлять соединениями через пул;
- задавать таймауты на сетевые операции с понятной политикой;
- ограничивать параллелизм, чтобы не разнести внешний сервис и не исчерпать ресурсы;
- правильно завершать фоновые корутины на shutdown (graceful shutdown), не теряя в обработке то, что можно довести до конца, и не роняя обработку при остановке.
Под конец — коротко о том, как углубиться в тему; в контекстном месте упомянем курс по async-интеграциям: «».
Почему async-интеграции ломаются: карта рисков
У асинхронного кода есть соблазнительная простота: «я вызвал await — значит всё будет работать». Но проблема в том, что await не управляет всем остальным: у вас остаются открытые соединения, фоновые задачи, очереди, конкурентный доступ к ресурсам и отдельные ветки обработки ошибок.
Ниже — наиболее частые причины инцидентов.
Утечки задач (task leaks): «вроде остановили, но всё равно работает»
Типичная картина: сервис делает фоновые задачи (например, consumer очереди), а при shutdown пытается «просто отменить» или вообще не отменяет. В итоге:
- задачи продолжают жить после остановки приложения;
- при очередных событиях они используют уже закрытые ресурсы (пул, клиенты, файлы);
- процесс может завершиться с исключениями в фоновых корутинах;
- или хуже: задачи остаются висеть, удерживая память/соединения.
В asyncio утечки часто проявляются предупреждениями вида Task was destroyed but it is pending!, но нередко проблемы всплывают позже — в деградации и росте потребления ресурсов.
Неправильные таймауты: «таймаут есть, но не туда»
Самый распространённый анти-паттерн: ставить таймаут на «весь запрос целиком», не различая этапы и не учитывая контекст отмены. Например:
- таймаут не применяется к конкретному сетевому вызову, а только оборачивает логику сверху;
- отмена по таймауту не приводит к корректному освобождению контекста (например, задача поймала
CancelledErrorи не завершилась); - у сервиса несколько слоёв таймаутов (HTTP-клиент, операционный лимит, общий limit на обработку запроса) — их нужно согласовывать, иначе получаются неожиданные отмены «раньше/позже».
Конкуренция за ресурсы: «всё параллельно, но не всё выдержит»
Даже если вы используете пул соединений, не факт, что ваш код не устроит шторм:
- слишком высокая параллельность порождает очередь ожиданий в пуле;
- фоновые задачи конкурируют за один внешний сервис и приводят к каскадным таймаутам;
- без лимитов вы можете быстро исчерпать file descriptors, порты, или получить лимиты внешнего API.
Важный нюанс: пул соединений — это не тот же инструмент, что контроль параллелизма. Пул ограничивает одновременно используемые соединения, но вы можете создать больше задач, чем способен «переварить» внешний сервис, и тогда эти задачи начнут ждать в очереди, удерживая память и контекст.
Отмена и graceful shutdown: «shutdown был, но задачи не успели договорить»
Остановку приложения нужно рассматривать как протокол:
- перестаём принимать новые запросы или новые события в фоновые конвейеры;
- корректно завершаем активные операции (или прерываем их по таймауту остановки);
- освобождаем ресурсы в правильном порядке.
Если сделать не так, можно:
- потерять запросы (например, очередь не успела обработаться);
- повредить целостность (частично выполненные операции без компенсации);
- получить исключения при закрытии клиента/пула, пока операции ещё идут.
Практический каркас: интеграционный сервис с пулом, таймаутами и shutdown
Давайте соберём шаблон, который можно адаптировать под разные интеграции (HTTP, БД, очереди). В примере будем использовать aiohttp для HTTP-запросов, но принципы одинаковы и для asyncpg, aiomysql, и т.п.: пул, таймауты, лимиты параллелизма и управляемое завершение.
Выбор библиотек и архитектурная идея
- HTTP-клиент:
aiohttp.ClientSessionобычно создаётся один раз и живёт долго (на время приложения), чтобы переиспользовать connection pooling и reduce latency. - Пул соединений: у aiohttp он встроен в рамках
ClientSession(с настройками connector). - Ограничение параллелизма:
asyncio.Semaphoreили отдельный «worker pool», который контролирует, сколько запросов реально выполняется одновременно. - Таймауты: лучше задавать на уровне операций (connect/read/total), а также иметь внешний таймаут на обработку интеграции.
- Очередь задач: для фонового конвейера (например, отправка событий во внешний сервис) удобно использовать
asyncio.Queueи конечный набор worker-корунин. - Graceful shutdown: в сигнал-обработчике выставляем флаг остановки, перестаём добавлять в очередь, уведомляем workers, ждём их завершения с таймаутом и только потом закрываем
ClientSession.
Шаблон: классы интегратора
Ниже — рабочий каркас сервиса.
import asyncio
import logging
import signal
from dataclasses import dataclass
from typing import Any, Optional
import aiohttp
logger = logging.getLogger(__name__)
@dataclass(frozen=True)
class IntegrationConfig:
# Максимум открытых соединений к внешнему сервису
max_connections: int = 50
# Лимит одновременных запросов "по смыслу" (может быть меньше, чем max_connections)
max_in_flight: int = 20
# Таймауты на уровень запроса
connect_timeout: float = 3.0
sock_read_timeout: float = 10.0
total_timeout: float = 20.0
# Таймаут на остановку сервиса
shutdown_timeout: float = 15.0
# Кол-во worker-корунин для фоновой обработки
workers: int = 5
class ExternalAPIError(Exception):
pass
Теперь сам клиент интеграции.
class ExternalAPIClient:
def __init__(self, session: aiohttp.ClientSession, cfg: IntegrationConfig) -> None:
self._session = session
self._cfg = cfg
async def post_event(self, url: str, payload: dict[str, Any]) -> None:
# aiohttp timeout можно задать точнее
timeout = aiohttp.ClientTimeout(
total=self._cfg.total_timeout,
connect=self._cfg.connect_timeout,
sock_read=self._cfg.sock_read_timeout,
)
try:
async with self._session.post(url, json=payload, timeout=timeout) as resp:
# Важно: обрабатывать не только 200, но и корректно вычитывать тело при ошибках
if 200 <= resp.status < 300:
await resp.read() # полезно, чтобы корректно завершить соединение
return
body = await resp.text()
raise ExternalAPIError(f"HTTP {resp.status}: {body[:500]}")
except asyncio.TimeoutError as e:
# Внешний таймаут: это ожидаемый сценарий, его нужно отличать от CancelledError
raise ExternalAPIError("Request timeout") from e
except aiohttp.ClientError as e:
raise ExternalAPIError("Client error") from e
Обратите внимание на порядок обработки:
CancelledErrorне ловим напрямую (в Python 3.8+ он наследуется отBaseException, а неException), значит отмена «пройдёт» корректно.- ошибки сети и таймауты превращаем в доменные исключения — так легче строить политику ретраев/компенсации.
Сервис с очередью и worker-корутинами
Смысл: у нас есть очередь событий для отправки во внешний сервис и фиксированное число worker-корунин, которые эти события отправляют. Плюсы:
- легко ограничить параллельность;
- управляемая структура для shutdown;
- меньше шанс «утечки» задач.
class IntegrationService:
def __init__(self, cfg: IntegrationConfig, api_url: str) -> None:
self._cfg = cfg
self._api_url = api_url
self._queue: asyncio.Queue[dict[str, Any]] = asyncio.Queue()
self._stopping = asyncio.Event()
self._workers: list[asyncio.Task[None]] = []
# Семафор ограничивает одновременные in-flight запросы.
self._in_flight = asyncio.Semaphore(cfg.max_in_flight)
self._session: Optional[aiohttp.ClientSession] = None
self._client: Optional[ExternalAPIClient] = None
async def start(self) -> None:
# Важно: создаём session один раз и используем connector с лимитами
connector = aiohttp.TCPConnector(
limit=self._cfg.max_connections,
limit_per_host=self._cfg.max_connections, # при необходимости можно уменьшить
enable_cleanup_closed=True,
)
self._session = aiohttp.ClientSession(connector=connector)
self._client = ExternalAPIClient(self._session, self._cfg)
# Запуск worker-корунин
for i in range(self._cfg.workers):
task = asyncio.create_task(self._worker_loop(worker_id=i), name=f"worker-{i}")
self._workers.append(task)
async def enqueue_event(self, payload: dict[str, Any]) -> None:
# При остановке новые события не добавляем
if self._stopping.is_set():
raise RuntimeError("Service is stopping; cannot enqueue new events")
await self._queue.put(payload)
async def _worker_loop(self, worker_id: int) -> None:
assert self._client is not None
assert self._session is not None
while True:
# Если мы в режиме остановки и очередь опустела — выходим
if self._stopping.is_set() and self._queue.empty():
return
try:
# Таймаут получения из очереди позволяет worker-у переоценить условие stop.
payload = await asyncio.wait_for(self._queue.get(), timeout=0.5)
except asyncio.TimeoutError:
continue
try:
# Ограничиваем число одновременных запросов.
async with self._in_flight:
await self._client.post_event(self._api_url, payload)
except ExternalAPIError as e:
# Тут нужна политика: логирование, ретраи, DLQ и т.д.
logger.warning("worker=%s failed to post event: %s payload=%s",
worker_id, e, payload)
# Пример: не делаем ретраи в этом шаблоне.
except Exception:
logger.exception("worker=%s unexpected error", worker_id)
finally:
self._queue.task_done()
async def graceful_shutdown(self) -> None:
# 1) Переходим в режим остановки: новые события не принимаем
self._stopping.set()
# 2) Ждём опустошения очереди и завершения активных worker-корунин.
# Смысл: работаем до тех пор, пока возможно, но не бесконечно.
try:
await asyncio.wait_for(self._wait_for_drain(), timeout=self._cfg.shutdown_timeout)
except asyncio.TimeoutError:
logger.warning("Shutdown timeout reached; cancelling workers")
# 3) На всякий случай отменяем задачи, если они всё ещё живы
for t in self._workers:
if not t.done():
t.cancel()
# 4) Дожидаемся завершения задач
await asyncio.gather(*self._workers, return_exceptions=True)
# 5) Закрываем session (после остановки worker-ов!)
if self._session is not None and not self._session.closed:
await self._session.close()
async def _wait_for_drain(self) -> None:
# Очередь task_done фиксирует прогресс. join ждёт, пока все task_done будут вызваны.
await self._queue.join()
Ключевые детали:
- Сигнал остановки —
self._stopping.set(). - Worker периодически проверяет
stopдаже если очередь пуста/почти пуста (черезwait_forс таймаутом наqueue.get). - Ограничение параллельности —
asyncio.Semaphore(max_in_flight), отдельно от лимита соединений вTCPConnector. - Graceful shutdown — сначала дождаться
queue.join()(когда все элементы обработаны), затем отмена оставшихся workers, затем закрытьClientSession.
Пример запуска со сбором сигналов
Чтобы shutdown был корректным в реальном процессе, добавим обработку SIGTERM/SIGINT.
async def main() -> None:
logging.basicConfig(level=logging.INFO)
cfg = IntegrationConfig(
max_connections=50,
max_in_flight=20,
workers=5,
shutdown_timeout=15.0,
total_timeout=20.0,
connect_timeout=3.0,
sock_read_timeout=10.0,
)
service = IntegrationService(cfg=cfg, api_url="https://example.com/api/events")
await service.start()
loop = asyncio.get_running_loop()
stop_event = asyncio.Event()
def _signal_handler() -> None:
stop_event.set()
for sig in (signal.SIGINT, signal.SIGTERM):
try:
loop.add_signal_handler(sig, _signal_handler)
except NotImplementedError:
# В некоторых окружениях (например, Windows) add_signal_handler может быть недоступен
pass
# Пример: генерируем события
async def producer():
i = 0
while not stop_event.is_set():
await service.enqueue_event({"id": i, "ts": asyncio.get_running_loop().time()})
i += 1
await asyncio.sleep(0.05)
prod_task = asyncio.create_task(producer(), name="producer")
await stop_event.wait()
# Останов: прекращаем продюсер, затем graceful shutdown интегратора
prod_task.cancel()
await asyncio.gather(prod_task, return_exceptions=True)
await service.graceful_shutdown()
if __name__ == "__main__":
asyncio.run(main())
Типовые ошибки в этой части и как их предотвращать
Ошибка 1: «Отмена отменяет всё» — и потеря обработки
Если вы в shutdown просто делаете for task in tasks: task.cancel() и не ждёте, можно:
- прервать операции на середине;
- потерять обработанные, но ещё не
task_done()элементы; - оставить очередь в состоянии, где
queue.join()никогда не завершится.
Правильный подход: сначала выставить флаг остановки, затем дождаться drain с таймаутом, и только в крайнем случае отменять.
Ошибка 2: «task_done вызывается не всегда»
В нашем worker мы кладём self._queue.task_done() в finally. Если забыть — queue.join() будет ждать вечно и shutdown превращается в зависание.
Ошибка 3: «Ставим таймаут на верхний уровень и игнорируем локальные»
Допустим, вы сделали await asyncio.wait_for(client.post_event(...), timeout=10) поверх того, что уже имеет ClientTimeout(total=20). Что произойдёт при нагрузке? Часто — отмена снаружи, но внутренний клиент может не завершить корректно cleanup. В результате — цепочка нестабильностей.
Практика: держите иерархию таймаутов:
- внутренняя операция (connect/read/total) — ближе к сетевому вызову;
- верхний таймаут — ограничение на весь workflow обработки запроса/события (если нужно), но его значения должны быть согласованы (обычно верхний таймаут чуть больше внутреннего, чтобы избежать преждевременных отмен).
Ошибка 4: «Параллельность контролируем только пулом»
Если пул соединений ограничен, но вы создаёте тысячи worker-корунин или тысячи задач на отправку, вы можете:
- утопить event loop в накладных расходах;
- разогнать память под контекст задач;
- усилить backpressure только на уровне сокетов, но не на уровне очередей.
Нормальная стратегия: лимитировать in-flight (семафор) и/или контролировать число worker-корунин.
Ошибка 5: «Закрываем session до завершения задач»
Порядок shutdown критичен. Если закрыть ClientSession до того, как worker-корутины перестали выполнять HTTP-запросы, получите Connection closed/RuntimeError: Session is closed и «грязное» поведение. В нашем каркасе session закрывается после завершения worker-ов.
Как делать таймауты осмысленными: политика и значения
Таймауты — это часть контракта системы. Хорошая практика:
- Connect timeout: ограничивает время на установление соедин
Комментарии
Пока нет комментариев