← Все темы
Asyncio python
Вопросов: 41
asyncio — это стандартная библиотека Python для асинхронного (неблокирующего) выполнения задач в одном потоке на основе цикла событий (
Какие задачи решает asyncio
- Параллельное по времени выполнение I/O-задач: сетевые запросы, работа с сокетами, HTTP/WebSocket, ожидание таймеров, чтение/запись (через асинхронные библиотеки).
- Высокая конкурентность при малых накладных расходах: тысячи соединений/задач в одном процессе, когда основное время уходит на ожидание I/O.
- Структурирование асинхронного кода с
- Построение асинхронных серверов: обработка большого числа клиентов (например, чат, API-шлюз, прокси).
- Оркестрация фоновых задач: периодические задания, пайплайны, конкурентные запросы к нескольким сервисам.
Что asyncio не решает напрямую
- CPU-bound задачи (тяжёлые вычисления) оно не ускоряет: для этого обычно используют
Мини-пример
event loop). Она позволяет эффективно обслуживать множество операций ввода‑вывода без создания большого числа потоков.Какие задачи решает asyncio
- Параллельное по времени выполнение I/O-задач: сетевые запросы, работа с сокетами, HTTP/WebSocket, ожидание таймеров, чтение/запись (через асинхронные библиотеки).
- Высокая конкурентность при малых накладных расходах: тысячи соединений/задач в одном процессе, когда основное время уходит на ожидание I/O.
- Структурирование асинхронного кода с
async/await: корутины, задачи (Task), отмена, таймауты, очереди, семафоры.- Построение асинхронных серверов: обработка большого числа клиентов (например, чат, API-шлюз, прокси).
- Оркестрация фоновых задач: периодические задания, пайплайны, конкурентные запросы к нескольким сервисам.
Что asyncio не решает напрямую
- CPU-bound задачи (тяжёлые вычисления) оно не ускоряет: для этого обычно используют
multiprocessing или вынос в отдельные процессы/воркеры (в asyncio можно лишь запускать такие задачи через executors).Мини-пример
import asyncio
async def fetch(name, delay):
await asyncio.sleep(delay) # имитация I/O ожидания
return f"{name} done"
async def main():
results = await asyncio.gather(
fetch("A", 1),
fetch("B", 1),
fetch("C", 1),
)
print(results)
asyncio.run(main())
Асинхронность — это модель конкурентности, где одна нить выполнения (обычно один поток) переключается между задачами, пока они ждут (I/O: сеть, диск, БД, таймеры). В Python чаще всего это
- хорошо для большого числа I/O‑задач
- плохо ускоряет CPU‑тяжёлые вычисления
Параллелизм — это реальное одновременное выполнение задач (на нескольких ядрах/процессорах) с целью ускорения. В Python для CPU‑bound обычно достигается через
- ускоряет CPU‑bound при использовании нескольких ядер
- требует межпроцессного обмена данными (IPC), выше накладные расходы
Многопоточность — это несколько потоков внутри одного процесса. В CPython есть GIL (Global Interpreter Lock), поэтому для CPU‑bound кода потоки обычно не дают параллельного ускорения, но полезны для I/O‑bound (когда поток блокируется на I/O, другой может выполняться).
- удобно для I/O и для обёрток над блокирующими библиотеками
- CPU‑bound почти не ускоряет в CPython из-за GIL
Ключевые различия в контексте Python
- Async: конкурентность в одном потоке, переключение на await, отлично для масштабируемого I/O.
- Threads: конкурентность через потоки; I/O‑bound ок, CPU‑bound ограничен GIL.
- Parallel: настоящая одновременность на ядрах; для CPU‑bound обычно нужны процессы (или нативные расширения, освобождающие GIL).
asyncio: кооперативное переключение, управляемое event loop. - хорошо для большого числа I/O‑задач
- плохо ускоряет CPU‑тяжёлые вычисления
Параллелизм — это реальное одновременное выполнение задач (на нескольких ядрах/процессорах) с целью ускорения. В Python для CPU‑bound обычно достигается через
multiprocessing (несколько процессов). - ускоряет CPU‑bound при использовании нескольких ядер
- требует межпроцессного обмена данными (IPC), выше накладные расходы
Многопоточность — это несколько потоков внутри одного процесса. В CPython есть GIL (Global Interpreter Lock), поэтому для CPU‑bound кода потоки обычно не дают параллельного ускорения, но полезны для I/O‑bound (когда поток блокируется на I/O, другой может выполняться).
- удобно для I/O и для обёрток над блокирующими библиотеками
- CPU‑bound почти не ускоряет в CPython из-за GIL
Ключевые различия в контексте Python
- Async: конкурентность в одном потоке, переключение на await, отлично для масштабируемого I/O.
- Threads: конкурентность через потоки; I/O‑bound ок, CPU‑bound ограничен GIL.
- Parallel: настоящая одновременность на ядрах; для CPU‑bound обычно нужны процессы (или нативные расширения, освобождающие GIL).
# Асинхронность (I/O-конкурентность)
import asyncio
async def fetch():
reader, writer = await asyncio.open_connection("example.com", 80)
writer.write(b"GET / HTTP/1.0\r\nHost: example.com\r\n\r\n")
await writer.drain()
data = await reader.read(1000)
writer.close()
await writer.wait_closed()
return data
async def main():
await asyncio.gather(*(fetch() for _ in range(100)))
asyncio.run(main())
# Многопоточность (I/O-bound, блокирующие вызовы)
from concurrent.futures import ThreadPoolExecutor
import requests
def fetch(url):
return requests.get(url, timeout=5).status_code
with ThreadPoolExecutor(max_workers=50) as ex:
list(ex.map(fetch, ["https://example.com"] * 100))
# Параллелизм (CPU-bound) через процессы
from concurrent.futures import ProcessPoolExecutor
def cpu_heavy(n):
s = 0
for i in range(n):
s += i * i
return s
with ProcessPoolExecutor() as ex:
list(ex.map(cpu_heavy, [10_000_00] * 8))
Asyncio лучше всего подходит
- IO-bound задачи: много ожиданий внешних ресурсов (сеть/диск), мало CPU. Почему:
- Высококонкурентные сетевые сервисы: HTTP/WebSocket/боты/прокси, тысячи одновременных соединений. Почему: один поток + неблокирующий ввод-вывод дают низкие накладные расходы по сравнению с множеством потоков.
- Пайплайны с множеством независимых ожиданий: параллельные запросы к API/БД/очередям. Почему: удобно управлять конкурентностью через
- Задачи с таймерами и реактивной логикой: планирование, ретраи, дедлайны. Почему: встроенные примитивы
Asyncio подходит плохо
- CPU-bound вычисления: тяжёлая математика, парсинг больших данных, сжатие, криптография. Почему: цикл событий блокируется, другие корутины не выполняются; GIL не даёт ускорения в одном процессе. Решение:
- Блокирующие библиотеки/драйверы (без async-API): синхронные HTTP/БД/файловые операции. Почему: блокируют event loop. Решение: async-аналоги или вынос в
- Простой линейный скрипт без конкурентности. Почему: усложнение кода (корутины, event loop) без выигрыша.
- Задачи, требующие жёсткого realtime. Почему: планирование в ОС и кооперативная многозадачность не гарантируют точные дедлайны.
Ключевое правило выбора
- Если время тратится на ожидание (I/O) и это ожидание можно сделать неблокирующим — asyncio даёт выигрыш.
- Если время тратится на вычисления (CPU) или на блокирующие вызовы — используйте процессы/потоки или заменяйте библиотеку на async-вариант.
- IO-bound задачи: много ожиданий внешних ресурсов (сеть/диск), мало CPU. Почему:
await освобождает цикл событий, позволяя выполнять другие корутины, пока текущая ждёт.- Высококонкурентные сетевые сервисы: HTTP/WebSocket/боты/прокси, тысячи одновременных соединений. Почему: один поток + неблокирующий ввод-вывод дают низкие накладные расходы по сравнению с множеством потоков.
- Пайплайны с множеством независимых ожиданий: параллельные запросы к API/БД/очередям. Почему: удобно управлять конкурентностью через
asyncio.gather, Semaphore, таймауты и отмену.- Задачи с таймерами и реактивной логикой: планирование, ретраи, дедлайны. Почему: встроенные примитивы
sleep, wait_for, отмена задач и обработка таймаутов.Asyncio подходит плохо
- CPU-bound вычисления: тяжёлая математика, парсинг больших данных, сжатие, криптография. Почему: цикл событий блокируется, другие корутины не выполняются; GIL не даёт ускорения в одном процессе. Решение:
multiprocessing или ProcessPoolExecutor.- Блокирующие библиотеки/драйверы (без async-API): синхронные HTTP/БД/файловые операции. Почему: блокируют event loop. Решение: async-аналоги или вынос в
ThreadPoolExecutor (как компромисс).- Простой линейный скрипт без конкурентности. Почему: усложнение кода (корутины, event loop) без выигрыша.
- Задачи, требующие жёсткого realtime. Почему: планирование в ОС и кооперативная многозадачность не гарантируют точные дедлайны.
Ключевое правило выбора
- Если время тратится на ожидание (I/O) и это ожидание можно сделать неблокирующим — asyncio даёт выигрыш.
- Если время тратится на вычисления (CPU) или на блокирующие вызовы — используйте процессы/потоки или заменяйте библиотеку на async-вариант.
Event loop — это центральный цикл выполнения в асинхронной модели, который планирует и координирует выполнение задач без блокировки потока.
Роль event loop в asyncio:
- Планирование корутин: запускает корутины, переключает их, когда они делают
- Неблокирующий I/O: следит за сокетами/каналами ввода-вывода и «будит» ожидающие корутины, когда данные доступны или запись возможна.
- Управление задачами: создаёт и выполняет
- Таймеры и задержки: обслуживает отложенные вызовы и ожидания вроде
- Колбэки: выполняет запланированные функции (callbacks) и интегрирует их с асинхронными событиями.
Ключевая идея: пока корутина ждёт I/O или таймер, event loop не простаивает — он выполняет другие готовые задачи, обеспечивая конкурентность в одном потоке.
Роль event loop в asyncio:
- Планирование корутин: запускает корутины, переключает их, когда они делают
await, и продолжает выполнение, когда операция готова.- Неблокирующий I/O: следит за сокетами/каналами ввода-вывода и «будит» ожидающие корутины, когда данные доступны или запись возможна.
- Управление задачами: создаёт и выполняет
asyncio.Task, отслеживает их состояние, обрабатывает отмену (cancel) и исключения.- Таймеры и задержки: обслуживает отложенные вызовы и ожидания вроде
asyncio.sleep().- Колбэки: выполняет запланированные функции (callbacks) и интегрирует их с асинхронными событиями.
Ключевая идея: пока корутина ждёт I/O или таймер, event loop не простаивает — он выполняет другие готовые задачи, обеспечивая конкурентность в одном потоке.
import asyncio
async def main():
await asyncio.sleep(1)
return "done"
# asyncio.run() создаёт event loop, запускает корутину и корректно закрывает цикл
result = asyncio.run(main())
print(result)
Event loop в
Ключевые сущности
- Coroutine (корутина): функция
- Task: обёртка над корутиной, которую event loop умеет планировать (создаётся через
- Future: объект «результат будет позже»;
- Awaitable: то, что можно
- Selector: механизм ожидания готовности сокетов/pipe (обычно
Что происходит при
- Корутина выполняется до первого
- На
- Loop регистрирует «условие продолжения»: например, «сокет станет читаемым», «таймер истечёт», «другая Task завершится».
- Когда условие выполнено, loop ставит продолжение корутины в очередь «готовых к выполнению» и позже снова запускает её с места после
Одна итерация event loop (упрощённо)
- Взять все «готовые» callbacks/Tasks из очереди ready и выполнить их кусками (каждая корутина бежит до следующего
- Посчитать ближайший таймер (sleep, call_later, timeout) и выбрать таймаут ожидания.
- Заблокироваться на короткое ожидание I/O через selector: «какие дескрипторы готовы на чтение/запись» или «истёк таймер».
- По событиям I/O/таймерам добавить соответствующие callbacks/Tasks в ready.
- Повторить.
Почему это не параллельность
- В одном loop’е по умолчанию выполняется один Python-поток; задачи чередуются кооперативно.
- Если внутри корутины сделать блокирующий вызов (например, обычный
- Для CPU-bound или блокирующих операций используют
Мини-пример, показывающий переключение
Как это выполняется
-
- Loop запускает
- Затем loop запускает
- Loop ждёт истечения таймеров (и/или I/O), затем продолжает обе корутины.
Практическое правило
- Каждое
- Чем чаще корутина делает
asyncio — это центральный цикл, который планирует и исполняет корутины, обрабатывает I/O и таймеры, переключая выполнение между задачами, когда они ждут.Ключевые сущности
- Coroutine (корутина): функция
async def, которая может await.- Task: обёртка над корутиной, которую event loop умеет планировать (создаётся через
asyncio.create_task()).- Future: объект «результат будет позже»;
Task тоже является Future.- Awaitable: то, что можно
await (корутина/Task/Future).- Selector: механизм ожидания готовности сокетов/pipe (обычно
selectors), т.е. «когда можно читать/писать без блокировки».Что происходит при
await- Корутина выполняется до первого
await.- На
await X она добровольно отдаёт управление loop’у.- Loop регистрирует «условие продолжения»: например, «сокет станет читаемым», «таймер истечёт», «другая Task завершится».
- Когда условие выполнено, loop ставит продолжение корутины в очередь «готовых к выполнению» и позже снова запускает её с места после
await.Одна итерация event loop (упрощённо)
- Взять все «готовые» callbacks/Tasks из очереди ready и выполнить их кусками (каждая корутина бежит до следующего
await или завершения).- Посчитать ближайший таймер (sleep, call_later, timeout) и выбрать таймаут ожидания.
- Заблокироваться на короткое ожидание I/O через selector: «какие дескрипторы готовы на чтение/запись» или «истёк таймер».
- По событиям I/O/таймерам добавить соответствующие callbacks/Tasks в ready.
- Повторить.
Почему это не параллельность
- В одном loop’е по умолчанию выполняется один Python-поток; задачи чередуются кооперативно.
- Если внутри корутины сделать блокирующий вызов (например, обычный
time.sleep() или синхронный запрос), loop «замрёт» и не сможет переключаться.- Для CPU-bound или блокирующих операций используют
asyncio.to_thread() или run_in_executor().Мини-пример, показывающий переключение
import asyncio
async def worker(name):
print("start", name)
await asyncio.sleep(1) # корутина отдаёт управление loop'у + ставится таймер
print("end", name)
async def main():
t1 = asyncio.create_task(worker("A"))
t2 = asyncio.create_task(worker("B"))
await t1
await t2
asyncio.run(main())
Как это выполняется
-
create_task ставит корутины в очередь ready.- Loop запускает
worker("A") до await asyncio.sleep(1); sleep регистрирует таймер и возвращает управление loop’у.- Затем loop запускает
worker("B") до await sleep.- Loop ждёт истечения таймеров (и/или I/O), затем продолжает обе корутины.
Практическое правило
- Каждое
await — потенциальная точка переключения.- Чем чаще корутина делает
await на неблокирующих операциях, тем «живее» и отзывчивее работает приложение.
Как получить и запустить event loop (Python 3.10+)
-
1) Самый рекомендуемый способ:
Создаёт новый event loop, запускает корутину до завершения, корректно закрывает loop, отменяет незавершённые задачи, завершает async generators и освобождает ресурсы.
-
2) Получение текущего loop внутри корутины:
Правильный способ получить loop, когда он уже запущен (внутри
-
3) Низкоуровневый ручной запуск (обычно не нужен)
Полезно для нестандартных сценариев, но легко допустить утечки/некорректное закрытие.
-
Почему часто рекомендуют
- Безопасное управление жизненным циклом: создаёт и закрывает loop правильно, не оставляя «висящих» задач/дескрипторов.
- Меньше ошибок: не нужно вручную вызывать
- Предсказуемость: один чёткий вход в async-мир из синхронного кода (идеально для
- Современная практика: вне корутин не следует полагаться на
-
1) Самый рекомендуемый способ:
asyncio.run()Создаёт новый event loop, запускает корутину до завершения, корректно закрывает loop, отменяет незавершённые задачи, завершает async generators и освобождает ресурсы.
import asyncio
async def main():
await asyncio.sleep(1)
return 42
result = asyncio.run(main())
print(result)
-
2) Получение текущего loop внутри корутины:
asyncio.get_running_loop()Правильный способ получить loop, когда он уже запущен (внутри
async def).
import asyncio
async def main():
loop = asyncio.get_running_loop()
print(loop)
asyncio.run(main())
-
3) Низкоуровневый ручной запуск (обычно не нужен)
Полезно для нестандартных сценариев, но легко допустить утечки/некорректное закрытие.
import asyncio
async def main():
await asyncio.sleep(0.1)
loop = asyncio.new_event_loop()
try:
asyncio.set_event_loop(loop)
loop.run_until_complete(main())
finally:
loop.close()
-
Почему часто рекомендуют
asyncio.run()- Безопасное управление жизненным циклом: создаёт и закрывает loop правильно, не оставляя «висящих» задач/дескрипторов.
- Меньше ошибок: не нужно вручную вызывать
new_event_loop(), set_event_loop(), run_until_complete(), close().- Предсказуемость: один чёткий вход в async-мир из синхронного кода (идеально для
if __name__ == "__main__", CLI, скриптов).- Современная практика: вне корутин не следует полагаться на
asyncio.get_event_loop() (в новых версиях поведение менялось и часто приводит к путанице).s
Coroutine (короутина) — это специальная функция, выполнение которой можно приостанавливать и возобновлять позже, сохраняя её состояние (локальные переменные, точку выполнения). В Python короутины используются для асинхронного кода и работают через
Чем отличается от обычной функции
- Обычная функция (
- Короутина (
Ключевые признаки короутины в Python
- объявляется как
- внутри использует
- требует event loop (например,
Практический смысл: короутины позволяют эффективно обрабатывать множество операций ожидания (I/O) в одном потоке, не блокируя выполнение на время ожидания.
await внутри async def.Чем отличается от обычной функции
- Обычная функция (
def) выполняется сразу и до конца при вызове (если не выбросит исключение/не сделает ранний return).- Короутина (
async def) при вызове не выполняется сразу: она возвращает объект короутины, который нужно запустить (обычно через await или планирование в event loop). Внутри она может отдавать управление в точках await, позволяя выполнять другие задачи, пока она ждёт (например, сеть/диск/таймер).Ключевые признаки короутины в Python
- объявляется как
async def- внутри использует
await для неблокирующего ожидания- требует event loop (например,
asyncio) для выполненияimport asyncio
def обычная():
return "я выполнилась сразу"
async def короутина():
await asyncio.sleep(1)
return "я выполнилась после await"
async def main():
print(обычная()) # сразу
c = короутина() # пока НЕ выполняется, это объект короутины
print(await c) # запускаем и ждём результат
asyncio.run(main())Практический смысл: короутины позволяют эффективно обрабатывать множество операций ожидания (I/O) в одном потоке, не блокируя выполнение на время ожидания.
await приостанавливает выполнение текущей асинхронной функции до тех пор, пока не завершится ожидаемый объект (обычно
Зачем нужно
- позволяет выполнять неблокирующее ожидание: пока идёт I/O (сеть, файлы, БД, таймер), event loop может выполнять другие задачи.
- упрощает код: вместо колбэков/ручного управления состояниями — линейный стиль.
Где можно использовать (Python)
- только внутри функции, объявленной как
- можно применять к объектам, поддерживающим протокол ожидания (
- нельзя использовать в обычной
Типичные ошибки
- написать
- забыть
coroutine/Task/Future), и затем возвращает его результат (или пробрасывает исключение).Зачем нужно
- позволяет выполнять неблокирующее ожидание: пока идёт I/O (сеть, файлы, БД, таймер), event loop может выполнять другие задачи.
- упрощает код: вместо колбэков/ручного управления состояниями — линейный стиль.
Где можно использовать (Python)
- только внутри функции, объявленной как
async def. - можно применять к объектам, поддерживающим протокол ожидания (
__await__): корутины, asyncio.Task, asyncio.Future и т.п. - нельзя использовать в обычной
def и на верхнем уровне файла (кроме интерактивной среды, где top-level await может поддерживаться).import asyncio
async def fetch():
await asyncio.sleep(1) # неблокирующее ожидание
return "ok"
async def main():
result = await fetch() # ждём завершения корутины и получаем результат
print(result)
asyncio.run(main())
Типичные ошибки
- написать
await вне async def → SyntaxError. - забыть
await при вызове корутины → получите объект-корутину, а не результат (и возможное предупреждение о “coroutine was never awaited”).
Awaitable-объекты — это объекты, которые можно использовать с
Основные виды awaitable в Python
- Coroutine (корутина): объект, возвращаемый вызовом async-функции.
- Task: обёртка над корутиной, запланированная на выполнение в event loop (обычно через
- Future: низкоуровневый объект-плейсхолдер результата, который будет установлен позже; awaitable (в
- Пользовательский awaitable: любой объект с методом
Важно: coroutine, Task и Future — самые частые awaitable в практике; Task и Future обычно относятся к библиотеке
await в async-коде. При await x управление передаётся циклу событий, пока x не завершится, после чего возвращается результат (или возбуждается исключение).Основные виды awaitable в Python
- Coroutine (корутина): объект, возвращаемый вызовом async-функции.
async def f():
return 123
coro = f() # coro — awaitable
res = await coro
- Task: обёртка над корутиной, запланированная на выполнение в event loop (обычно через
asyncio.create_task); тоже awaitable.import asyncio
task = asyncio.create_task(f()) # task — awaitable
res = await task
- Future: низкоуровневый объект-плейсхолдер результата, который будет установлен позже; awaitable (в
asyncio это asyncio.Future).import asyncio
loop = asyncio.get_running_loop()
fut = loop.create_future() # fut — awaitable
# где-то позже: fut.set_result(123)
res = await fut
- Пользовательский awaitable: любой объект с методом
__await__(), возвращающим итератор (обычно генератор), тем самым делая объект совместимым с await.class MyAwaitable:
def __await__(self):
async def _impl():
return "ok"
return _impl().__await__()
res = await MyAwaitable()
Важно: coroutine, Task и Future — самые частые awaitable в практике; Task и Future обычно относятся к библиотеке
asyncio, а корутины — к языковой конструкции async/await.
Корутина — это функция, объявленная через
Пример: создать и выполнить простую короутину
Что будет, если короутину не await-ить
- При вызове
- Если объект корутины так и не будет
- Обычно при завершении программы/сборке мусора появится предупреждение:
Мини-пример “не await”
Правильно “запустить без await” (через задачу)
async def, которая возвращает объект корутины и выполняется только при await (или если её запланировать в event loop).Пример: создать и выполнить простую короутину
import asyncio
async def hello():
print("start")
await asyncio.sleep(0.1)
print("end")
return 42
async def main():
result = await hello() # корутина реально выполняется здесь
print("result:", result)
asyncio.run(main())
Что будет, если короутину не await-ить
- При вызове
hello() код внутри корутины не выполняется; создаётся только объект корутины.- Если объект корутины так и не будет
await-нут (и не будет запланирован как task), то она не выполнится вообще.- Обычно при завершении программы/сборке мусора появится предупреждение:
RuntimeWarning: coroutine 'hello' was never awaited.Мини-пример “не await”
import asyncio
async def hello():
print("I will not run")
async def main():
c = hello() # создали корутину, но не await
print("created:", c)
asyncio.run(main())
Правильно “запустить без await” (через задачу)
import asyncio
async def hello():
await asyncio.sleep(0.1)
print("ran")
async def main():
task = asyncio.create_task(hello()) # запланировали выполнение
await task # дождались результата (или можно не ждать, если это осознанно)
asyncio.run(main())
Coroutine object — это объект-результат вызова async-функции, например:
Task (asyncio.Task) — это обёртка над coroutine object, которая планирует его выполнение в event loop и запускает конкурентно (кооперативно) с другими задачами. Обычно создаётся через
Ключевые отличия
- Запуск: coroutine object не начнёт выполняться, пока его не
- Конкурентность: Task позволяет “параллельно” (в рамках одного потока через переключения на await) выполнять корутину, пока текущая корутина продолжает работу.
- Управление: у Task есть состояние (
- Получение результата: у Task результат можно получить через
- Многократное ожидание: Task можно безопасно
coro = foo(). Он описывает выполнение, но сам по себе не запускается; чтобы он выполнялся, его нужно await-ить или передать циклу событий.Task (asyncio.Task) — это обёртка над coroutine object, которая планирует его выполнение в event loop и запускает конкурентно (кооперативно) с другими задачами. Обычно создаётся через
asyncio.create_task(coro).Ключевые отличия
- Запуск: coroutine object не начнёт выполняться, пока его не
await-нут или не превратят в Task; Task начинает выполняться после постановки в цикл событий.- Конкурентность: Task позволяет “параллельно” (в рамках одного потока через переключения на await) выполнять корутину, пока текущая корутина продолжает работу.
- Управление: у Task есть состояние (
pending/done/cancelled), методы cancel(), done(), получение результата/исключения; coroutine object — это просто вычисление “в ожидании”.- Получение результата: у Task результат можно получить через
await task (или task.result() после завершения); coroutine object — через непосредственный await coro.- Многократное ожидание: Task можно безопасно
await-ить из разных мест (результат кэшируется после завершения); coroutine object нельзя “переиспользовать” — повторный await одной и той же корутины приведёт к ошибке.import asyncio
async def work():
await asyncio.sleep(1)
return 42
async def main():
coro = work() # coroutine object: не запущен
task = asyncio.create_task(work()) # Task: запущен в фоне
r1 = await coro # запуск через await
r2 = await task # дождаться уже запущенного
print(r1, r2)
asyncio.run(main())
Task в asyncio — это объект, который планирует выполнение корутины в цикле событий и позволяет ей выполняться параллельно (конкурентно) с другими корутинами в одном потоке.
Зачем создавать Task в реальном коде
- Запустить работу “в фоне”, не блокируя текущую корутину (например, периодическая синхронизация, отправка метрик, обработка очереди).
- Параллелить независимые I/O-операции (несколько запросов/чтений/записей одновременно) для снижения общего времени ожидания.
- Управлять жизненным циклом: отмена через
- Корректно обработать исключения: если “фоновую” корутину просто вызвать без ожидания, ошибки могут потеряться/всплыть поздно; Task позволяет их собрать и обработать.
Как создать Task
- Внутри работающего event loop:
- Явно через loop (реже):
Пример “фоновой” задачи + правильная остановка
Практическое правило: создавай Task, когда корутину нужно запустить сейчас и дать ей выполняться независимо (а не ждать её немедленно). Если корутину нужно выполнить строго последовательно — достаточно
Зачем создавать Task в реальном коде
- Запустить работу “в фоне”, не блокируя текущую корутину (например, периодическая синхронизация, отправка метрик, обработка очереди).
- Параллелить независимые I/O-операции (несколько запросов/чтений/записей одновременно) для снижения общего времени ожидания.
- Управлять жизненным циклом: отмена через
task.cancel(), ожидание через await task, таймауты.- Корректно обработать исключения: если “фоновую” корутину просто вызвать без ожидания, ошибки могут потеряться/всплыть поздно; Task позволяет их собрать и обработать.
Как создать Task
- Внутри работающего event loop:
asyncio.create_task(coro())- Явно через loop (реже):
loop.create_task(coro())import asyncio
async def fetch(n: int) -> str:
await asyncio.sleep(1)
return f"ok:{n}"
async def main():
# 1) Запускаем две операции конкурентно
t1 = asyncio.create_task(fetch(1))
t2 = asyncio.create_task(fetch(2))
# делаем что-то ещё, пока они выполняются
await asyncio.sleep(0.1)
# 2) Дожидаемся результатов
r1 = await t1
r2 = await t2
print(r1, r2)
asyncio.run(main())Пример “фоновой” задачи + правильная остановка
import asyncio
async def heartbeat():
try:
while True:
print("tick")
await asyncio.sleep(2)
except asyncio.CancelledError:
# здесь можно сделать финализацию (закрыть ресурсы и т.п.)
raise
async def main():
hb_task = asyncio.create_task(heartbeat())
try:
await asyncio.sleep(5) # основная работа
finally:
hb_task.cancel()
with contextlib.suppress(asyncio.CancelledError):
await hb_task
if __name__ == "__main__":
import contextlib
asyncio.run(main())Практическое правило: создавай Task, когда корутину нужно запустить сейчас и дать ей выполняться независимо (а не ждать её немедленно). Если корутину нужно выполнить строго последовательно — достаточно
await coro() без Task.
Последовательный
- Что происходит: вы ждёте завершения первой корутины, только потом запускается/продолжается следующая.
- Эффект: нет перекрытия по времени; общая длительность примерно
- Когда уместно: есть зависимость по данным (результат первой нужен второй) или нужен строгий порядок.
Конкурентный запуск через
- Что происходит: задачи планируются сразу, начинают выполняться при первой возможности, а вы можете ждать их позже.
- Эффект: время близко к
- Важно:
Ключевые отличия
- Время: последовательный
- Управление: с задачей можно отменять (
- Ошибки: при последовательном
- CPU-bound: конкурентность через
await (корутины выполняются одна за другой)- Что происходит: вы ждёте завершения первой корутины, только потом запускается/продолжается следующая.
- Эффект: нет перекрытия по времени; общая длительность примерно
t1 + t2 + ....- Когда уместно: есть зависимость по данным (результат первой нужен второй) или нужен строгий порядок.
import asyncio
async def a():
await asyncio.sleep(1)
return "A"
async def b():
await asyncio.sleep(1)
return "B"
async def main():
r1 = await a()
r2 = await b()
print(r1, r2) # ~2 секунды
asyncio.run(main())
Конкурентный запуск через
asyncio.create_task (корутины выполняются “параллельно” в рамках одного event loop)- Что происходит: задачи планируются сразу, начинают выполняться при первой возможности, а вы можете ждать их позже.
- Эффект: время близко к
max(t1, t2, ...) (если они в основном I/O-bound и часто await-ят).- Важно:
create_task сам по себе не ждёт завершения; если забыть дождаться — задача может остаться “в фоне”, а исключения станут “Task exception was never retrieved”.import asyncio
async def a():
await asyncio.sleep(1)
return "A"
async def b():
await asyncio.sleep(1)
return "B"
async def main():
t1 = asyncio.create_task(a())
t2 = asyncio.create_task(b())
r1 = await t1
r2 = await t2
print(r1, r2) # ~1 секунда
asyncio.run(main())
Ключевые отличия
- Время: последовательный
await суммирует задержки; create_task позволяет перекрывать ожидания.- Управление: с задачей можно отменять (
task.cancel()), проверять состояние (task.done()), обрабатывать исключения централизованно.- Ошибки: при последовательном
await исключение всплывает сразу в точке ожидания; при create_task исключение проявится когда вы await-нете задачу (или будет предупреждение, если не дождались).- CPU-bound: конкурентность через
create_task не ускоряет тяжёлые вычисления на CPU (нужны процессы/потоки или вынос в executor).
Задача: дождаться завершения нескольких
asyncio.gather
- Назначение: собрать результаты группы корутин/тасков, как “join + collect results”.
- Возврат: список результатов в том же порядке, что и входные аргументы.
- Ошибки: по умолчанию первая поднятая ошибка пробрасывается наружу, остальные корутины при этом не “автоматически отменяются” как обязательное правило; но если вы не обработаете исключение, выполнение вашего кода прервётся, и вы можете не дождаться остальных результатов.
- return_exceptions=True: исключения возвращаются как элементы списка, что удобно для “собрать всё и потом разобрать”.
- Отмена: если отменить сам
Когда выбирать gather:
- нужен список результатов в фиксированном порядке
- типичный “запустить пачку и дождаться всех”
- нужна опция
asyncio.wait
- Назначение: ждать по условию готовности задач и управлять ими (низкоуровневее).
- Вход: обычно уже созданные
- Возврат: два множества:
- Условие ожидания:
-
-
-
- Результаты/исключения: нужно доставать вручную:
- Важно:
Когда выбирать wait:
- нужно частичное ожидание (первая готовая/первая ошибка) и работа с
- нужно реализовать “гонку” задач, таймауты, ранний выход, отмену оставшихся
- нужен контроль над жизненным циклом задач на уровне множеств
Практическое правило выбора
- Нужно просто дождаться всех и получить результаты →
- Нужно дождаться “кого-то первого”, или управлять оставшимися, или разделить done/pending →
Дополнение: для таймаутов часто удобнее использовать
asyncio-задач. Основные инструменты: asyncio.gather и asyncio.wait.asyncio.gather
- Назначение: собрать результаты группы корутин/тасков, как “join + collect results”.
- Возврат: список результатов в том же порядке, что и входные аргументы.
- Ошибки: по умолчанию первая поднятая ошибка пробрасывается наружу, остальные корутины при этом не “автоматически отменяются” как обязательное правило; но если вы не обработаете исключение, выполнение вашего кода прервётся, и вы можете не дождаться остальных результатов.
- return_exceptions=True: исключения возвращаются как элементы списка, что удобно для “собрать всё и потом разобрать”.
- Отмена: если отменить сам
gather, он отменит дочерние задачи.import asyncio
async def work(x):
await asyncio.sleep(x)
return x
async def main():
results = await asyncio.gather(work(1), work(2), work(0.5))
print(results) # [1, 2, 0.5] (порядок как в аргументах)
asyncio.run(main())
Когда выбирать gather:
- нужен список результатов в фиксированном порядке
- типичный “запустить пачку и дождаться всех”
- нужна опция
return_exceptions=True для аккуратного сбора ошибокasyncio.wait
- Назначение: ждать по условию готовности задач и управлять ими (низкоуровневее).
- Вход: обычно уже созданные
Task/Future (корутины лучше предварительно оборачивать в create_task).- Возврат: два множества:
(done, pending).- Условие ожидания:
return_when:-
asyncio.ALL_COMPLETED (по умолчанию) — дождаться всех-
asyncio.FIRST_COMPLETED — дождаться первой завершившейся-
asyncio.FIRST_EXCEPTION — дождаться первой с исключением (или всех, если без исключений)- Результаты/исключения: нужно доставать вручную:
task.result() или task.exception().- Важно:
wait не “собирает” результаты и не пробрасывает исключения сам по себе — это ваша ответственность.import asyncio
async def work(name, t):
await asyncio.sleep(t)
return name
async def main():
tasks = {asyncio.create_task(work("a", 1)),
asyncio.create_task(work("b", 2))}
done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)
for t in done:
print(t.result()) # результат первой завершившейся
for t in pending:
t.cancel()
await asyncio.gather(*pending, return_exceptions=True) # корректно “дочистить” отмену
asyncio.run(main())
Когда выбирать wait:
- нужно частичное ожидание (первая готовая/первая ошибка) и работа с
pending- нужно реализовать “гонку” задач, таймауты, ранний выход, отмену оставшихся
- нужен контроль над жизненным циклом задач на уровне множеств
done/pendingПрактическое правило выбора
- Нужно просто дождаться всех и получить результаты →
asyncio.gather- Нужно дождаться “кого-то первого”, или управлять оставшимися, или разделить done/pending →
asyncio.waitДополнение: для таймаутов часто удобнее использовать
asyncio.wait_for (таймаут на конкретное ожидание) или asyncio.timeout (Python 3.11+), а для “потоковой” обработки результатов по мере готовности — asyncio.as_completed.
asyncio.gather запускает несколько awaitable параллельно и возвращает результаты в порядке аргументов. Обработка исключений зависит от
1) Поведение по умолчанию:
- Первое возникшее исключение пробрасывается наружу из
- Остальные задачи, которые уже выполняются, обычно продолжают выполняться (они не “откатываются” автоматически). Если они позже упадут, и их исключения никто не заберёт, можно получить предупреждения вида Task exception was never retrieved.
- Практика: передавайте в
2) Режим
-
- Вместо этого возвращает список, где на месте упавших корутин будет объект исключения (например,
- Это удобно для “частичного успеха” и агрегации ошибок.
3) Правильная реакция на отмену (
-
- Если используете
4) Короткие рекомендации
- Если нужна стратегия “всё или ничего” и ошибка должна падать наружу:
- Если нужен “best effort” и сбор всех ошибок:
- Для контроля жизненного цикла используйте
return_exceptions.1) Поведение по умолчанию:
return_exceptions=False- Первое возникшее исключение пробрасывается наружу из
await gather(...).- Остальные задачи, которые уже выполняются, обычно продолжают выполняться (они не “откатываются” автоматически). Если они позже упадут, и их исключения никто не заберёт, можно получить предупреждения вида Task exception was never retrieved.
- Практика: передавайте в
gather именно Task (через asyncio.create_task), чтобы затем при ошибке аккуратно дождаться/отменить и “собрать” ошибки остальных.import asyncio
async def worker(n: int):
await asyncio.sleep(0.1 * n)
if n == 2:
raise ValueError("boom")
return n
async def main():
tasks = [asyncio.create_task(worker(i)) for i in range(5)]
try:
res = await asyncio.gather(*tasks) # return_exceptions=False
print(res)
except Exception as e:
# Важно: "собрать" результаты/исключения остальных, чтобы не было un-retrieved
await asyncio.gather(*tasks, return_exceptions=True)
raise # или обработать e
asyncio.run(main())
2) Режим
return_exceptions=True-
gather не бросает исключения наружу.- Вместо этого возвращает список, где на месте упавших корутин будет объект исключения (например,
ValueError(...)).- Это удобно для “частичного успеха” и агрегации ошибок.
import asyncio
async def worker(n: int):
if n == 2:
raise ValueError("boom")
return n
async def main():
results = await asyncio.gather(*(worker(i) for i in range(5)),
return_exceptions=True)
ok = []
errors = []
for r in results:
if isinstance(r, Exception):
errors.append(r)
else:
ok.append(r)
print("ok:", ok)
print("errors:", [type(e).__name__ + ": " + str(e) for e in errors])
asyncio.run(main())
3) Правильная реакция на отмену (
CancelledError)-
asyncio.CancelledError — особое исключение: обычно его не “глотают”, а пробрасывают дальше.- Если используете
return_exceptions=True, отмена тоже придёт как объект исключения; часто её стоит отдельно распознавать и завершать корректно.import asyncio
async def main():
tasks = [asyncio.create_task(asyncio.sleep(10)) for _ in range(3)]
try:
await asyncio.sleep(0.1)
for t in tasks:
t.cancel()
results = await asyncio.gather(*tasks, return_exceptions=True)
for r in results:
if isinstance(r, asyncio.CancelledError):
# обычно это нормальный исход при остановке
pass
except asyncio.CancelledError:
raise
asyncio.run(main())
4) Короткие рекомендации
- Если нужна стратегия “всё или ничего” и ошибка должна падать наружу:
return_exceptions=False + в except обязательно сделайте await gather(..., return_exceptions=True), чтобы собрать хвосты.- Если нужен “best effort” и сбор всех ошибок:
return_exceptions=True и далее фильтруйте Exception.- Для контроля жизненного цикла используйте
asyncio.create_task (а не только голые корутины), чтобы при необходимости отменять/дожидаться явно.
Отмена задач в asyncio: что происходит и как делать правильно
1) Что такое отмена и где она живёт
-
- Отмена в asyncio кооперативная: корутина должна дойти до
2) Когда возникает
- Когда кто-то вызывает
- Когда отменяют ожидание (например,
- При завершении
- При отмене родительской задачи отмена может “протечь” в дочерние, если они создавались и ожидаются в группах (
3) Как выглядит “корректная” отмена задачи
- Инициировать отмену:
- Обязательно дождаться завершения задачи:
- Поймать
4) Почему важно “пробрасывать”
- Если вы поймали
- Допустимый шаблон: сделать cleanup и
5) Что делать с
- Для закрытия ресурсов используйте
6) “Опасные” места: подавление отмены и долгий cleanup
- Если в
- Если вам нужно сделать “неотменяемый” короткий участок cleanup, используйте
7) Рекомендованный способ управления группой задач:
- В Python 3.11+ используйте
8) Практические правила
- Отмена:
- Внутри корутины: ловите
- Не делайте долгий/блокирующий код без
- Для множества задач предпочитайте
1) Что такое отмена и где она живёт
-
Task.cancel() не убивает задачу мгновенно. Он помечает задачу как отменённую и планирует доставку исключения asyncio.CancelledError внутрь корутины при ближайшей “точке ожидания” (await).- Отмена в asyncio кооперативная: корутина должна дойти до
await (или сама проверить флаг) чтобы реально прерваться.2) Когда возникает
asyncio.CancelledError- Когда кто-то вызывает
task.cancel(), и корутина достигает следующего await (например await asyncio.sleep(), await queue.get(), await stream.read() и т.д.).- Когда отменяют ожидание (например,
await asyncio.wait_for(...) по таймауту): отмена прокидывается во внутреннюю задачу/ожидание.- При завершении
asyncio.run() все оставшиеся задачи отменяются: внутри них также прилетает CancelledError.- При отмене родительской задачи отмена может “протечь” в дочерние, если они создавались и ожидаются в группах (
TaskGroup) или если вы явно отменяете их сами.3) Как выглядит “корректная” отмена задачи
- Инициировать отмену:
task.cancel()- Обязательно дождаться завершения задачи:
await task- Поймать
asyncio.CancelledError у места, где вы ожидаете задачу, если вам нужно продолжать работу дальше.import asyncio
async def worker():
try:
while True:
await asyncio.sleep(1) # точка, где придёт CancelledError
except asyncio.CancelledError:
# освободить ресурсы / сделать минимальный cleanup
raise # почти всегда нужно пробросить дальше
async def main():
task = asyncio.create_task(worker())
await asyncio.sleep(2)
task.cancel()
try:
await task
except asyncio.CancelledError:
pass # ожидаемая отмена
asyncio.run(main())
4) Почему важно “пробрасывать”
CancelledError- Если вы поймали
CancelledError и не сделали raise, вы “съедаете” отмену: задача продолжит жить или завершится как успешная, что ломает логику shutdown, таймаутов и отмены.- Допустимый шаблон: сделать cleanup и
raise.5) Что делать с
finally и ресурсами- Для закрытия ресурсов используйте
try/finally (или async with, contextlib.aclosing), т.к. finally выполнится и при отмене.async def use_resource(res):
try:
await res.open()
await asyncio.sleep(10)
finally:
await res.close()
6) “Опасные” места: подавление отмены и долгий cleanup
- Если в
except CancelledError вы делаете долгие операции, вы задерживаете остановку.- Если вам нужно сделать “неотменяемый” короткий участок cleanup, используйте
asyncio.shield() точечно.import asyncio
async def worker():
try:
await asyncio.sleep(999)
except asyncio.CancelledError:
# гарантированно довести важный cleanup до конца
await asyncio.shield(asyncio.sleep(0.1))
raise
7) Рекомендованный способ управления группой задач:
TaskGroup- В Python 3.11+ используйте
asyncio.TaskGroup: он отменяет остальные задачи при ошибке и корректно ждёт завершения всех.import asyncio
async def main():
async with asyncio.TaskGroup() as tg:
tg.create_task(asyncio.sleep(10))
tg.create_task(asyncio.sleep(20))
8) Практические правила
- Отмена:
task.cancel() + await task (и обработка CancelledError снаружи).- Внутри корутины: ловите
CancelledError только для cleanup и затем raise.- Не делайте долгий/блокирующий код без
await: иначе отмена “не доедет”.- Для множества задач предпочитайте
TaskGroup; для таймаутов аккуратно используйте asyncio.wait_for или asyncio.timeout() (3.11+).
Graceful shutdown — это управляемая остановка: перестаём принимать новую работу, даём активным задачам завершиться (или ограничиваем временем), затем гарантированно закрываем ресурсы (соединения, файлы, пулы, сессии).
Рекомендуемый шаблон для Python asyncio
- Один общий сигнал остановки:
- Все фоновые задачи периодически проверяют
- На сигнал ОС (
- Останавливаем «вход» (сервер, consumer, планировщик) чтобы не появлялись новые задачи
- Ждём завершения задач с таймаутом; затем отменяем оставшиеся
- Закрываем ресурсы в
Ключевые практики
- Не блокируйте event loop: вместо долгих операций используйте
- Для внешних ресурсов используйте
- После таймаута делайте
- Для сервера/консьюмера сначала остановите приём, затем дождитесь обработки текущих задач, и только потом закрывайте соединения/пулы.
Если скажете стек (FastAPI/Starlette, aiohttp, Celery, APScheduler, Kafka/RabbitMQ, asyncpg и т.п.), подстрою пример под ваш конкретный сервер и ресурсы.
Рекомендуемый шаблон для Python asyncio
- Один общий сигнал остановки:
asyncio.Event (stop_event)- Все фоновые задачи периодически проверяют
stop_event и корректно выходят- На сигнал ОС (
SIGINT/SIGTERM) выставляем stop_event- Останавливаем «вход» (сервер, consumer, планировщик) чтобы не появлялись новые задачи
- Ждём завершения задач с таймаутом; затем отменяем оставшиеся
- Закрываем ресурсы в
finally или через async withimport asyncio
import signal
from contextlib import AsyncExitStack
async def worker(name: str, stop: asyncio.Event):
try:
while not stop.is_set():
# делаем работу кусками, чтобы можно было остановиться
await asyncio.sleep(0.2)
except asyncio.CancelledError:
# если отменили — быстро завершаем, но можно сделать минимальную уборку
raise
async def main():
stop = asyncio.Event()
loop = asyncio.get_running_loop()
def request_stop():
stop.set()
for sig in (signal.SIGINT, signal.SIGTERM):
try:
loop.add_signal_handler(sig, request_stop)
except NotImplementedError:
# Windows/некоторые окружения: fallback — ловите KeyboardInterrupt снаружи
pass
tasks: set[asyncio.Task] = set()
async with AsyncExitStack() as stack:
# Пример ресурсов (подставьте реальные):
# session = await stack.enter_async_context(aiohttp.ClientSession())
# pool = await stack.enter_async_context(asyncpg.create_pool(...))
# server = await stack.enter_async_context(start_server(...)) # чтобы закрывался автоматически
# Запуск фоновых задач
for i in range(3):
t = asyncio.create_task(worker(f"w{i}", stop))
tasks.add(t)
t.add_done_callback(tasks.discard)
try:
await stop.wait() # ждём сигнал остановки
finally:
# 1) прекращаем принимать новую работу (закройте сервер/consumer/очередь)
# пример: server.close(); await server.wait_closed()
# 2) даём задачам время завершиться сами
try:
await asyncio.wait_for(asyncio.gather(*tasks, return_exceptions=True), timeout=10)
except asyncio.TimeoutError:
# 3) принудительная отмена оставшихся
for t in tasks:
t.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
# 4) ресурсы закроются через AsyncExitStack (LIFO) автоматически
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
pass
Ключевые практики
- Не блокируйте event loop: вместо долгих операций используйте
await, разбивайте работу на итерации, проверяйте stop_event.- Для внешних ресурсов используйте
async with/AsyncExitStack, чтобы закрытие было гарантировано.- После таймаута делайте
cancel() и собирайте результаты через gather(..., return_exceptions=True), чтобы не потерять исключения и не оставить «висящие» задачи.- Для сервера/консьюмера сначала остановите приём, затем дождитесь обработки текущих задач, и только потом закрывайте соединения/пулы.
Если скажете стек (FastAPI/Starlette, aiohttp, Celery, APScheduler, Kafka/RabbitMQ, asyncpg и т.п.), подстрою пример под ваш конкретный сервер и ресурсы.
Таймауты в asyncio — это ограничения по времени на выполнение await-операций (ожидание корутины, Future, I/O). Если операция не завершилась за заданный срок, её обычно отменяют, чтобы не зависать бесконечно и корректно обрабатывать медленные/подвисшие внешние ресурсы (HTTP, БД, очереди, сокеты).
asyncio.wait_for — обёртка, которая ждёт завершения awaitable не дольше
-
- Важно: отмена может занять немного времени, если корутина подавляет отмену или делает долгую работу без точек
Базовый шаблон
Прикладные сценарии использования
1) Таймаут на внешний запрос (HTTP/DB/RPC)
Когда библиотека не даёт удобных таймаутов или нужен общий “зонтик”:
2) Таймаут на один шаг пайплайна/бизнес-операции
Например, расчёт, который не должен блокировать обработку запроса пользователя:
3) Ограничение ожидания элемента из очереди
Чтобы воркер мог периодически проверять флаг остановки/делать housekeeping:
4) Таймаут при “гонке” задач
Если результат нужен быстро, иначе возвращаем частичный/пустой:
Практические советы
- Ловите именно
- Если внутри корутины есть ресурсы (соединение, файл, локи), используйте
- Для сетевых клиентов чаще лучше использовать встроенные таймауты клиента (они точнее и могут различать connect/read/write), а
asyncio.wait_for — обёртка, которая ждёт завершения awaitable не дольше
timeout секунд. Если время вышло, возбуждается asyncio.TimeoutError, а ожидаемая задача отменяется (в неё пробрасывается CancelledError). -
timeout=None — ждать без ограничения. - Важно: отмена может занять немного времени, если корутина подавляет отмену или делает долгую работу без точек
await.Базовый шаблон
import asyncio
async def do_work():
await asyncio.sleep(2)
return "ok"
async def main():
try:
result = await asyncio.wait_for(do_work(), timeout=1.0)
print(result)
except asyncio.TimeoutError:
print("timeout")
asyncio.run(main())Прикладные сценарии использования
1) Таймаут на внешний запрос (HTTP/DB/RPC)
Когда библиотека не даёт удобных таймаутов или нужен общий “зонтик”:
import asyncio
async def fetch_data():
# здесь мог бы быть вызов клиента БД/HTTP и т.п.
await asyncio.sleep(5)
return {"data": 123}
async def main():
try:
data = await asyncio.wait_for(fetch_data(), timeout=0.8)
return data
except asyncio.TimeoutError:
# логирование, метрики, деградация сервиса, фолбэк
return {"data": None, "error": "timeout"}2) Таймаут на один шаг пайплайна/бизнес-операции
Например, расчёт, который не должен блокировать обработку запроса пользователя:
import asyncio
async def compute():
await asyncio.sleep(3)
return 42
async def handle_request():
try:
return await asyncio.wait_for(compute(), timeout=0.2)
except asyncio.TimeoutError:
return "try again later"3) Ограничение ожидания элемента из очереди
Чтобы воркер мог периодически проверять флаг остановки/делать housekeeping:
import asyncio
async def worker(q: asyncio.Queue, stop: asyncio.Event):
while not stop.is_set():
try:
item = await asyncio.wait_for(q.get(), timeout=0.5)
except asyncio.TimeoutError:
continue
try:
# обработка item
await asyncio.sleep(0.1)
finally:
q.task_done()4) Таймаут при “гонке” задач
Если результат нужен быстро, иначе возвращаем частичный/пустой:
import asyncio
async def fast():
await asyncio.sleep(0.1)
return "fast"
async def slow():
await asyncio.sleep(2)
return "slow"
async def main():
task = asyncio.create_task(slow())
try:
return await asyncio.wait_for(task, timeout=0.3)
except asyncio.TimeoutError:
return await fast()Практические советы
- Ловите именно
asyncio.TimeoutError вокруг wait_for. - Если внутри корутины есть ресурсы (соединение, файл, локи), используйте
try/finally, чтобы корректно освобождать их при отмене. - Для сетевых клиентов чаще лучше использовать встроенные таймауты клиента (они точнее и могут различать connect/read/write), а
wait_for применять как общий верхний предел на всю операцию.
Идея:
Замечания:
-
- Если важно не создавать сразу 1000 задач (экономия памяти), скажи — покажу вариант с очередью (
asyncio.Semaphore(N) ограничивает число одновременно выполняющихся корутин, которые входят в критическую секцию через async with sem:. Для 1000 URL и лимита 50 создаём семафор на 50 и оборачиваем сетевой запрос.import asyncio
import aiohttp
CONCURRENCY = 50
async def fetch(url: str, session: aiohttp.ClientSession, sem: asyncio.Semaphore) -> tuple[str, int]:
async with sem: # гарантирует не более CONCURRENCY одновременных fetch
async with session.get(url, timeout=aiohttp.ClientTimeout(total=30)) as resp:
await resp.read() # или resp.text()/json()
return url, resp.status
async def main(urls: list[str]) -> list[tuple[str, int]]:
sem = asyncio.Semaphore(CONCURRENCY)
async with aiohttp.ClientSession() as session:
tasks = [asyncio.create_task(fetch(url, session, sem)) for url in urls]
results = await asyncio.gather(*tasks, return_exceptions=True)
return results
if __name__ == "__main__":
urls = [f"https://example.com/?q={i}" for i in range(1000)]
out = asyncio.run(main(urls))
print(out[:5])
Замечания:
-
tasks можно создать сразу на 1000: семафор ограничит реальную параллельность до 50.- Если важно не создавать сразу 1000 задач (экономия памяти), скажи — покажу вариант с очередью (
asyncio.Queue) и воркерами.
asyncio.Lock — это асинхронный мьютекс для защиты критической секции от одновременного доступа нескольких корутин. Нужен, когда есть общий ресурс (память/кэш/файл/соединение), который нельзя корректно изменять параллельно, иначе будут гонки данных.
- Гарантия: в критическую секцию одновременно входит только одна корутина
- Применение: атомарное обновление общей структуры, ограничение параллельной записи, последовательный доступ к одному клиенту/соединению
asyncio.Event — это асинхронный сигнал (флаг) для координации: одни корутины ждут наступления события, другая корутина его устанавливает. Нужен не для защиты ресурса, а для оповещения “можно продолжать”.
- Семантика:
- Применение: “сервис готов”, “получены данные”, “старт/стоп” для воркеров, graceful shutdown
Ключевое различие
- Lock: взаимоисключение при доступе к ресурсу (1 корутина за раз)
- Event: сигнализация о наступлении условия (многие корутины могут “проснуться” одновременно)
- Гарантия: в критическую секцию одновременно входит только одна корутина
- Применение: атомарное обновление общей структуры, ограничение параллельной записи, последовательный доступ к одному клиенту/соединению
import asyncio
balance = 0
lock = asyncio.Lock()
async def deposit(amount: int):
global balance
async with lock: # критическая секция
tmp = balance
await asyncio.sleep(0) # имитация переключения контекста
balance = tmp + amount
async def main():
await asyncio.gather(*(deposit(1) for _ in range(1000)))
print(balance) # всегда 1000
asyncio.run(main())
asyncio.Event — это асинхронный сигнал (флаг) для координации: одни корутины ждут наступления события, другая корутина его устанавливает. Нужен не для защиты ресурса, а для оповещения “можно продолжать”.
- Семантика:
wait() блокирует до set(); после set() все ожидающие просыпаются; флаг остаётся установленным, пока не вызвать clear() - Применение: “сервис готов”, “получены данные”, “старт/стоп” для воркеров, graceful shutdown
import asyncio
ready = asyncio.Event()
async def worker():
await ready.wait() # ждём, пока система будет готова
print("worker: started work")
async def initializer():
await asyncio.sleep(0.2)
ready.set() # оповещаем всех ожидающих
async def main():
await asyncio.gather(worker(), worker(), initializer())
asyncio.run(main())
Ключевое различие
- Lock: взаимоисключение при доступе к ресурсу (1 корутина за раз)
- Event: сигнализация о наступлении условия (многие корутины могут “проснуться” одновременно)
Разница между
- Модель конкурентности:
- API ожидания: у
- Нельзя смешивать напрямую:
- Сигнал завершения: в обоих случаях часто используют “sentinel” (специальный объект) или отмену задач; в
Producer-consumer с
- Producer генерирует элементы и делает
- Consumer в цикле делает
- Для корректного завершения обычно:
- либо отправляют по одному “sentinel” на каждого consumer,
- либо отменяют consumer-задачи после
Ключевые моменты:
- Backpressure: задавайте
- q.join(): работает вместе с
- Не используйте
asyncio.Queue и queue.Queue:- Модель конкурентности:
queue.Queue рассчитана на потоки (threading) и блокирует поток при ожидании; asyncio.Queue рассчитана на корутины в одном event loop и при ожидании не блокирует поток, а отдаёт управление циклу событий.- API ожидания: у
queue.Queue методы put()/get() блокирующие (или с timeout); у asyncio.Queue await q.put()/await q.get().- Нельзя смешивать напрямую:
asyncio.Queue не предназначена для безопасного использования из разных потоков; queue.Queue не “awaitable” и при использовании в async-коде будет блокировать event loop (плохо).- Сигнал завершения: в обоих случаях часто используют “sentinel” (специальный объект) или отмену задач; в
asyncio удобнее применять q.join() + task_done() для ожидания обработки всех элементов.Producer-consumer с
asyncio.Queue:- Producer генерирует элементы и делает
await q.put(item).- Consumer в цикле делает
item = await q.get(), обрабатывает, затем вызывает q.task_done().- Для корректного завершения обычно:
- либо отправляют по одному “sentinel” на каждого consumer,
- либо отменяют consumer-задачи после
await q.join().import asyncio
SENTINEL = object()
async def producer(name: str, q: asyncio.Queue, n: int) -> None:
for i in range(n):
item = (name, i)
await q.put(item) # не блокирует поток, а "ждёт" в event loop
await asyncio.sleep(0.05) # имитация I/O
# producer завершился (sentinel отправим снаружи, когда завершатся все producers)
async def consumer(name: str, q: asyncio.Queue) -> None:
while True:
item = await q.get()
try:
if item is SENTINEL:
return
# обработка
await asyncio.sleep(0.1) # имитация I/O
# print(f"{name} processed {item}")
finally:
q.task_done()
async def main() -> None:
q: asyncio.Queue = asyncio.Queue(maxsize=100) # maxsize даёт backpressure
consumers = [asyncio.create_task(consumer(f"c{i}", q)) for i in range(3)]
producers = [asyncio.create_task(producer(f"p{i}", q, n=5)) for i in range(2)]
# ждём, пока producers закончат наполнение
await asyncio.gather(*producers)
# ждём, пока очередь будет полностью обработана consumers
await q.join()
# корректно останавливаем consumers: по одному sentinel на consumer
for _ in consumers:
await q.put(SENTINEL)
await asyncio.gather(*consumers)
if __name__ == "__main__":
asyncio.run(main())
Ключевые моменты:
- Backpressure: задавайте
maxsize, чтобы producer “притормаживал” через await q.put(), когда consumers не успевают.- q.join(): работает вместе с
task_done() и позволяет дождаться обработки всех задач до остановки consumers.- Не используйте
queue.Queue внутри async-кода без вынесения в отдельный поток (иначе блокируется event loop).
Идея: конвейер из стадий download → parse → write, между стадиями asyncio.Queue. На каждой стадии несколько воркеров. Завершение — через sentinel (специальный объект), который прокидывается дальше по конвейеру.
Ключевые правила
-
- Каждый воркер делает:
- Для корректного завершения: кладём в входную очередь
- вызывает
- пересылает sentinel дальше (если стадия не последняя)
- завершает цикл.
Практические замечания
- Если скачивание через
- Если запись/парсинг блокируют CPU/IO синхронно: выносите в
- Оборачивайте сетевые операции в
Ключевые правила
-
Queue(maxsize=...) даёт backpressure (ограничивает память и скорость).- Каждый воркер делает:
item = await q.get() → обработка → q.task_done().- Для корректного завершения: кладём в входную очередь
N sentinel-объектов, где N — число воркеров стадии; получив sentinel, воркер:- вызывает
task_done()- пересылает sentinel дальше (если стадия не последняя)
- завершает цикл.
import asyncio
from dataclasses import dataclass
from typing import Iterable, Optional
SENTINEL = object()
@dataclass(frozen=True)
class Task:
url: str
@dataclass(frozen=True)
class Page:
url: str
html: str
@dataclass(frozen=True)
class Record:
url: str
title: str
async def fetch(url: str) -> str:
# Заглушка: здесь обычно aiohttp/anyio/httpx
await asyncio.sleep(0.05)
return f"<html><title>{url}</title></html>"
def parse_html(url: str, html: str) -> Record:
# Заглушка: реальный парсинг через lxml/bs4/regex
title = url
return Record(url=url, title=title)
async def write_record(rec: Record) -> None:
# Заглушка: запись в БД/файл. Для синхронной БД — asyncio.to_thread(...)
await asyncio.sleep(0.01)
async def downloader(in_q: asyncio.Queue, out_q: asyncio.Queue) -> None:
while True:
item = await in_q.get()
try:
if item is SENTINEL:
await out_q.put(SENTINEL)
return
task: Task = item
html = await fetch(task.url)
await out_q.put(Page(url=task.url, html=html))
finally:
in_q.task_done()
async def parser(in_q: asyncio.Queue, out_q: asyncio.Queue) -> None:
while True:
item = await in_q.get()
try:
if item is SENTINEL:
await out_q.put(SENTINEL)
return
page: Page = item
rec = parse_html(page.url, page.html)
await out_q.put(rec)
finally:
in_q.task_done()
async def writer(in_q: asyncio.Queue) -> None:
while True:
item = await in_q.get()
try:
if item is SENTINEL:
return
rec: Record = item
await write_record(rec)
finally:
in_q.task_done()
async def run_pipeline(
urls: Iterable[str],
n_downloaders: int = 10,
n_parsers: int = 4,
n_writers: int = 2,
) -> None:
q_tasks: asyncio.Queue = asyncio.Queue(maxsize=1000)
q_pages: asyncio.Queue = asyncio.Queue(maxsize=200)
q_records: asyncio.Queue = asyncio.Queue(maxsize=200)
downloaders = [asyncio.create_task(downloader(q_tasks, q_pages)) for _ in range(n_downloaders)]
parsers = [asyncio.create_task(parser(q_pages, q_records)) for _ in range(n_parsers)]
writers = [asyncio.create_task(writer(q_records)) for _ in range(n_writers)]
for url in urls:
await q_tasks.put(Task(url=url))
for _ in range(n_downloaders):
await q_tasks.put(SENTINEL)
await q_tasks.join()
await q_pages.join()
await q_records.join()
await asyncio.gather(*downloaders)
await asyncio.gather(*parsers)
await asyncio.gather(*writers)
async def main() -> None:
urls = [f"https://example.com/{i}" for i in range(100)]
await run_pipeline(urls)
if __name__ == "__main__":
asyncio.run(main())
Практические замечания
- Если скачивание через
aiohttp: держите ClientSession один на все downloader-воркеры и используйте TCPConnector(limit=...).- Если запись/парсинг блокируют CPU/IO синхронно: выносите в
asyncio.to_thread(...) или отдельный пул.- Оборачивайте сетевые операции в
asyncio.wait_for и добавляйте ретраи, иначе воркеры могут зависать на одном URL.
Асинхронный контекстный менеджер — это объект, который гарантирует корректное выделение и освобождение ресурса в асинхронном коде (даже при исключениях), не блокируя цикл событий. Используется с
Зачем нужен
- Надёжное закрытие/освобождение ресурсов: соединений, транзакций, файлов, локов, семафоров.
- Автоматическая обработка ошибок: гарантированный выход из контекста при исключении.
- Неблокирующий вход/выход: если открытие/закрытие требует await (сетевые операции), это делается корректно.
Как это работает
Объект для
-
-
Эквивалент по смыслу:
-
Практический пример: соединение/клиент
Важно: с
Пример: пул соединений (типичный паттерн)
- берём соединение из пула
- делаем работу
- автоматически возвращаем соединение в пул при выходе из
Типичные случаи
-
-
-
Ключевая выгода: ресурс не “утечёт”, а освобождение произойдёт всегда и асинхронно, что критично для соединений и IO в высоконагруженных сервисах.
async with.Зачем нужен
async with- Надёжное закрытие/освобождение ресурсов: соединений, транзакций, файлов, локов, семафоров.
- Автоматическая обработка ошибок: гарантированный выход из контекста при исключении.
- Неблокирующий вход/выход: если открытие/закрытие требует await (сетевые операции), это делается корректно.
Как это работает
Объект для
async with должен реализовать:-
__aenter__(self) — выполняется при входе (с await).-
__aexit__(self, exc_type, exc, tb) — выполняется при выходе (тоже с await), где обычно делается закрытие/rollback/release.Эквивалент по смыслу:
-
async with cm as x: ... ≈ x = await cm.__aenter__() и затем await cm.__aexit__(...) в finally.class AsyncCM:
async def __aenter__(self):
# выделяем ресурс (может быть сеть/IO)
return self
async def __aexit__(self, exc_type, exc, tb):
# освобождаем ресурс (закрытие/rollback/release)
return False # не подавлять исключение
async def main():
async with AsyncCM() as cm:
passПрактический пример: соединение/клиент
Важно: с
async with соединение закроется даже если внутри возникнет ошибка.import aiohttp
import asyncio
async def fetch(url: str) -> str:
async with aiohttp.ClientSession() as session:
async with session.get(url) as resp:
resp.raise_for_status()
return await resp.text()
async def main():
html = await fetch("https://example.com")
print(len(html))
asyncio.run(main())Пример: пул соединений (типичный паттерн)
- берём соединение из пула
- делаем работу
- автоматически возвращаем соединение в пул при выходе из
async withclass Pool:
async def acquire(self): ...
async def release(self, conn): ...
class Acquire:
def __init__(self, pool: Pool):
self.pool = pool
self.conn = None
async def __aenter__(self):
self.conn = await self.pool.acquire()
return self.conn
async def __aexit__(self, exc_type, exc, tb):
await self.pool.release(self.conn)
return False
async def use_pool(pool: Pool):
async with Acquire(pool) as conn:
await conn.execute("SELECT 1")Типичные случаи
-
async with lock: — безопасная синхронизация в asyncio.-
async with session/connection: — корректное открытие/закрытие сетевых ресурсов.-
async with transaction: — commit/rollback автоматически (rollback при исключении).Ключевая выгода: ресурс не “утечёт”, а освобождение произойдёт всегда и асинхронно, что критично для соединений и IO в высоконагруженных сервисах.
Асинхронные итераторы — это протокол итерации, где получение следующего элемента может требовать await (I/O, таймеры, сеть), поэтому цикл не блокирует event loop.
Как работает
-
- вызывается
- далее по кругу вызывается
- когда
- внутри цикла тело выполняется как обычно; точки ожидания — в
Минимальная реализация
Чем отличается от обычного
- обычный итератор:
- асинхронный итератор:
Где встречается в реальных библиотеках
- WebSockets (
- пример:
- HTTP/Streaming (
- пример:
- Базы данных (async-драйверы, напр.
- пример: итерация по курсору, где строки подтягиваются порциями по сети
- Очереди и каналы (
- паттерн: генератор/итератор, который ждёт данные и отдаёт их по мере появления
- Фреймворки событий/сообщений (клиенты Kafka/RabbitMQ, SSE/long-polling): поток событий как асинхронная последовательность
Когда использовать
- когда элементы появляются со временем и получение следующего элемента — это I/O или ожидание события (сеть, файлы, очереди, подписки)
- когда нужно обрабатывать поток данных без загрузки всего набора в память (streaming/backpressure)
Практическая заметка
- асинхронный генератор (
Как работает
async for-
async for x in it: ожидает, что it реализует асинхронный итератор:- вызывается
aiter = it.__aiter__()- далее по кругу вызывается
await aiter.__anext__()- когда
__anext__ выбрасывает StopAsyncIteration, цикл завершён- внутри цикла тело выполняется как обычно; точки ожидания — в
await __anext__() и внутри тела циклаМинимальная реализация
class Counter:
def __init__(self, n):
self.n = n
self.i = 0
def __aiter__(self):
return self
async def __anext__(self):
if self.i >= self.n:
raise StopAsyncIteration
self.i += 1
return self.i
async def main():
async for x in Counter(3):
print(x)Чем отличается от обычного
for- обычный итератор:
__iter__/__next__ без await- асинхронный итератор:
__aiter__/__anext__, где __anext__ — coroutine и требует awaitГде встречается в реальных библиотеках
- WebSockets (
websockets): чтение входящих сообщений потоково- пример:
async for msg in websocket: — каждое сообщение приходит из сети асинхронно- HTTP/Streaming (
aiohttp, httpx): потоковое чтение тела ответа- пример:
async for chunk in response.content.iter_chunked(...): или async for chunk in response.aiter_bytes():- Базы данных (async-драйверы, напр.
asyncpg, async-ORM/клиенты): курсоры/стриминг результатов- пример: итерация по курсору, где строки подтягиваются порциями по сети
- Очереди и каналы (
asyncio, anyio): потребление элементов из async-очередей/стримов- паттерн: генератор/итератор, который ждёт данные и отдаёт их по мере появления
- Фреймворки событий/сообщений (клиенты Kafka/RabbitMQ, SSE/long-polling): поток событий как асинхронная последовательность
Когда использовать
- когда элементы появляются со временем и получение следующего элемента — это I/O или ожидание события (сеть, файлы, очереди, подписки)
- когда нужно обрабатывать поток данных без загрузки всего набора в память (streaming/backpressure)
Практическая заметка
- асинхронный генератор (
async def ...: yield ...) автоматически является асинхронным итератором и часто удобнее ручной реализации класса.
Async-генератор — это функция, объявленная как
Минимальный пример
Когда async-генераторы удобнее, чем возврат списка целиком
- Стриминг результатов: можно отдавать элементы по мере готовности, не дожидаясь окончания всей работы (например, чтение построчно из сети/файла, поток событий, пагинация API).
- Экономия памяти: список требует хранить все элементы сразу; генератор держит только текущее состояние итерации.
- Низкая задержка (latency): потребитель начинает обрабатывать первые элементы сразу (pipeline), что особенно полезно для ETL, парсинга больших данных, отправки результатов в клиент (SSE/WebSocket).
- Естественная интеграция с I/O: внутри можно делать
- Бесконечные/долго живущие источники: события, подписки, очереди, tail логов — список в принципе не подходит.
Когда лучше вернуть список
- Нужны повторные проходы по данным или произвольный доступ по индексу.
- Требуется агрегировать/отсортировать/перемешать всё целиком перед использованием.
- Объём данных небольшой, и важнее простота интерфейса.
Практический ориентир: если результат можно потреблять по частям и он появляется постепенно (I/O, поток), выбирай async-генератор; если результат нужен целиком и сразу — список.
async def, которая внутри использует yield и возвращает асинхронный итератор. Его элементы получают через async for, а при необходимости — вручную через __anext__(). Такой генератор может приостанавливать выполнение на await между выдачей элементов, не блокируя event loop.Минимальный пример
import asyncio
async def gen():
for i in range(3):
await asyncio.sleep(1) # имитация I/O
yield i
async def main():
async for x in gen():
print(x)
asyncio.run(main())Когда async-генераторы удобнее, чем возврат списка целиком
- Стриминг результатов: можно отдавать элементы по мере готовности, не дожидаясь окончания всей работы (например, чтение построчно из сети/файла, поток событий, пагинация API).
- Экономия памяти: список требует хранить все элементы сразу; генератор держит только текущее состояние итерации.
- Низкая задержка (latency): потребитель начинает обрабатывать первые элементы сразу (pipeline), что особенно полезно для ETL, парсинга больших данных, отправки результатов в клиент (SSE/WebSocket).
- Естественная интеграция с I/O: внутри можно делать
await между yield, сохраняя отзывчивость приложения.- Бесконечные/долго живущие источники: события, подписки, очереди, tail логов — список в принципе не подходит.
Когда лучше вернуть список
- Нужны повторные проходы по данным или произвольный доступ по индексу.
- Требуется агрегировать/отсортировать/перемешать всё целиком перед использованием.
- Объём данных небольшой, и важнее простота интерфейса.
Практический ориентир: если результат можно потреблять по частям и он появляется постепенно (I/O, поток), выбирай async-генератор; если результат нужен целиком и сразу — список.
Задача: встроить блокирующий (sync) код в asyncio, не останавливая event loop. Для этого блокирующее выполняют в другом потоке (реже в процессе), а в asyncio ждут результат через
1)
- Что делает: запускает функцию в потоке из стандартного thread pool, возвращает awaitable.
- Плюсы: самый простой API; автоматически пробрасывает
- Когда выбирать: Python 3.9+, нужен быстрый и простой способ «вынести в поток» блокирующую I/O-логику или библиотеку без async API.
2)
- Что делает: запускает функцию в указанном executor (обычно
- Плюсы: полный контроль: свой пул потоков (размер, именование, жизненный цикл), можно использовать процессы для CPU-bound, можно шарить executor между частями приложения.
- Когда выбирать: нужен контроль пула (ограничить параллелизм, изолировать «тяжёлые» блокирующие задачи), или требуется ProcessPoolExecutor.
Как выбирать:
- Просто вынести sync I/O в поток и забыть:
- Нужно ограничить/разделить ресурсы (свой пул, max_workers), управлять жизненным циклом:
- CPU-bound вычисления:
Практические замечания:
- Ограничивайте параллелизм (иначе легко создать сотни потоковых задач): используйте
- Отмена:
- Если есть async-библиотека для задачи (HTTP, DB) — используйте её вместо потоков: это эффективнее и проще в сопровождении.
await.1)
asyncio.to_thread()- Что делает: запускает функцию в потоке из стандартного thread pool, возвращает awaitable.
- Плюсы: самый простой API; автоматически пробрасывает
contextvars; хороший выбор по умолчанию для редких/умеренных блокирующих вызовов.- Когда выбирать: Python 3.9+, нужен быстрый и простой способ «вынести в поток» блокирующую I/O-логику или библиотеку без async API.
import asyncio
import time
def blocking_call(x: int) -> int:
time.sleep(1)
return x * 2
async def main():
result = await asyncio.to_thread(blocking_call, 21)
print(result)
asyncio.run(main())2)
loop.run_in_executor()- Что делает: запускает функцию в указанном executor (обычно
ThreadPoolExecutor, иногда ProcessPoolExecutor), возвращает Future, который await’ится.- Плюсы: полный контроль: свой пул потоков (размер, именование, жизненный цикл), можно использовать процессы для CPU-bound, можно шарить executor между частями приложения.
- Когда выбирать: нужен контроль пула (ограничить параллелизм, изолировать «тяжёлые» блокирующие задачи), или требуется ProcessPoolExecutor.
import asyncio
from concurrent.futures import ThreadPoolExecutor
import time
def blocking_io():
time.sleep(1)
return "ok"
async def main():
loop = asyncio.get_running_loop()
with ThreadPoolExecutor(max_workers=4) as ex:
result = await loop.run_in_executor(ex, blocking_io)
print(result)
asyncio.run(main())Как выбирать:
- Просто вынести sync I/O в поток и забыть:
asyncio.to_thread().- Нужно ограничить/разделить ресурсы (свой пул, max_workers), управлять жизненным циклом:
run_in_executor() с явным ThreadPoolExecutor.- CPU-bound вычисления:
run_in_executor() с ProcessPoolExecutor (потоки не ускорят из-за GIL, но помогут не блокировать loop).import asyncio
from concurrent.futures import ProcessPoolExecutor
def cpu_bound(n: int) -> int:
s = 0
for i in range(n):
s += i * i
return s
async def main():
loop = asyncio.get_running_loop()
with ProcessPoolExecutor() as ex:
result = await loop.run_in_executor(ex, cpu_bound, 10_000_00)
print(result)
asyncio.run(main())Практические замечания:
- Ограничивайте параллелизм (иначе легко создать сотни потоковых задач): используйте
ThreadPoolExecutor(max_workers=...) или семафор вокруг to_thread.- Отмена:
await можно отменить, но сам блокирующий вызов в потоке обычно не прерывается мгновенно; проектируйте функции с таймаутами/флагами остановки.- Если есть async-библиотека для задачи (HTTP, DB) — используйте её вместо потоков: это эффективнее и проще в сопровождении.
Почему вредят
- Event loop в asyncio по сути однопоточный: он должен быстро переключаться между задачами и обрабатывать готовые события (I/O, таймеры, колбэки).
- Любой блокирующий участок кода останавливает поток, в котором крутится loop, поэтому:
-
-
- Тяжёлые вычисления держат GIL и CPU → loop не успевает “прокрутиться”, растёт задержка реакции.
- Итог: рост latency, таймауты, скачки нагрузки, “залипания” веб-сервера, неверная работа периодических задач.
Как диагностировать
- Включить режим отладки asyncio и ловить “медленные” колбэки:
- Включить WARNING про долгие шаги loop через переменную окружения:
- Смотреть задержку тиков (простой мониторинг “дрейфа” планировщика): если loop блокируется, фактическая пауза будет намного больше ожидаемой.
- Профилировать где тратится CPU/время:
- для CPU:
- для “кто блокирует поток”: снятие стека по сигналу (например,
- Логи веб-сервера (uvicorn/gunicorn): рост времени ответа, зависание keep-alive, “timeout handling request” указывает на блокировку event loop.
Что делать вместо
-
-
- CPU-bound → вынос в пул:
- Event loop в asyncio по сути однопоточный: он должен быстро переключаться между задачами и обрабатывать готовые события (I/O, таймеры, колбэки).
- Любой блокирующий участок кода останавливает поток, в котором крутится loop, поэтому:
-
time.sleep() блокирует поток целиком → не выполняются другие корутины, таймеры, обработка сокетов.-
requests (синхронный HTTP) делает блокирующий I/O → loop “мертвеет” на время запроса.- Тяжёлые вычисления держат GIL и CPU → loop не успевает “прокрутиться”, растёт задержка реакции.
- Итог: рост latency, таймауты, скачки нагрузки, “залипания” веб-сервера, неверная работа периодических задач.
Как диагностировать
- Включить режим отладки asyncio и ловить “медленные” колбэки:
import asyncio
import logging
logging.basicConfig(level=logging.WARNING)
loop = asyncio.get_event_loop()
loop.set_debug(True)
loop.slow_callback_duration = 0.05 # 50ms
- Включить WARNING про долгие шаги loop через переменную окружения:
# в окружении процесса
PYTHONASYNCIODEBUG=1
- Смотреть задержку тиков (простой мониторинг “дрейфа” планировщика): если loop блокируется, фактическая пауза будет намного больше ожидаемой.
import asyncio, time
async def monitor_loop(period=0.1):
while True:
t0 = time.perf_counter()
await asyncio.sleep(period)
dt = time.perf_counter() - t0
lag = dt - period
if lag > 0.05:
print(f"loop lag: {lag:.3f}s (dt={dt:.3f}s)")
asyncio.run(monitor_loop())
- Профилировать где тратится CPU/время:
- для CPU:
py-spy top/py-spy record (без остановки процесса) или cProfile- для “кто блокирует поток”: снятие стека по сигналу (например,
faulthandler) и поиск мест с time.sleep/requests/долгими циклами- Логи веб-сервера (uvicorn/gunicorn): рост времени ответа, зависание keep-alive, “timeout handling request” указывает на блокировку event loop.
Что делать вместо
-
time.sleep → await asyncio.sleep-
requests → httpx.AsyncClient или aiohttp- CPU-bound → вынос в пул:
import asyncio
from concurrent.futures import ProcessPoolExecutor
def heavy():
# CPU-bound
...
async def main():
loop = asyncio.get_running_loop()
with ProcessPoolExecutor() as ex:
result = await loop.run_in_executor(ex, heavy)
asyncio.run(main())
Подходы к cron-like задачам внутри asyncio (без внешнего планировщика)
- Простой периодический цикл: фиксированный интервал между запусками (подходит для every N seconds/minutes).
- Выравнивание по “стеночным” границам времени: запуск, например, ровно в начале минуты/часа/суток (ближе к cron).
- Устойчивость: не допускать одновременных запусков одной задачи, обрабатывать отмену, логировать исключения, не накапливать дрейф.
1) Периодическая задача с фиксированным интервалом (без дрейфа)
2) “Cron-подобно”: запуск в точное время (например, каждую минуту в 00 секунд)
Идея: считать, сколько осталось до следующей границы, спать до неё, затем запускать.
3) Защита от параллельных запусков (если задача может выполняться дольше периода)
- Skip-if-running: пропускать запуск, если предыдущий ещё идёт.
Практические рекомендации
- Используйте
- Всегда обрабатывайте
- Не давайте исключениям “убивать” планировщик: ловите и логируйте их внутри цикла.
- Если нужен настоящий cron-синтаксис (5 0 * * * и т.п.) без внешнего демона — чаще берут библиотеку уровня
Если скажете, какие выражения вам нужны (каждые N секунд, “в 03:00”, “по будням”, таймзона), подберу компактную реализацию именно под ваш шаблон.
- Простой периодический цикл: фиксированный интервал между запусками (подходит для every N seconds/minutes).
- Выравнивание по “стеночным” границам времени: запуск, например, ровно в начале минуты/часа/суток (ближе к cron).
- Устойчивость: не допускать одновременных запусков одной задачи, обрабатывать отмену, логировать исключения, не накапливать дрейф.
1) Периодическая задача с фиксированным интервалом (без дрейфа)
import asyncio
import time
from contextlib import suppress
async def run_periodic(interval_s: float, job, *, name: str = "job"):
next_t = time.monotonic()
while True:
next_t += interval_s
try:
await job()
except Exception as e:
# замените на нормальный логгер
print(f"{name} failed: {e!r}")
delay = next_t - time.monotonic()
if delay > 0:
await asyncio.sleep(delay)
else:
# если отстали, не пытаемся "догонять" серией запусков
next_t = time.monotonic()
# пример работы
async def cleanup():
print("cleanup...")
async def main():
task = asyncio.create_task(run_periodic(60, cleanup, name="cleanup"))
try:
await asyncio.Event().wait()
finally:
task.cancel()
with suppress(asyncio.CancelledError):
await task
asyncio.run(main())
2) “Cron-подобно”: запуск в точное время (например, каждую минуту в 00 секунд)
Идея: считать, сколько осталось до следующей границы, спать до неё, затем запускать.
import asyncio
from datetime import datetime, timedelta, timezone
from contextlib import suppress
def seconds_until_next_minute(now: datetime) -> float:
next_min = (now.replace(second=0, microsecond=0) + timedelta(minutes=1))
return (next_min - now).total_seconds()
async def run_each_minute(job, *, tz=timezone.utc, name="job"):
while True:
now = datetime.now(tz)
await asyncio.sleep(seconds_until_next_minute(now))
try:
await job()
except Exception as e:
print(f"{name} failed: {e!r}")
async def job():
print("minute tick", datetime.now(timezone.utc).isoformat())
async def main():
t = asyncio.create_task(run_each_minute(job, name="minute_job"))
try:
await asyncio.Event().wait()
finally:
t.cancel()
with suppress(asyncio.CancelledError):
await t
asyncio.run(main())
3) Защита от параллельных запусков (если задача может выполняться дольше периода)
- Skip-if-running: пропускать запуск, если предыдущий ещё идёт.
import asyncio
class SingleRunner:
def __init__(self):
self._lock = asyncio.Lock()
async def run(self, coro):
if self._lock.locked():
return # пропускаем
async with self._lock:
await coro
runner = SingleRunner()
async def job():
await asyncio.sleep(90)
async def wrapped_job():
await runner.run(job())
Практические рекомендации
- Используйте
time.monotonic() для интервалов (устойчиво к изменению системного времени), а datetime.now(tz) — для “календарных” границ (cron-like).- Всегда обрабатывайте
asyncio.CancelledError при остановке приложения и отменяйте фоновые задачи.- Не давайте исключениям “убивать” планировщик: ловите и логируйте их внутри цикла.
- Если нужен настоящий cron-синтаксис (5 0 * * * и т.п.) без внешнего демона — чаще берут библиотеку уровня
APScheduler, но это уже “встроенный” планировщик, не отдельный сервис.Если скажете, какие выражения вам нужны (каждые N секунд, “в 03:00”, “по будням”, таймзона), подберу компактную реализацию именно под ваш шаблон.
Retry с экспоненциальной задержкой в asyncio
- Идея: оборачиваешь сетевую операцию в функцию, делаешь несколько попыток, между попытками
- Таймауты ставятся на двух уровнях:
- На уровне библиотеки/сокета (предпочтительно): connect/read/total таймауты (например, в
- На уровне корутины:
Пример с aiohttp: где именно ставить таймауты
- 1) Таймауты aiohttp: лучший способ для сетевых операций (connect/read/total).
- 2) Retry: вокруг “логической” операции запроса; повторять только на временных ошибках (timeout, DNS, reset, 5xx).
- 3) asyncio.timeout: обычно не нужен, если таймауты aiohttp настроены, но полезен как общий “guard” на весь шаг.
Практические правила
- Ставь таймауты как можно ближе к I/O (в клиенте/сокете): connect + read + total.
- Разделяй таймаут на попытку и общий дедлайн: “на попытку” ограничивает зависание, “общий” ограничивает весь процесс (включая задержки backoff).
- Не ретраь детерминированные ошибки: 400/401/403/404, ошибки валидации, некорректный URL.
- Ретраь временные: timeouts, connection reset, DNS временные, 429 (с учётом
- Добавляй jitter, чтобы избежать синхронных повторов под нагрузкой.
- Идея: оборачиваешь сетевую операцию в функцию, делаешь несколько попыток, между попытками
await asyncio.sleep() с экспоненциальным backoff и небольшим jitter, чтобы снизить эффект “стадного” повтора.- Таймауты ставятся на двух уровнях:
- На уровне библиотеки/сокета (предпочтительно): connect/read/total таймауты (например, в
aiohttp) — это корректно прерывает операции ввода-вывода.- На уровне корутины:
asyncio.timeout()/asyncio.wait_for() как “страховка” вокруг конкретного await, если библиотека не даёт нужных таймаутов или нужно ограничить общий сценарий.import asyncio
import random
from typing import Callable, TypeVar, Awaitable
T = TypeVar("T")
async def retry_async(
op: Callable[[], Awaitable[T]],
*,
attempts: int = 5,
base_delay: float = 0.2,
max_delay: float = 5.0,
timeout: float | None = None,
retry_on: tuple[type[BaseException], ...] = (TimeoutError, OSError),
) -> T:
last_exc: BaseException | None = None
for i in range(1, attempts + 1):
try:
if timeout is None:
return await op()
async with asyncio.timeout(timeout):
return await op()
except retry_on as e:
last_exc = e
if i == attempts:
raise
# exp backoff: base * 2^(i-1), capped + jitter
delay = min(max_delay, base_delay * (2 ** (i - 1)))
jitter = random.uniform(0, delay * 0.2)
await asyncio.sleep(delay + jitter)
assert last_exc is not None
raise last_excПример с aiohttp: где именно ставить таймауты
- 1) Таймауты aiohttp: лучший способ для сетевых операций (connect/read/total).
- 2) Retry: вокруг “логической” операции запроса; повторять только на временных ошибках (timeout, DNS, reset, 5xx).
- 3) asyncio.timeout: обычно не нужен, если таймауты aiohttp настроены, но полезен как общий “guard” на весь шаг.
import asyncio
import aiohttp
async def fetch_json(session: aiohttp.ClientSession, url: str) -> dict:
async with session.get(url) as resp:
# 4xx обычно не ретраят, 5xx часто ретраят
if 500 <= resp.status < 600:
raise aiohttp.ClientResponseError(
request_info=resp.request_info,
history=resp.history,
status=resp.status,
message="server error",
headers=resp.headers,
)
resp.raise_for_status()
return await resp.json()
async def main(url: str) -> dict:
timeout = aiohttp.ClientTimeout(
total=10, # общий потолок на запрос
connect=3, # установка соединения
sock_read=5, # чтение данных
)
async with aiohttp.ClientSession(timeout=timeout) as session:
async def op():
return await fetch_json(session, url)
return await retry_async(
op,
attempts=5,
base_delay=0.3,
max_delay=4.0,
# timeout на попытку (опционально, часто достаточно aiohttp таймаутов)
timeout=None,
retry_on=(aiohttp.ClientError, asyncio.TimeoutError, OSError),
)Практические правила
- Ставь таймауты как можно ближе к I/O (в клиенте/сокете): connect + read + total.
- Разделяй таймаут на попытку и общий дедлайн: “на попытку” ограничивает зависание, “общий” ограничивает весь процесс (включая задержки backoff).
- Не ретраь детерминированные ошибки: 400/401/403/404, ошибки валидации, некорректный URL.
- Ретраь временные: timeouts, connection reset, DNS временные, 429 (с учётом
Retry-After), 5xx.- Добавляй jitter, чтобы избежать синхронных повторов под нагрузкой.
Практика asyncio для сетевого I/O: идея — не создавать потоки на каждый запрос, а запускать много конкурентных операций ввода-вывода в одном потоке через event loop. Типовой стек: asyncio + http-клиент (aiohttp или httpx) + семафор/лимит соединений + таймауты + ретраи.
Типовой пример (httpx AsyncClient): конкурентная загрузка URL с лимитом, таймаутами и базовой обработкой ошибок
Почему так
- Один AsyncClient на все запросы: переиспользование соединений (keep-alive), меньше TLS-рукопожатий, меньше накладных расходов.
- Limits + Semaphore: пул соединений ограничивает количество одновременных TCP; семафор дополнительно защищает внешний сервис и вашу память/CPU.
- Timeout: без таймаутов задачи могут зависать на неопределённое время.
- Чтение ответа (
Сложность
- Время:
- Память:
Частые ошибки в asyncio-сетевом I/O
- Создавать новый клиент на каждый запрос (например, новый
- Забывать закрывать клиент: не использовать
- Отсутствие таймаутов: зависшие корутины, накопление задач, деградация сервиса.
- Слишком высокая конкурентность:
- Неправильная обработка исключений: скрывать ошибки через
- Блокирующие вызовы внутри корутин:
- Непонимание backpressure: читать/парсить огромные ответы целиком без ограничений. Для больших данных используйте потоковое чтение и лимиты размера.
- Смешивание sync и async клиентов: вызывать синхронный
Если нужен пример на aiohttp: стек аналогичен —
Типовой пример (httpx AsyncClient): конкурентная загрузка URL с лимитом, таймаутами и базовой обработкой ошибок
import asyncio
from typing import Iterable, List, Tuple
import httpx
async def fetch(
client: httpx.AsyncClient,
sem: asyncio.Semaphore,
url: str,
) -> Tuple[str, int, int]:
async with sem:
r = await client.get(url)
r.raise_for_status()
# Важно: читаем контент целиком, чтобы соединение корректно вернулось в пул.
content = r.content
return url, r.status_code, len(content)
async def main(urls: Iterable[str]) -> List[Tuple[str, int, int]]:
limits = httpx.Limits(
max_connections=100,
max_keepalive_connections=20,
keepalive_expiry=30.0,
)
timeout = httpx.Timeout(connect=5.0, read=20.0, write=10.0, pool=5.0)
sem = asyncio.Semaphore(50) # дополнительный лимит конкурентности поверх пула
async with httpx.AsyncClient(limits=limits, timeout=timeout, follow_redirects=True) as client:
tasks = [asyncio.create_task(fetch(client, sem, url)) for url in urls]
# return_exceptions=False: пусть ошибки всплывают (обычно лучше для наблюдаемости)
return await asyncio.gather(*tasks)
if __name__ == "__main__":
sample_urls = [
"https://example.com",
"https://httpbin.org/get",
"https://httpbin.org/uuid",
]
results = asyncio.run(main(sample_urls))
for url, status, size in results:
print(url, status, size)
Почему так
- Один AsyncClient на все запросы: переиспользование соединений (keep-alive), меньше TLS-рукопожатий, меньше накладных расходов.
- Limits + Semaphore: пул соединений ограничивает количество одновременных TCP; семафор дополнительно защищает внешний сервис и вашу память/CPU.
- Timeout: без таймаутов задачи могут зависать на неопределённое время.
- Чтение ответа (
r.content): гарантирует, что соединение корректно возвращается в пул (в большинстве случаев httpx сам управляет, но явное чтение безопасно, особенно если вы дальше не используете ответ).Сложность
- Время:
O(N) по числу URL, но реальная стенка зависит от сети; конкурентность ограничена min(sem, max_connections).- Память:
O(N) на результаты и задачи; плюс объём загруженных данных (в примере читаем ответы целиком). Если ответы большие — лучше стримить и не держать всё в памяти.Частые ошибки в asyncio-сетевом I/O
- Создавать новый клиент на каждый запрос (например, новый
AsyncClient/ClientSession): теряется пул соединений, падает производительность, растёт нагрузка.- Забывать закрывать клиент: не использовать
async with → утечки соединений/варнинги.- Отсутствие таймаутов: зависшие корутины, накопление задач, деградация сервиса.
- Слишком высокая конкурентность:
gather на десятки тысяч URL без лимитов → истощение файловых дескрипторов, memory pressure, бан от внешнего API.- Неправильная обработка исключений: скрывать ошибки через
return_exceptions=True и не логировать; или наоборот — не перехватывать ожидаемые сетевые ошибки (DNS, timeouts) там, где нужны ретраи.- Блокирующие вызовы внутри корутин:
time.sleep, тяжёлый CPU, синхронный DNS/IO → блокирует event loop. Нужны await asyncio.sleep(...), вынос CPU в asyncio.to_thread или процессы.- Непонимание backpressure: читать/парсить огромные ответы целиком без ограничений. Для больших данных используйте потоковое чтение и лимиты размера.
- Смешивание sync и async клиентов: вызывать синхронный
requests внутри async-кода — типичный антипаттерн.Если нужен пример на aiohttp: стек аналогичен —
asyncio + один aiohttp.ClientSession + TCPConnector(limit=...) + ClientTimeout + семафор.
Connection pooling в асинхронных HTTP-клиентах — это механизм повторного использования TCP/TLS соединений между запросами, чтобы не тратить время и ресурсы на постоянные DNS/TCP handshake/TLS handshake и не перегружать систему большим числом короткоживущих соединений.
Как это устроено
- Пул хранит набор открытых соединений, обычно раздельно по ключу вида
- При запросе клиент:
- ищет idle-соединение в пуле под нужный хост;
- если нашёл — берёт его и отправляет запрос;
- если нет — открывает новое (если не превышены лимиты);
- если лимиты превышены — ждёт освобождения (или падает по таймауту).
- После ответа соединение:
- возвращается в пул (если его можно переиспользовать);
- или закрывается (например,
- Keep-Alive: для HTTP/1.1 соединения по умолчанию живут и переиспользуются; пул следит за таймаутами простоя и “протухшими” сокетами.
- HTTP/2: одно соединение может обслуживать много параллельных запросов (multiplexing). Пул чаще лимитирует кол-во соединений, а не “один запрос = одно соединение”.
- Ограничения:
- лимит на все соединения;
- лимит на один хост (per-host);
- таймаут ожидания свободного соединения;
- максимальный размер пула/очереди.
Почему важно переиспользовать сессию/клиент
- Сессия (или клиент) обычно владеет:
- пулом соединений;
- DNS-кэшем/резолвером;
- cookies (если включены);
- настройками TLS/прокси/таймаутов.
Если создавать “клиент на каждый запрос”, вы ломаете pooling: получите лишние handshakes, больше нагрузку на CPU/сеть, выше латентность и риск упереться в лимиты файловых дескрипторов.
Правильные практики
- Делайте один клиент/сессию на жизненный цикл приложения (или на модуль/сервис), и переиспользуйте его для всех запросов.
- Закрывайте клиент при завершении приложения, чтобы корректно закрыть сокеты.
- Настраивайте лимиты под параллелизм: общий и per-host, иначе либо будет слишком много соединений, либо лишние ожидания.
- Ставьте таймауты: connect/read/write/pool — чтобы не зависать при проблемах сети/сервера.
- Не делитесь одним клиентом между разными event loop (если вы создаёте несколько циклов/потоков) — клиент привязан к своему loop.
- Если есть разные прокси/SSL-политики/разные наборы заголовков “по умолчанию” — иногда логично иметь несколько клиентов, но каждый переиспользовать.
Примеры (Python)
aiohttp (одна
httpx (один
Типичные ошибки
- Создавать
- Не закрывать клиент — утечки сокетов/предупреждения, подвисание при завершении.
- Слишком маленький
- Слишком большой лимит — много одновременных коннектов, давление на ОС и удалённый сервис.
Рекомендация “в одном предложении”: создайте один асинхронный HTTP-клиент на приложение, настройте лимиты/таймауты и переиспользуйте его для всех запросов, закрывая при shutdown.
Как это устроено
- Пул хранит набор открытых соединений, обычно раздельно по ключу вида
(scheme, host, port) + параметры (прокси, SNI/SSL-контекст и т.п.).- При запросе клиент:
- ищет idle-соединение в пуле под нужный хост;
- если нашёл — берёт его и отправляет запрос;
- если нет — открывает новое (если не превышены лимиты);
- если лимиты превышены — ждёт освобождения (или падает по таймауту).
- После ответа соединение:
- возвращается в пул (если его можно переиспользовать);
- или закрывается (например,
Connection: close, ошибка протокола, истёк keep-alive, сервер закрыл сокет, HTTP/1.0 без keep-alive).- Keep-Alive: для HTTP/1.1 соединения по умолчанию живут и переиспользуются; пул следит за таймаутами простоя и “протухшими” сокетами.
- HTTP/2: одно соединение может обслуживать много параллельных запросов (multiplexing). Пул чаще лимитирует кол-во соединений, а не “один запрос = одно соединение”.
- Ограничения:
- лимит на все соединения;
- лимит на один хост (per-host);
- таймаут ожидания свободного соединения;
- максимальный размер пула/очереди.
Почему важно переиспользовать сессию/клиент
- Сессия (или клиент) обычно владеет:
- пулом соединений;
- DNS-кэшем/резолвером;
- cookies (если включены);
- настройками TLS/прокси/таймаутов.
Если создавать “клиент на каждый запрос”, вы ломаете pooling: получите лишние handshakes, больше нагрузку на CPU/сеть, выше латентность и риск упереться в лимиты файловых дескрипторов.
Правильные практики
- Делайте один клиент/сессию на жизненный цикл приложения (или на модуль/сервис), и переиспользуйте его для всех запросов.
- Закрывайте клиент при завершении приложения, чтобы корректно закрыть сокеты.
- Настраивайте лимиты под параллелизм: общий и per-host, иначе либо будет слишком много соединений, либо лишние ожидания.
- Ставьте таймауты: connect/read/write/pool — чтобы не зависать при проблемах сети/сервера.
- Не делитесь одним клиентом между разными event loop (если вы создаёте несколько циклов/потоков) — клиент привязан к своему loop.
- Если есть разные прокси/SSL-политики/разные наборы заголовков “по умолчанию” — иногда логично иметь несколько клиентов, но каждый переиспользовать.
Примеры (Python)
aiohttp (одна
ClientSession + TCPConnector):import aiohttp
import asyncio
async def main():
timeout = aiohttp.ClientTimeout(total=30, connect=5, sock_read=20)
connector = aiohttp.TCPConnector(
limit=200, # общий лимит соединений
limit_per_host=50, # на хост
ttl_dns_cache=300, # DNS cache
enable_cleanup_closed=True,
)
async with aiohttp.ClientSession(connector=connector, timeout=timeout) as session:
async with session.get("https://example.com") as r:
data = await r.text()
print(r.status, len(data))
asyncio.run(main())httpx (один
AsyncClient, pooling внутри):import asyncio
import httpx
async def main():
limits = httpx.Limits(max_connections=200, max_keepalive_connections=50)
timeout = httpx.Timeout(30.0, connect=5.0)
async with httpx.AsyncClient(limits=limits, timeout=timeout, http2=True) as client:
r = await client.get("https://example.com")
print(r.status_code, len(r.text))
asyncio.run(main())Типичные ошибки
- Создавать
ClientSession()/AsyncClient() внутри функции “на один запрос” и сразу закрывать — пул не успевает принести пользу.- Не закрывать клиент — утечки сокетов/предупреждения, подвисание при завершении.
- Слишком маленький
limit_per_host при высоком параллелизме — рост очередей и латентности.- Слишком большой лимит — много одновременных коннектов, давление на ОС и удалённый сервис.
Рекомендация “в одном предложении”: создайте один асинхронный HTTP-клиент на приложение, настройте лимиты/таймауты и переиспользуйте его для всех запросов, закрывая при shutdown.
Rate limiting — это ограничение частоты: сколько запросов можно сделать за единицу времени (например,
Чем отличаются на практике
- Rate limiting защищает внешний сервис/вас от превышения квот, антибота, 429, и выравнивает поток во времени.
- Concurrency limiting защищает ваши ресурсы (CPU, сокеты, файловые дескрипторы) и не даёт раздувать число одновременно «висящих» I/O задач.
- Часто нужно оба: сначала дождаться «токена» по скорости, потом войти под семафор конкурентности.
Простейший rate limiter в asyncio (token bucket)
Ограничение конкурентности (semaphore)
Комбинация: и скорость, и конкурентность
Ключевые нюансы
- Rate limiting должен быть общим для всех задач, которые делят лимит (один объект bucket на общий поток).
- Для нескольких разных лимитов (например, по доменам/API-ключам) делайте отдельные limiter’ы.
- Semaphore не гарантирует «N в секунду»: если запросы быстрые, семафор пропустит много запросов за секунду; он ограничивает только число одновременно выполняющихся.
10 req/s). Ограничение конкурентности — это ограничение параллелизма: сколько задач одновременно выполняются (например, max 20 in-flight). Это разные оси контроля: можно делать мало параллельных запросов, но слишком часто (нарушая лимит), и наоборот — много параллельных, но редко (не нарушая лимит).Чем отличаются на практике
- Rate limiting защищает внешний сервис/вас от превышения квот, антибота, 429, и выравнивает поток во времени.
- Concurrency limiting защищает ваши ресурсы (CPU, сокеты, файловые дескрипторы) и не даёт раздувать число одновременно «висящих» I/O задач.
- Часто нужно оба: сначала дождаться «токена» по скорости, потом войти под семафор конкурентности.
Простейший rate limiter в asyncio (token bucket)
import asyncio
import time
class AsyncTokenBucket:
def __init__(self, rate: float, capacity: int):
"""
rate: токенов в секунду (напр. 5.0 = 5 req/s)
capacity: максимальный "запас" токенов (burst)
"""
self.rate = float(rate)
self.capacity = int(capacity)
self._tokens = float(capacity)
self._updated = time.monotonic()
self._lock = asyncio.Lock()
def _refill(self) -> None:
now = time.monotonic()
delta = now - self._updated
self._updated = now
self._tokens = min(self.capacity, self._tokens + delta * self.rate)
async def acquire(self, n: float = 1.0) -> None:
while True:
async with self._lock:
self._refill()
if self._tokens >= n:
self._tokens -= n
return
need = n - self._tokens
wait = need / self.rate if self.rate > 0 else 3600.0
await asyncio.sleep(wait)Ограничение конкурентности (semaphore)
import asyncio
sem = asyncio.Semaphore(20)
async def limited_call(fn, *args, **kwargs):
async with sem:
return await fn(*args, **kwargs)Комбинация: и скорость, и конкурентность
import asyncio
bucket = AsyncTokenBucket(rate=10, capacity=20) # 10 req/s, burst до 20
sem = asyncio.Semaphore(50) # не больше 50 одновременно
async def fetch(session, url):
await bucket.acquire(1) # ограничили частоту
async with sem: # ограничили параллелизм
async with session.get(url) as r:
return await r.text()Ключевые нюансы
- Rate limiting должен быть общим для всех задач, которые делят лимит (один объект bucket на общий поток).
- Для нескольких разных лимитов (например, по доменам/API-ключам) делайте отдельные limiter’ы.
- Semaphore не гарантирует «N в секунду»: если запросы быстрые, семафор пропустит много запросов за секунду; он ограничивает только число одновременно выполняющихся.
Цель: корректно поймать
Ключевые принципы:
-
- На Unix используйте
- Завершение делайте в одном месте: выставить флаг остановки, отменить задачи, дождаться
Практические советы:
- Внутри задач регулярно делайте
- Для долгих операций используйте
- В
SIGINT/SIGTERM, инициировать graceful shutdown, отменить фоновые задачи, дождаться их завершения и закрыть ресурсы.Ключевые принципы:
-
asyncio сам по себе не “останавливает” задачи: нужно явно инициировать завершение (через Event и/или отмену задач).- На Unix используйте
loop.add_signal_handler(...); на Windows для SIGTERM поддержка ограничена, но KeyboardInterrupt (Ctrl+C) обычно работает.- Завершение делайте в одном месте: выставить флаг остановки, отменить задачи, дождаться
gather, закрыть соединения.import asyncio
import signal
from contextlib import suppress
async def worker(name: str, stop: asyncio.Event) -> None:
try:
while not stop.is_set():
# делаем полезную работу
await asyncio.sleep(0.5)
except asyncio.CancelledError:
# быстрые действия перед отменой (при необходимости)
raise
finally:
# гарантированное освобождение ресурсов (файлы, соединения и т.п.)
pass
async def shutdown(tasks: set[asyncio.Task], stop: asyncio.Event) -> None:
stop.set() # просим задачи завершиться сами
for t in tasks:
t.cancel() # на случай "зависших" await/IO
# ждём завершения, подавляя CancelledError
with suppress(asyncio.CancelledError):
await asyncio.gather(*tasks, return_exceptions=False)
async def main() -> None:
loop = asyncio.get_running_loop()
stop = asyncio.Event()
tasks: set[asyncio.Task] = set()
tasks.add(asyncio.create_task(worker("A", stop)))
tasks.add(asyncio.create_task(worker("B", stop)))
def _on_signal(sig: int) -> None:
# важно: не блокировать тут, только поставить "флаг"
stop.set()
# Unix: корректная обработка сигналов event loop'ом
for s in (signal.SIGINT, signal.SIGTERM):
try:
loop.add_signal_handler(s, _on_signal, s)
except NotImplementedError:
# например, Windows/некоторые окружения
pass
try:
# ждём сигнала/остановки
await stop.wait()
finally:
await shutdown(tasks, stop)
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
# если сигналы не доставлены в loop (часто на Windows), всё равно выходим
pass
Практические советы:
- Внутри задач регулярно делайте
await (или таймауты), чтобы отмена могла “пробиться”.- Для долгих операций используйте
asyncio.wait_for(..., timeout) или проверяйте stop.is_set().- В
finally закрывайте ресурсы; для сетевых клиентов (aiohttp/asyncpg и т.п.) делайте явное await client.close()/await pool.close() перед завершением.
Asyncio Streams API (
Как устроено
- Клиент:
- Сервер:
- Reader читает буферизованные байты из сокета, предоставляя методы вроде
- Writer пишет в сокет через
- Закрытие:
Ключевые детали поведения
- Границы сообщений отсутствуют: TCP — поток байт. Вам нужно самостоятельно определить фрейминг: длина+тело, разделитель, фиксированные блоки и т.п.
- Конкурентность: один event loop обслуживает много соединений; на каждое соединение обычно создают задачу-обработчик, внутри которой читают/пишут.
- Таймауты делаются снаружи через
- Flow control на чтение/запись: чтение может копить буфер до лимита (настраивается в деталях), запись регулируется
- SSL поддерживается параметром
Мини-примеры
Когда выбирают Streams API вместо высокоуровневых библиотек
- Свой протокол поверх TCP: бинарный/текстовый протокол, где нужно самому управлять фреймингом, состоянием соединения, а высокоуровневый HTTP/WebSocket/gRPC избыточен.
- Минимальные зависимости и контроль: требуется простой, лёгкий, полностью управляемый IO-слой без «фреймворка».
- Прокси/туннели/шлюзы: прозрачная перекачка байт, backpressure, управление буферами, half-close.
- Обучение/отладка: проще, чем transport/protocol, но ближе к «голому» сокету, чем aiohttp и аналоги.
- Нестандартные требования: кастомный handshake, multiplexing, специфичная схема таймаутов/ретраев, особая политика закрытия.
Когда лучше не Streams
- HTTP(S)/REST: обычно выгоднее aiohttp/httpx (клиент) и aiohttp/fastapi+uvicorn (сервер), чтобы не реализовывать парсинг, keep-alive, chunked, заголовки, компрессию и т.д.
- WebSocket: проще взять готовую библиотеку (корректный handshake, ping/pong, фреймы, маскирование).
- gRPC/AMQP/MQTT: протоколы сложные; готовые реализации экономят время и снижают риск ошибок.
- Нужны богатые middleware/роутинг/валидация/авторизация: это уровень веб-фреймворков, а не streams.
Практическое правило выбора
- Если задача — «обмен байтами по TCP по своему простому протоколу» и важны контроль и лёгкость — берите
- Если задача — стандартный прикладной протокол (HTTP/WebSocket/gRPC и т.п.) — выбирайте специализированную библиотеку.
asyncio.start_server, asyncio.open_connection) — это низко-/среднеуровневая обёртка над transport/protocol, дающая удобную работу с TCP-потоками через пары StreamReader/StreamWriter (клиент) и callback на подключение (сервер).Как устроено
- Клиент:
await asyncio.open_connection(host, port) создаёт TCP-соединение и возвращает (reader, writer).- Сервер:
await asyncio.start_server(client_connected_cb, host, port) слушает порт; на каждое подключение вызывает client_connected_cb(reader, writer) в отдельной задаче (логически параллельно, в рамках event loop).- Reader читает буферизованные байты из сокета, предоставляя методы вроде
read(n), readexactly(n), readline(), readuntil(sep).- Writer пишет в сокет через
write(data), а await drain() реализует backpressure: при переполнении буфера отправки корутина приостанавливается, пока ОС/транспорт не освободит место.- Закрытие:
writer.close() + await writer.wait_closed() (корректное завершение); для half-close есть write_eof() (не везде поддерживается).Ключевые детали поведения
- Границы сообщений отсутствуют: TCP — поток байт. Вам нужно самостоятельно определить фрейминг: длина+тело, разделитель, фиксированные блоки и т.п.
- Конкурентность: один event loop обслуживает много соединений; на каждое соединение обычно создают задачу-обработчик, внутри которой читают/пишут.
- Таймауты делаются снаружи через
asyncio.wait_for() или asyncio.timeout().- Flow control на чтение/запись: чтение может копить буфер до лимита (настраивается в деталях), запись регулируется
drain().- SSL поддерживается параметром
ssl=... (и на клиенте, и на сервере).Мини-примеры
import asyncio
async def handle(reader: asyncio.StreamReader, writer: asyncio.StreamWriter):
try:
while True:
line = await reader.readline()
if not line:
break
writer.write(line)
await writer.drain()
finally:
writer.close()
await writer.wait_closed()
async def main():
server = await asyncio.start_server(handle, "127.0.0.1", 8888)
async with server:
await server.serve_forever()
asyncio.run(main())import asyncio
async def main():
reader, writer = await asyncio.open_connection("127.0.0.1", 8888)
writer.write(b"hello\n")
await writer.drain()
resp = await reader.readline()
writer.close()
await writer.wait_closed()
print(resp)
asyncio.run(main())Когда выбирают Streams API вместо высокоуровневых библиотек
- Свой протокол поверх TCP: бинарный/текстовый протокол, где нужно самому управлять фреймингом, состоянием соединения, а высокоуровневый HTTP/WebSocket/gRPC избыточен.
- Минимальные зависимости и контроль: требуется простой, лёгкий, полностью управляемый IO-слой без «фреймворка».
- Прокси/туннели/шлюзы: прозрачная перекачка байт, backpressure, управление буферами, half-close.
- Обучение/отладка: проще, чем transport/protocol, но ближе к «голому» сокету, чем aiohttp и аналоги.
- Нестандартные требования: кастомный handshake, multiplexing, специфичная схема таймаутов/ретраев, особая политика закрытия.
Когда лучше не Streams
- HTTP(S)/REST: обычно выгоднее aiohttp/httpx (клиент) и aiohttp/fastapi+uvicorn (сервер), чтобы не реализовывать парсинг, keep-alive, chunked, заголовки, компрессию и т.д.
- WebSocket: проще взять готовую библиотеку (корректный handshake, ping/pong, фреймы, маскирование).
- gRPC/AMQP/MQTT: протоколы сложные; готовые реализации экономят время и снижают риск ошибок.
- Нужны богатые middleware/роутинг/валидация/авторизация: это уровень веб-фреймворков, а не streams.
Практическое правило выбора
- Если задача — «обмен байтами по TCP по своему простому протоколу» и важны контроль и лёгкость — берите
start_server/open_connection.- Если задача — стандартный прикладной протокол (HTTP/WebSocket/gRPC и т.п.) — выбирайте специализированную библиотеку.
Ключевая проблема: в
Практические правила
- Делай изменения состояния атомарными: собирай результат в локальные переменные, а общий state обновляй в одном месте и под блокировкой.
- Защищай критические секции через
- Всегда освобождай ресурсы: используй
- Если внутри критической секции есть await, отмена может прервать её. Либо не делай await внутри lock, либо на короткие критические участки применяй
- Старайся избегать общего mutable state: модель «один владелец состояния» (actor) через
1) Атомарный commit под Lock (рекомендуется)
2) Двухфазное изменение + откат (если без await внутри lock нельзя)
Идея: фиксируем «черновик»/маркер, затем выполняем работу, затем завершаем; при отмене откатываем.
3) Короткая неотменяемая критическая секция через shield (осторожно)
Используй только для очень коротких операций, иначе ухудшишь отзывчивость отмены.
4) Самый надёжный подход: «владелец состояния» (actor) через Queue
Только одна корутина мутирует state; остальные шлют ей команды. Тогда отмена клиентов не ломает state.
Итого
- Лучший базовый паттерн: вычисления и I/O вне lock, затем один быстрый commit под lock без await.
- Если нужны await в процессе изменения state: двухфазность + откат.
- Для сложного общего состояния: actor через Queue (минимум гонок и проблем с отменой).
asyncio отмена (CancelledError) может прилететь на любом await, поэтому общий state нельзя менять «пошагово» без защиты: задача может завершиться между частичными изменениями и оставить состояние неконсистентным.Практические правила
- Делай изменения состояния атомарными: собирай результат в локальные переменные, а общий state обновляй в одном месте и под блокировкой.
- Защищай критические секции через
asyncio.Lock/Semaphore. Не держи lock вокруг долгих I/O; держи только на время быстрого изменения state.- Всегда освобождай ресурсы: используй
async with lock и try/finally.- Если внутри критической секции есть await, отмена может прервать её. Либо не делай await внутри lock, либо на короткие критические участки применяй
asyncio.shield / подавляй отмену с обязательным восстановлением флага отмены.- Старайся избегать общего mutable state: модель «один владелец состояния» (actor) через
asyncio.Queue обычно проще и безопаснее.1) Атомарный commit под Lock (рекомендуется)
import asyncio
state = {"balance": 0}
lock = asyncio.Lock()
async def add_money(amount: int):
# Долгие/отменяемые операции — вне lock
await asyncio.sleep(0) # имитация I/O, тут может прилететь cancel
# Быстрый атомарный commit — под lock, без await внутри
async with lock:
state["balance"] += amount
2) Двухфазное изменение + откат (если без await внутри lock нельзя)
Идея: фиксируем «черновик»/маркер, затем выполняем работу, затем завершаем; при отмене откатываем.
import asyncio
state = {"items": set(), "inflight": set()}
lock = asyncio.Lock()
async def add_item(item: str):
async with lock:
state["inflight"].add(item)
try:
# Любые await-операции могут быть отменены
await asyncio.sleep(1)
async with lock:
state["items"].add(item)
except asyncio.CancelledError:
# Откат/очистка маркера
async with lock:
state["inflight"].discard(item)
raise
else:
async with lock:
state["inflight"].discard(item)
3) Короткая неотменяемая критическая секция через shield (осторожно)
Используй только для очень коротких операций, иначе ухудшишь отзывчивость отмены.
import asyncio
lock = asyncio.Lock()
state = {"n": 0}
async def critical_commit(delta: int):
async def _commit():
async with lock:
state["n"] += delta
try:
await asyncio.shield(_commit())
except asyncio.CancelledError:
# Важно: shield предотвращает отмену _commit, но отмена внешней задачи всё равно должна проявиться
raise
4) Самый надёжный подход: «владелец состояния» (actor) через Queue
Только одна корутина мутирует state; остальные шлют ей команды. Тогда отмена клиентов не ломает state.
import asyncio
class StateActor:
def __init__(self):
self._q = asyncio.Queue()
self._state = {"balance": 0}
async def run(self):
while True:
cmd, arg, fut = await self._q.get()
try:
if cmd == "add":
self._state["balance"] += arg
fut.set_result(self._state["balance"])
else:
raise ValueError("unknown cmd")
except Exception as e:
fut.set_exception(e)
async def add(self, amount: int) -> int:
fut = asyncio.get_running_loop().create_future()
await self._q.put(("add", amount, fut))
return await fut
Итого
- Лучший базовый паттерн: вычисления и I/O вне lock, затем один быстрый commit под lock без await.
- Если нужны await в процессе изменения state: двухфазность + откат.
- Для сложного общего состояния: actor через Queue (минимум гонок и проблем с отменой).
Цель: понятная слоистая архитектура, управляемый lifecycle, безопасные фоновые задачи.
1) Разделение на слои
- Domain: чистая бизнес-логика (entities, value objects, правила, исключения). Без asyncio/БД/HTTP.
- Application: use-cases (оркестрация домена, транзакции, политики ретраев/таймаутов), интерфейсы (ports) к внешним зависимостям: репозитории, брокеры, кэш, http-клиенты.
- Infrastructure: реализации портов: SQL/Redis/Kafka/HTTP, адаптеры, сериализация, миграции.
- Presentation/Transport: HTTP/gRPC/CLI/consumer, валидация входа, маппинг DTO ↔ domain, обработка ошибок, авторизация.
- Cross-cutting: конфиг, логирование, метрики, трассировка, DI, healthchecks.
Ключевое правило: зависимости направлены внутрь (transport → app → domain; infra подключается через интерфейсы).
2) Управление жизненным циклом (startup/shutdown)
- Composition root (обычно
- Все долгоживущие ресурсы создавай на startup и закрывай на shutdown через
- Везде прокидывай явные зависимости (параметрами/DI), не используй глобальные синглтоны.
- Для отмены используй
3) Фоновые задачи (background jobs)
- Запускай их из lifecycle-менеджера, а не «где-то в обработчике».
- Держи TaskGroup (Python 3.11+) или свою группу задач: при shutdown отменяй все задачи и жди завершения.
- Каждая задача должна:
- периодически await-ить (не блокировать loop),
- корректно реагировать на cancel,
- иметь таймауты/ретраи для I/O,
- не терять исключения (логировать/прокидывать).
- Для очередей/воркеров: предпочитай модель consumer loop + пул воркеров (Semaphore) + graceful shutdown.
4) Практические договоренности
- Границы: transport слой не содержит бизнес-логики, только маппинг/валидацию/вызов use-case.
- Таймауты: любой внешний I/O оборачивай в
- Изоляция блокировок: CPU/блокирующее — в
- Ошибки: единый формат доменных исключений; в transport слой — маппинг в HTTP/gRPC коды.
- Тесты: domain и application тестируются без инфраструктуры; infra — интеграционные.
Если скажешь какой транспорт (FastAPI/aiohttp/gRPC), какая БД и есть ли очереди, дам конкретный шаблон папок, DI и примеры lifespan под твой стек.
1) Разделение на слои
- Domain: чистая бизнес-логика (entities, value objects, правила, исключения). Без asyncio/БД/HTTP.
- Application: use-cases (оркестрация домена, транзакции, политики ретраев/таймаутов), интерфейсы (ports) к внешним зависимостям: репозитории, брокеры, кэш, http-клиенты.
- Infrastructure: реализации портов: SQL/Redis/Kafka/HTTP, адаптеры, сериализация, миграции.
- Presentation/Transport: HTTP/gRPC/CLI/consumer, валидация входа, маппинг DTO ↔ domain, обработка ошибок, авторизация.
- Cross-cutting: конфиг, логирование, метрики, трассировка, DI, healthchecks.
Ключевое правило: зависимости направлены внутрь (transport → app → domain; infra подключается через интерфейсы).
2) Управление жизненным циклом (startup/shutdown)
- Composition root (обычно
main.py): собирает контейнер зависимостей, создает клиентов (DB pool, redis, http sessions), запускает сервер, стартует фоновые задачи.- Все долгоживущие ресурсы создавай на startup и закрывай на shutdown через
async with/aclose()/close().- Везде прокидывай явные зависимости (параметрами/DI), не используй глобальные синглтоны.
- Для отмены используй
asyncio.CancelledError: ловить, выполнять cleanup, затем обязательно пробрасывать дальше.3) Фоновые задачи (background jobs)
- Запускай их из lifecycle-менеджера, а не «где-то в обработчике».
- Держи TaskGroup (Python 3.11+) или свою группу задач: при shutdown отменяй все задачи и жди завершения.
- Каждая задача должна:
- периодически await-ить (не блокировать loop),
- корректно реагировать на cancel,
- иметь таймауты/ретраи для I/O,
- не терять исключения (логировать/прокидывать).
- Для очередей/воркеров: предпочитай модель consumer loop + пул воркеров (Semaphore) + graceful shutdown.
import asyncio
from contextlib import asynccontextmanager
class App:
def __init__(self, db, http):
self.db = db
self.http = http
async def close(self):
await self.http.aclose()
await self.db.aclose()
async def poller(app: App):
try:
while True:
# пример периодической работы
await asyncio.sleep(1)
except asyncio.CancelledError:
# cleanup при необходимости
raise
async def worker(app: App, queue: asyncio.Queue[int]):
try:
while True:
item = await queue.get()
try:
# обработка item (I/O только с таймаутами)
await asyncio.wait_for(asyncio.sleep(0.1), timeout=5)
finally:
queue.task_done()
except asyncio.CancelledError:
raise
@asynccontextmanager
async def lifespan():
# startup: создаем ресурсы
db = ... # async pool/client
http = ... # http client session
app = App(db=db, http=http)
queue: asyncio.Queue[int] = asyncio.Queue()
try:
async with asyncio.TaskGroup() as tg:
tg.create_task(poller(app))
tg.create_task(worker(app, queue))
# отдаём управление «серверу»
yield app
# shutdown начнется после выхода из yield
# TaskGroup отменит задачи и дождется завершения
finally:
await app.close()
async def main():
async with lifespan() as app:
# здесь запускается HTTP сервер / consumer / CLI loop
await asyncio.sleep(10)
if __name__ == "__main__":
asyncio.run(main())
4) Практические договоренности
- Границы: transport слой не содержит бизнес-логики, только маппинг/валидацию/вызов use-case.
- Таймауты: любой внешний I/O оборачивай в
asyncio.timeout() (3.11+) или wait_for.- Изоляция блокировок: CPU/блокирующее — в
asyncio.to_thread() или отдельный process.- Ошибки: единый формат доменных исключений; в transport слой — маппинг в HTTP/gRPC коды.
- Тесты: domain и application тестируются без инфраструктуры; infra — интеграционные.
Если скажешь какой транспорт (FastAPI/aiohttp/gRPC), какая БД и есть ли очереди, дам конкретный шаблон папок, DI и примеры lifespan под твой стек.
Структурированная конкурентность — это подход, где жизненный цикл конкурентных задач жёстко ограничен областью видимости (scope): задачи создаются внутри контекста и гарантированно завершаются/отменяются при выходе из него, а ошибки агрегируются предсказуемо.
Ключевые паттерны в современном Python (asyncio, Python 3.11+)
- Nursery/Scope через
- Fail-fast (отмена остальных при первой ошибке): если любая задача падает, остальные отменяются, а ошибка(и) возвращаются как
- Агрегация ошибок: несколько исключений из разных задач не теряются, а поднимаются вместе (PEP 654).
- Гарантированная очистка: при выходе из
- Ограниченный параллелизм (часто в связке): структурирование + лимиты через
Какие проблемы решает
- Утечки задач: при ручном
- Потерянные исключения: “Task exception was never retrieved” и скрытые падения в фоне;
- Неопределённый порядок отмены: при ошибке в одной задаче остальные продолжают работать; в
- Сложная ручная координация: вместо сборки
- Проблемы с корректным shutdown: при отмене верхнего уровня (например, сигнал/таймаут) группа корректно отменяет дочерние задачи и дожидается их завершения.
Базовый пример
Паттерн: ограниченный параллелизм внутри структурированного scope
Идея: задачи структурированы (не протекают наружу) и одновременно контролируется нагрузка.
Связанные элементы структурированного стиля
-
-
Итог:
Ключевые паттерны в современном Python (asyncio, Python 3.11+)
- Nursery/Scope через
asyncio.TaskGroup: групповой запуск задач с гарантированным завершением.- Fail-fast (отмена остальных при первой ошибке): если любая задача падает, остальные отменяются, а ошибка(и) возвращаются как
ExceptionGroup.- Агрегация ошибок: несколько исключений из разных задач не теряются, а поднимаются вместе (PEP 654).
- Гарантированная очистка: при выходе из
TaskGroup не остаётся “висячих” фоновых задач.- Ограниченный параллелизм (часто в связке): структурирование + лимиты через
asyncio.Semaphore или пул, чтобы не “взрывать” число задач.Какие проблемы решает
TaskGroup- Утечки задач: при ручном
create_task() легко забыть дождаться/отменить задачу; TaskGroup гарантирует завершение всех.- Потерянные исключения: “Task exception was never retrieved” и скрытые падения в фоне;
TaskGroup поднимает исключения при выходе из контекста.- Неопределённый порядок отмены: при ошибке в одной задаче остальные продолжают работать; в
TaskGroup типичный паттерн — fail-fast.- Сложная ручная координация: вместо сборки
gather + обработка отмен/ошибок/cleanup — единый структурированный блок.- Проблемы с корректным shutdown: при отмене верхнего уровня (например, сигнал/таймаут) группа корректно отменяет дочерние задачи и дожидается их завершения.
Базовый пример
TaskGroup (fail-fast, агрегирование исключений)import asyncio
async def work(n: int) -> int:
await asyncio.sleep(0.1)
if n == 2:
raise ValueError("boom")
return n * 10
async def main():
try:
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(work(i)) for i in range(5)]
# Если ошибок нет — здесь все задачи завершены
results = [t.result() for t in tasks]
print(results)
except* ValueError as eg:
# ExceptionGroup: можно фильтровать по типам
print("caught:", eg)
asyncio.run(main())
Паттерн: ограниченный параллелизм внутри структурированного scope
import asyncio
sem = asyncio.Semaphore(10)
async def bounded(coro):
async with sem:
return await coro
async def main(urls):
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(bounded(fetch(u))) for u in urls]
return [t.result() for t in tasks]
Идея: задачи структурированы (не протекают наружу) и одновременно контролируется нагрузка.
Связанные элементы структурированного стиля
-
asyncio.timeout() (3.11+): структурированный таймаут, который отменяет вложенные операции в рамках блока, упрощая корректную отмену.-
ExceptionGroup и except* (3.11+): нативная модель обработки множественных исключений из конкурентных задач.Итог:
asyncio.TaskGroup (и связка timeout + ExceptionGroup) дают в Python структурированную конкурентность: предсказуемое время жизни задач, корректную отмену и полную, управляемую обработку ошибок без “фонового хаоса”.
asyncio.TaskGroup (Python 3.11+) — это структурированная конкурентность: вы создаёте группу, добавляете задачи, и при выходе из блока гарантируется, что все задачи либо завершились успешно, либо были корректно отменены при ошибке.
Базовое использование
Чем отличается от asyncio.gather
- Контроль жизненного цикла:
- Поведение при исключениях (по умолчанию):
- Агрегация ошибок:
- API результата: у
Как TaskGroup ведёт себя при исключениях
- если любая задача в группе падает исключением, TaskGroup отменяет остальные задачи (им прилетит
- если исключений несколько, они группируются в один
- если внутри задач вы перехватили
Пример: одна задача падает, остальные отменяются, ловим ExceptionGroup
Сравнение с gather по исключениям
-
Когда что выбирать
- TaskGroup: когда важны строгие границы, отмена остальных при ошибке и корректное “закрытие” параллельных подзадач (рекомендуемый стиль в 3.11+).
- gather: когда нужно просто собрать результаты из набора awaitable’ов, особенно если вы явно хотите получить ошибки как значения через
Базовое использование
import asyncio
async def job(n: int) -> int:
await asyncio.sleep(0.1)
return n
async def main():
async with asyncio.TaskGroup() as tg:
t1 = tg.create_task(job(1))
t2 = tg.create_task(job(2))
# здесь гарантированно: обе задачи завершены
print(t1.result(), t2.result())
asyncio.run(main())Чем отличается от asyncio.gather
- Контроль жизненного цикла:
TaskGroup привязывает задачи к блоку async with и не даёт “утечь” фоновым задачам; gather просто ждёт переданные awaitable’ы и сам по себе не задаёт структурных границ.- Поведение при исключениях (по умолчанию):
TaskGroup при первой ошибке отменяет остальные задачи и затем выбрасывает агрегированное исключение; gather обычно выбрасывает первое исключение, а судьба остальных зависит от момента/отмены, и часто требуется ручная обработка.- Агрегация ошибок:
TaskGroup поднимает ExceptionGroup (если ошибок несколько); gather либо поднимает одно исключение, либо (с return_exceptions=True) возвращает исключения как значения.- API результата: у
gather результат — список/кортеж значений; у TaskGroup вы работаете с объектами Task и берёте .result()/.exception() после выхода из блока.Как TaskGroup ведёт себя при исключениях
- если любая задача в группе падает исключением, TaskGroup отменяет остальные задачи (им прилетит
asyncio.CancelledError), ждёт их завершения и затем выбрасывает ExceptionGroup.- если исключений несколько, они группируются в один
ExceptionGroup.- если внутри задач вы перехватили
CancelledError и подавили его, задача может “пережить” отмену — но это плохая практика: отмена должна корректно проходить вверх.Пример: одна задача падает, остальные отменяются, ловим ExceptionGroup
import asyncio
async def ok():
try:
await asyncio.sleep(1)
return "ok"
except asyncio.CancelledError:
# корректно пробрасываем отмену
raise
async def boom():
await asyncio.sleep(0.2)
raise ValueError("fail")
async def main():
try:
async with asyncio.TaskGroup() as tg:
t1 = tg.create_task(ok())
t2 = tg.create_task(boom())
t3 = tg.create_task(ok())
except* ValueError as eg:
# eg.exceptions содержит все ValueError из группы
print("caught:", [str(e) for e in eg.exceptions])
asyncio.run(main())Сравнение с gather по исключениям
-
await asyncio.gather(a, b): при исключении обычно выбросит его наружу; если нужны “все результаты/ошибки”, используют return_exceptions=True.import asyncio
async def main():
res = await asyncio.gather(
asyncio.sleep(0.1, result=1),
asyncio.sleep(0.2, result=2),
return_exceptions=True,
)
print(res)
asyncio.run(main())Когда что выбирать
- TaskGroup: когда важны строгие границы, отмена остальных при ошибке и корректное “закрытие” параллельных подзадач (рекомендуемый стиль в 3.11+).
- gather: когда нужно просто собрать результаты из набора awaitable’ов, особенно если вы явно хотите получить ошибки как значения через
return_exceptions=True.
Инструменты наблюдаемости и отладки для asyncio (что реально использую)
- Логирование (structlog / logging)
- Корреляция запросов:
- Логи задач: имя/идентификатор asyncio-task (
- Логи медленных операций: оборачиваю I/O и внешние вызовы в измерители длительности; пороги (например >100–300ms) логирую как warning.
- Для продакшена предпочитаю структурные JSON-логи, чтобы их можно было агрегировать.
- asyncio debug mode
- Включаю на стендах/в тестах:
- Настраиваю
- Совмещаю с расширенными warnings и строгим закрытием ресурсов, чтобы находить “висящие” задачи и незакрытые transport/session.
- tracemalloc (+ поиск утечек)
- Включаю ранний старт:
- Снимаю снапшоты до/после нагрузки, сравниваю
- Полезно для утечек: неосвобождённые буферы, рост кэшей, накопление задач/фьюч.
- Мониторинг latency event loop (loop lag)
- Замеряю “дрейф” планирования: периодический
- Метрика: max/avg/p95 лаг; алерты при росте (часто причина — CPU-bound в event loop, синхронные блокировки, слишком тяжёлые колбэки).
- Интеграция с метриками/трейсингом
- Prometheus/OpenTelemetry: длительность обработчиков, очередь задач, активные таски, таймауты, ретраи, loop lag.
- Распределённый трейсинг: видно, где “застревает” await-цепочка (БД/HTTP/очереди).
- Профилирование CPU в асинхронном коде
- py-spy (без внедрения, удобно в проде), yappi (умеет wall/CPU и потоки), cProfile точечно.
- Ищу CPU-bound участки в event loop и выношу в executor/отдельный сервис.
- Отладка задач и “висяков”
- Дамп всех задач:
- На SIGUSR1/SIGTERM можно печатать состояние задач в лог, чтобы понять, где блокировка.
- Практический минимум, который почти всегда включаю
- Структурное логирование + корреляция запросов.
- Метрика loop lag + метрики таймаутов/ошибок I/O.
- Debug mode и
-
Если скажешь стек (FastAPI/aiohttp, asyncpg, httpx и т.п.) и где ловишь проблему (утечки/лаг/висяки/таймауты), подскажу конкретную конфигурацию и метрики под твой кейс.
- Логирование (structlog / logging)
- Корреляция запросов:
request_id/trace_id, добавляю в контекст (contextvars) и в каждый лог.- Логи задач: имя/идентификатор asyncio-task (
asyncio.current_task()), время выполнения, таймауты, отмены (CancelledError).- Логи медленных операций: оборачиваю I/O и внешние вызовы в измерители длительности; пороги (например >100–300ms) логирую как warning.
- Для продакшена предпочитаю структурные JSON-логи, чтобы их можно было агрегировать.
- asyncio debug mode
- Включаю на стендах/в тестах:
PYTHONASYNCIODEBUG=1 или loop.set_debug(True).- Настраиваю
loop.slow_callback_duration (например 0.05–0.2), чтобы ловить “долгие” колбэки/таски.- Совмещаю с расширенными warnings и строгим закрытием ресурсов, чтобы находить “висящие” задачи и незакрытые transport/session.
import asyncio
loop = asyncio.get_event_loop()
loop.set_debug(True)
loop.slow_callback_duration = 0.1
- tracemalloc (+ поиск утечек)
- Включаю ранний старт:
tracemalloc.start(nframes) (обычно 25–50 кадров).- Снимаю снапшоты до/после нагрузки, сравниваю
snapshot.compare_to.- Полезно для утечек: неосвобождённые буферы, рост кэшей, накопление задач/фьюч.
import tracemalloc
tracemalloc.start(50)
# ... нагрузка ...
s1 = tracemalloc.take_snapshot()
# ... ещё нагрузка ...
s2 = tracemalloc.take_snapshot()
top = s2.compare_to(s1, "lineno")[:20]
for stat in top:
print(stat)
- Мониторинг latency event loop (loop lag)
- Замеряю “дрейф” планирования: периодический
await asyncio.sleep(dt) и измерение фактической задержки.- Метрика: max/avg/p95 лаг; алерты при росте (часто причина — CPU-bound в event loop, синхронные блокировки, слишком тяжёлые колбэки).
import asyncio, time
async def loop_lag_monitor(dt: float = 0.1):
next_t = time.perf_counter() + dt
while True:
await asyncio.sleep(dt)
now = time.perf_counter()
lag = max(0.0, now - next_t)
# отправить lag в метрики/лог
print("loop_lag_sec", lag)
next_t = now + dt
- Интеграция с метриками/трейсингом
- Prometheus/OpenTelemetry: длительность обработчиков, очередь задач, активные таски, таймауты, ретраи, loop lag.
- Распределённый трейсинг: видно, где “застревает” await-цепочка (БД/HTTP/очереди).
- Профилирование CPU в асинхронном коде
- py-spy (без внедрения, удобно в проде), yappi (умеет wall/CPU и потоки), cProfile точечно.
- Ищу CPU-bound участки в event loop и выношу в executor/отдельный сервис.
- Отладка задач и “висяков”
- Дамп всех задач:
asyncio.all_tasks() + их stack (task.get_stack()) при зависаниях.- На SIGUSR1/SIGTERM можно печатать состояние задач в лог, чтобы понять, где блокировка.
import asyncio, traceback
def dump_tasks():
for t in asyncio.all_tasks():
print("TASK:", t, "done=", t.done())
for frame in t.get_stack(limit=10):
traceback.print_stack(frame)
- Практический минимум, который почти всегда включаю
- Структурное логирование + корреляция запросов.
- Метрика loop lag + метрики таймаутов/ошибок I/O.
- Debug mode и
slow_callback_duration на тестовых окружениях.-
tracemalloc по необходимости при подозрении на утечку памяти.Если скажешь стек (FastAPI/aiohttp, asyncpg, httpx и т.п.) и где ловишь проблему (утечки/лаг/висяки/таймауты), подскажу конкретную конфигурацию и метрики под твой кейс.
Loop lag — это разница между ожидаемым временем выполнения callback/таймера в event loop и фактическим. В продакшене важно одновременно измерять лаг (метрики) и локализовать источник (профилирование и трассировка).
Как измерить в продакшене
- Периодический тикер: ставим
- Метрики: публикуем max/avg/p95 lag в мониторинг (Prometheus/StatsD/OpenTelemetry metrics).
- Корреляция: параллельно собираем метрики CPU, GC, количество активных тасков, длины очередей, время блокирующих I/O, чтобы лаг связывать с нагрузкой.
- Алерты: пороги по p95/p99 и max (например, p95 > 50–100ms для latency-sensitive сервисов).
Почему так
- Используется
-
- Скользящее окно даёт стабильные p95/max без хранения всей истории.
Сложность
- По времени: O(1) на тик (добавление в deque). Печать/квантили в примере периодически: O(W) на расчёт, где
- По памяти: O(W) на окно.
Частые причины loop lag
- CPU-bound код в event loop: тяжёлые вычисления, большие JSON (де)сериализации, компрессия, криптография, сортировки, большие регулярки.
- Блокирующий I/O внутри async: синхронный DNS, файловые операции, обращения к БД через sync-драйвер, HTTP через requests, ожидание локов.
- Большие монолитные корутины без yield: долгие циклы без
- Stop-the-world эффекты: GC-паузы (особенно при аллокациях), создание множества объектов, фрагментация.
- Перегрузка планировщика: слишком много задач/таймеров, высокая частота мелких callback, чрезмерные контекстные переключения.
- Инфраструктура: CPU throttling в контейнерах, noisy neighbor, нехватка CPU, частые page faults, медленный диск/сетевой FS.
Как уменьшить loop lag
- Вынести CPU-bound:
- использовать
- применить более быстрые библиотеки (например, быстрый JSON), сократить аллокации.
- Убрать блокирующие вызовы:
- только async-драйверы (HTTP, БД, Redis), async DNS при необходимости.
- файловые операции через thread pool либо асинхронные подходы.
- Дробить большие работы:
- обрабатывать батчи порциями и периодически делать
- Ограничивать конкуренцию:
- семафоры, backpressure, лимиты на in-flight запросы, очереди с bounded size.
- Снизить давление на GC:
- уменьшить количество временных объектов, реюз буферов, следить за большими списками/строками, при необходимости тюнить GC (осторожно, только после измерений).
- Наблюдаемость для локализации:
- профили (py-spy/austin) на CPU, трассировка медленных корутин, логирование длительных обработчиков, метрики очередей и времени ответов зависимостей.
Практический минимум для продакшена
- метрика loop_lag (p95/max) + алерт
- профилирование CPU по сигналу/по расписанию
- аудит на блокирующие вызовы (линтеры/код-ревью) и ограничение конкуренции на горячих путях
Как измерить в продакшене
- Периодический тикер: ставим
loop.call_later(dt, ...) и сравниваем реальное время с ожидаемым. Это даёт прямую оценку lag.- Метрики: публикуем max/avg/p95 lag в мониторинг (Prometheus/StatsD/OpenTelemetry metrics).
- Корреляция: параллельно собираем метрики CPU, GC, количество активных тасков, длины очередей, время блокирующих I/O, чтобы лаг связывать с нагрузкой.
- Алерты: пороги по p95/p99 и max (например, p95 > 50–100ms для latency-sensitive сервисов).
import asyncio
import time
import statistics
from collections import deque
class LoopLagMonitor:
def __init__(self, loop: asyncio.AbstractEventLoop, interval: float = 0.1, window: int = 600):
self.loop = loop
self.interval = interval
self.samples = deque(maxlen=window)
self._handle = None
self._stopped = False
def start(self) -> None:
self._stopped = False
self._schedule()
def stop(self) -> None:
self._stopped = True
if self._handle is not None:
self._handle.cancel()
self._handle = None
def _schedule(self) -> None:
expected = self.loop.time() + self.interval
self._handle = self.loop.call_at(expected, self._tick, expected)
def _tick(self, expected: float) -> None:
if self._stopped:
return
now = self.loop.time()
lag = max(0.0, now - expected)
self.samples.append(lag)
if len(self.samples) % 50 == 0:
data = list(self.samples)
p95 = statistics.quantiles(data, n=20)[18] if len(data) >= 20 else max(data)
print(f"loop_lag_sec avg={statistics.fmean(data):.6f} p95={p95:.6f} max={max(data):.6f}")
self._schedule()
async def main():
loop = asyncio.get_running_loop()
mon = LoopLagMonitor(loop, interval=0.1, window=600)
mon.start()
while True:
await asyncio.sleep(1)
if __name__ == "__main__":
asyncio.run(main())
Почему так
- Используется
loop.time() (монотонные часы event loop), чтобы измерение не ломалось от NTP/перевода системного времени.-
call_at/call_later измеряет именно задержку планирования и выполнения callback в loop, то есть целевой loop lag.- Скользящее окно даёт стабильные p95/max без хранения всей истории.
Сложность
- По времени: O(1) на тик (добавление в deque). Печать/квантили в примере периодически: O(W) на расчёт, где
W — размер окна (можно заменить на потоковые квантиль-оценки/гистограмму в метриках, чтобы было строго O(1)).- По памяти: O(W) на окно.
Частые причины loop lag
- CPU-bound код в event loop: тяжёлые вычисления, большие JSON (де)сериализации, компрессия, криптография, сортировки, большие регулярки.
- Блокирующий I/O внутри async: синхронный DNS, файловые операции, обращения к БД через sync-драйвер, HTTP через requests, ожидание локов.
- Большие монолитные корутины без yield: долгие циклы без
await, обработка больших батчей одним куском.- Stop-the-world эффекты: GC-паузы (особенно при аллокациях), создание множества объектов, фрагментация.
- Перегрузка планировщика: слишком много задач/таймеров, высокая частота мелких callback, чрезмерные контекстные переключения.
- Инфраструктура: CPU throttling в контейнерах, noisy neighbor, нехватка CPU, частые page faults, медленный диск/сетевой FS.
Как уменьшить loop lag
- Вынести CPU-bound:
- использовать
asyncio.to_thread для умеренных задач или ProcessPoolExecutor для тяжёлых/CPU-насыщенных.- применить более быстрые библиотеки (например, быстрый JSON), сократить аллокации.
- Убрать блокирующие вызовы:
- только async-драйверы (HTTP, БД, Redis), async DNS при необходимости.
- файловые операции через thread pool либо асинхронные подходы.
- Дробить большие работы:
- обрабатывать батчи порциями и периодически делать
await asyncio.sleep(0) или отдавать управление через очереди.- Ограничивать конкуренцию:
- семафоры, backpressure, лимиты на in-flight запросы, очереди с bounded size.
- Снизить давление на GC:
- уменьшить количество временных объектов, реюз буферов, следить за большими списками/строками, при необходимости тюнить GC (осторожно, только после измерений).
- Наблюдаемость для локализации:
- профили (py-spy/austin) на CPU, трассировка медленных корутин, логирование длительных обработчиков, метрики очередей и времени ответов зависимостей.
Практический минимум для продакшена
- метрика loop_lag (p95/max) + алерт
- профилирование CPU по сигналу/по расписанию
- аудит на блокирующие вызовы (линтеры/код-ревью) и ограничение конкуренции на горячих путях
Кейс
Сервис агрегирует котировки/статусы заказов из 3–5 внешних API (разные лимиты: 5 rps, 100 req/min, суточные квоты). Нужно обновлять данные для 200k сущностей каждые N минут, писать в БД батчами, не терять события, выдерживать всплески, корректно обрабатывать 429/5xx/таймауты.
Архитектура (Python, async)
- Два контура: сбор (I/O bound) и запись (DB bound), разделены очередью.
- Планировщик формирует задания (entity_id, provider, due_at, attempt) и кладёт в очередь запросов.
- Воркеры-сборщики (asyncio) берут задания, дергают API через
- Воркеры-записывальщики читают результаты, делают идемпотентный upsert в БД (asyncpg/SQLAlchemy async) батчами, подтверждают обработку.
- Хранилище очередей: для надежности лучше брокер (RabbitMQ/Redis Streams/Kafka). Для простоты/малых объемов возможен
- Идемпотентность: ключ
- Ограничение нагрузки: глобальная и per-provider квоты (rps, concurrency, burst), уважение
Лимиты, ретраи, таймауты
- Rate limit: per-provider лимитер (token bucket/leaky bucket). Дополнительно semaphores на concurrency.
- Ретраи: только на временные ошибки: таймауты, 429, 5xx, сетевые ошибки. Экспоненциальная пауза + джиттер, max_attempts, max_elapsed. Для 429 учитывать
- Таймауты: отдельные connect/read/write/pool таймауты, общий deadline на запрос.
- Circuit breaker: если провайдер деградировал (много ошибок), временно останавливаем запросы и быстро фейлим/откладываем задачи.
- Backpressure: если очередь на запись растёт, сбор замедляется (уменьшаем concurrency/приостанавливаем планировщик).
Очередь на запись в БД
- Пакетирование: копим до
- Разделение по таблицам/шардам при необходимости.
- Подтверждение сообщений после успешного коммита (at-least-once). Дубликаты режем идемпотентностью.
Набросок кода (упрощённо)
Проектирование устойчивости
- Семантика доставки: at-least-once + идемпотентная запись. Это практичнее, чем exactly-once.
- Планирование: хранить next_due в БД и поднимать задачи на горизонте (например, 5 минут) в брокер; при падении воркеров задачи не теряются.
- Наблюдаемость: метрики по провайдерам (rps, p95 latency, error rate, 429 rate, depth очередей, время от due_at до записи), трассировка, structured logs с correlation id.
- Конфигурация: per-provider профили (лимиты, таймауты, max_attempts, backoff).
Тестирование
- Unit:
- лимитер (не превышает rps/burst; корректно ждёт).
- backoff (экспонента, джиттер, уважение Retry-After).
- нормализация/валидация payload (pydantic-схемы).
- идемпотентный ключ и генерация upsert-запросов.
- Интеграционные:
- мок внешних API через
- тест БД через testcontainers (PostgreSQL) и проверка: дубликаты не плодятся, транзакции, батчи, конфликтные обновления.
- тест брокера (Redis Streams/RabbitMQ) в контейнере: подтверждения, переигрывание после падения воркера.
- Нагрузочные/хаос:
- прогон 10–50k задач, проверка p95, глубины очередей, времени обновления.
- искусственно ронять воркер/DB на 30–60 секунд: система должна восстановиться, не потеряв задачи, и догнать backlog.
- симуляция жёстких 429: проверка, что rps снижается, нет трешинга, очередь стабилизируется.
- Контрактные:
- фиксировать ожидаемые поля/типы ответов провайдера; при изменении схемы тест падает до продакшена.
Критерии готовности
- при 429/5xx нет лавинообразного ретрая, соблюдаются лимиты.
- при рестарте/краше воркеров данные не теряются, максимум дубликаты, которые подавляются upsert-ом.
- запись в БД не становится бутылочным горлом благодаря батчам и backpressure.
- метрики позволяют быстро увидеть деградацию конкретного провайдера и глубину очередей.
Сервис агрегирует котировки/статусы заказов из 3–5 внешних API (разные лимиты: 5 rps, 100 req/min, суточные квоты). Нужно обновлять данные для 200k сущностей каждые N минут, писать в БД батчами, не терять события, выдерживать всплески, корректно обрабатывать 429/5xx/таймауты.
Архитектура (Python, async)
- Два контура: сбор (I/O bound) и запись (DB bound), разделены очередью.
- Планировщик формирует задания (entity_id, provider, due_at, attempt) и кладёт в очередь запросов.
- Воркеры-сборщики (asyncio) берут задания, дергают API через
httpx.AsyncClient с таймаутами и ретраями, нормализуют ответ и кладут результат в очередь на запись.- Воркеры-записывальщики читают результаты, делают идемпотентный upsert в БД (asyncpg/SQLAlchemy async) батчами, подтверждают обработку.
- Хранилище очередей: для надежности лучше брокер (RabbitMQ/Redis Streams/Kafka). Для простоты/малых объемов возможен
asyncio.Queue, но без гарантии при падении процесса.- Идемпотентность: ключ
(provider, entity_id, version/observed_at); запись через INSERT ... ON CONFLICT DO UPDATE. Это позволяет безопасно ретраить и переигрывать сообщения.- Ограничение нагрузки: глобальная и per-provider квоты (rps, concurrency, burst), уважение
Retry-After.Лимиты, ретраи, таймауты
- Rate limit: per-provider лимитер (token bucket/leaky bucket). Дополнительно semaphores на concurrency.
- Ретраи: только на временные ошибки: таймауты, 429, 5xx, сетевые ошибки. Экспоненциальная пауза + джиттер, max_attempts, max_elapsed. Для 429 учитывать
Retry-After.- Таймауты: отдельные connect/read/write/pool таймауты, общий deadline на запрос.
- Circuit breaker: если провайдер деградировал (много ошибок), временно останавливаем запросы и быстро фейлим/откладываем задачи.
- Backpressure: если очередь на запись растёт, сбор замедляется (уменьшаем concurrency/приостанавливаем планировщик).
Очередь на запись в БД
- Пакетирование: копим до
batch_size или flush_interval и пишем одним запросом/транзакцией.- Разделение по таблицам/шардам при необходимости.
- Подтверждение сообщений после успешного коммита (at-least-once). Дубликаты режем идемпотентностью.
Набросок кода (упрощённо)
import asyncio
import random
from dataclasses import dataclass
from datetime import datetime, timezone
import httpx
@dataclass(frozen=True)
class Task:
provider: str
entity_id: str
attempt: int = 0
@dataclass(frozen=True)
class Result:
provider: str
entity_id: str
payload: dict
observed_at: datetime
class RateLimiter:
def __init__(self, rps: float, burst: int):
self._rps = rps
self._burst = burst
self._tokens = burst
self._updated = asyncio.get_event_loop().time()
self._lock = asyncio.Lock()
async def acquire(self):
async with self._lock:
now = asyncio.get_event_loop().time()
delta = now - self._updated
self._updated = now
self._tokens = min(self._burst, self._tokens + delta * self._rps)
if self._tokens < 1:
sleep_for = (1 - self._tokens) / self._rps
await asyncio.sleep(sleep_for)
self._tokens = 0
self._tokens -= 1
async def fetch_with_retry(client: httpx.AsyncClient, url: str, max_attempts=5):
base = 0.5
for attempt in range(max_attempts):
try:
resp = await client.get(url)
if resp.status_code == 429:
ra = resp.headers.get("Retry-After")
delay = float(ra) if ra else base * (2 ** attempt)
await asyncio.sleep(delay + random.random() * 0.2)
continue
if 500 <= resp.status_code <= 599:
raise httpx.HTTPStatusError("5xx", request=resp.request, response=resp)
resp.raise_for_status()
return resp.json()
except (httpx.TimeoutException, httpx.TransportError, httpx.HTTPStatusError):
if attempt == max_attempts - 1:
raise
delay = base * (2 ** attempt) + random.random() * 0.2
await asyncio.sleep(delay)
async def collector(name, in_q: asyncio.Queue, out_q: asyncio.Queue, limiter: RateLimiter):
timeout = httpx.Timeout(connect=2.0, read=5.0, write=2.0, pool=5.0)
async with httpx.AsyncClient(timeout=timeout) as client:
while True:
task = await in_q.get()
try:
await limiter.acquire()
url = f"https://api.example.com/{task.provider}/obj/{task.entity_id}"
data = await fetch_with_retry(client, url)
res = Result(task.provider, task.entity_id, data, datetime.now(timezone.utc))
await out_q.put(res)
finally:
in_q.task_done()
async def db_writer(out_q: asyncio.Queue, batch_size=200, flush_interval=0.5):
buf = []
last = asyncio.get_event_loop().time()
while True:
try:
res = await asyncio.wait_for(out_q.get(), timeout=flush_interval)
buf.append(res)
out_q.task_done()
except asyncio.TimeoutError:
pass
now = asyncio.get_event_loop().time()
if buf and (len(buf) >= batch_size or (now - last) >= flush_interval):
# upsert_batch(buf) # транзакция + ON CONFLICT
buf.clear()
last = now
Проектирование устойчивости
- Семантика доставки: at-least-once + идемпотентная запись. Это практичнее, чем exactly-once.
- Планирование: хранить next_due в БД и поднимать задачи на горизонте (например, 5 минут) в брокер; при падении воркеров задачи не теряются.
- Наблюдаемость: метрики по провайдерам (rps, p95 latency, error rate, 429 rate, depth очередей, время от due_at до записи), трассировка, structured logs с correlation id.
- Конфигурация: per-provider профили (лимиты, таймауты, max_attempts, backoff).
Тестирование
- Unit:
- лимитер (не превышает rps/burst; корректно ждёт).
- backoff (экспонента, джиттер, уважение Retry-After).
- нормализация/валидация payload (pydantic-схемы).
- идемпотентный ключ и генерация upsert-запросов.
- Интеграционные:
- мок внешних API через
respx или поднятый stub-сервер (fastapi) с сценариями 200/429/5xx/slow responses.- тест БД через testcontainers (PostgreSQL) и проверка: дубликаты не плодятся, транзакции, батчи, конфликтные обновления.
- тест брокера (Redis Streams/RabbitMQ) в контейнере: подтверждения, переигрывание после падения воркера.
- Нагрузочные/хаос:
- прогон 10–50k задач, проверка p95, глубины очередей, времени обновления.
- искусственно ронять воркер/DB на 30–60 секунд: система должна восстановиться, не потеряв задачи, и догнать backlog.
- симуляция жёстких 429: проверка, что rps снижается, нет трешинга, очередь стабилизируется.
- Контрактные:
- фиксировать ожидаемые поля/типы ответов провайдера; при изменении схемы тест падает до продакшена.
Критерии готовности
- при 429/5xx нет лавинообразного ретрая, соблюдаются лимиты.
- при рестарте/краше воркеров данные не теряются, максимум дубликаты, которые подавляются upsert-ом.
- запись в БД не становится бутылочным горлом благодаря батчам и backpressure.
- метрики позволяют быстро увидеть деградацию конкретного провайдера и глубину очередей.