← Все темы

Asyncio python

Вопросов: 41

asyncio — это стандартная библиотека Python для асинхронного (неблокирующего) выполнения задач в одном потоке на основе цикла событий (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 чаще всего это 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. Почему: 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:
- Планирование корутин: запускает корутины, переключает их, когда они делают 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 в 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) Самый рекомендуемый способ: 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 короутины используются для асинхронного кода и работают через 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 приостанавливает выполнение текущей асинхронной функции до тех пор, пока не завершится ожидаемый объект (обычно 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 defSyntaxError.
- забыть await при вызове корутины → получите объект-корутину, а не результат (и возможное предупреждение о “coroutine was never awaited”).

Awaitable-объекты — это объекты, которые можно использовать с 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.

Корутина — это функция, объявленная через 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-функции, например: 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.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.

Последовательный 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-задач. Основные инструменты: 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/pendingasyncio.wait

Дополнение: для таймаутов часто удобнее использовать asyncio.wait_for (таймаут на конкретное ожидание) или asyncio.timeout (Python 3.11+), а для “потоковой” обработки результатов по мере готовности — asyncio.as_completed.

asyncio.gather запускает несколько awaitable параллельно и возвращает результаты в порядке аргументов. Обработка исключений зависит от 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) Что такое отмена и где она живёт
- 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
- Один общий сигнал остановки: asyncio.Event (stop_event)
- Все фоновые задачи периодически проверяют stop_event и корректно выходят
- На сигнал ОС (SIGINT/SIGTERM) выставляем stop_event
- Останавливаем «вход» (сервер, consumer, планировщик) чтобы не появлялись новые задачи
- Ждём завершения задач с таймаутом; затем отменяем оставшиеся
- Закрываем ресурсы в finally или через async with

import 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 не дольше 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 применять как общий верхний предел на всю операцию.

Идея: 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 — это асинхронный мьютекс для защиты критической секции от одновременного доступа нескольких корутин. Нужен, когда есть общий ресурс (память/кэш/файл/соединение), который нельзя корректно изменять параллельно, иначе будут гонки данных.
- Гарантия: в критическую секцию одновременно входит только одна корутина
- Применение: атомарное обновление общей структуры, ограничение параллельной записи, последовательный доступ к одному клиенту/соединению

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: сигнализация о наступлении условия (многие корутины могут “проснуться” одновременно)

Разница между 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).

Идея: конвейер из стадий downloadparsewrite, между стадиями asyncio.Queue. На каждой стадии несколько воркеров. Завершение — через sentinel (специальный объект), который прокидывается дальше по конвейеру.

Ключевые правила
- 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.

Асинхронный контекстный менеджер — это объект, который гарантирует корректное выделение и освобождение ресурса в асинхронном коде (даже при исключениях), не блокируя цикл событий. Используется с 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 with

class 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.

Как работает 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 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 ждут результат через 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, поэтому:
- 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.sleepawait asyncio.sleep
- requestshttpx.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) Периодическая задача с фиксированным интервалом (без дрейфа)

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
- Идея: оборачиваешь сетевую операцию в функцию, делаешь несколько попыток, между попытками 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 с лимитом, таймаутами и базовой обработкой ошибок
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 и не перегружать систему большим числом короткоживущих соединений.

Как это устроено
- Пул хранит набор открытых соединений, обычно раздельно по ключу вида (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 — это ограничение частоты: сколько запросов можно сделать за единицу времени (например, 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 в секунду»: если запросы быстрые, семафор пропустит много запросов за секунду; он ограничивает только число одновременно выполняющихся.

Цель: корректно поймать 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 (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 и т.п.) — выбирайте специализированную библиотеку.

Ключевая проблема: в 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.

Ключевое правило: зависимости направлены внутрь (transportappdomain; 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 через 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+) — это структурированная конкурентность: вы создаёте группу, добавляете задачи, и при выходе из блока гарантируется, что все задачи либо завершились успешно, либо были корректно отменены при ошибке.

Базовое использование
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)
- Корреляция запросов: 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 и фактическим. В продакшене важно одновременно измерять лаг (метрики) и локализовать источник (профилирование и трассировка).

Как измерить в продакшене
- Периодический тикер: ставим 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 через 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.

- метрики позволяют быстро увидеть деградацию конкретного провайдера и глубину очередей.