У обычного кода есть свойство, которое мы не замечаем: функция не возвращает управление, пока не закончила всё, что начала. У кода с create_task это свойство теряется: функция вернулась, а запущенные ею задачи живут дальше, падают в никуда и держат ресурсы. Структурная конкурентность возвращает это свойство: задачи живут внутри блока и не переживают его. В asyncio её воплощает TaskGroup, появившийся в Python 3.11, и в новом коде он заменяет gather почти везде.
Задача не переживает блок
async def enrich_order(order_id: int) -> Order:
async with asyncio.TaskGroup() as tg:
customer_task = tg.create_task(customers.get(order_id))
items_task = tg.create_task(items.list_for(order_id))
return Order(customer_task.result(), items_task.result())
Блок async with не завершится, пока обе задачи не закончатся. На выходе из блока гарантировано: все задачи либо вернули результат, либо упали, либо отменены, и ни одна не осталась висеть. Ссылки хранить не нужно, группа держит их сама. Это и есть структурность: дерево задач повторяет структуру кода, и по стеку вызовов видно, кто кого породил.
Первая ошибка отменяет остальных
Главное отличие от gather: если одна задача группы падает с исключением, TaskGroup отменяет все остальные задачи группы, дожидается их завершения и только потом поднимает исключение наружу. Проверено на Python 3.14: при ValueError в одной задаче соседняя получает CancelledError на своём await.
Это правильное поведение для обработчика запроса: если не удалось получить клиента, нет смысла ждать позиции заказа ещё две секунды. У gather соседи продолжили бы работать после того, как обработчик уже вернул 500, и держали бы соединения из пула.
Исключения поднимаются как ExceptionGroup, даже если упала одна задача, потому что упасть могли несколько одновременно. Ловят их синтаксисом except*:
try:
async with asyncio.TaskGroup() as tg:
tg.create_task(customers.get(order_id))
tg.create_task(items.list_for(order_id))
except* CustomerNotFound as eg:
raise HTTPException(404) from eg.exceptions[0]
except* httpx.HTTPError as eg:
raise ServiceUnavailable() from eg.exceptions[0]
except* срабатывает для каждого типа, который есть в группе, и может выполниться несколько раз для одной группы. Обычный except ValueError группу не поймает, и это второй источник сюрпризов при переходе на TaskGroup: глобальный обработчик ошибок FastAPI, который ждёт AppError, получит ExceptionGroup и ответит 500. Обработчик либо разворачивает группу (eg.exceptions), либо регистрирует обработчик для ExceptionGroup, который достаёт первое прикладное исключение.
Когда gather всё ещё уместен
gather остаётся для двух случаев. Первый: нужны результаты всех задач независимо от ошибок, то есть gather(..., return_exceptions=True), когда десять вызовов к разным поставщикам должны вернуть «что получилось», а не упасть на первом. Второй: совместимость с кодом до Python 3.11. В остальных случаях разница в пользу TaskGroup: отмена соседей, гарантированное завершение, явные ошибки.
Сравнение в одной таблице:
gather | TaskGroup | |
|---|---|---|
| Ошибка в одной задаче | перебрасывается первая, остальные работают дальше | остальные отменяются, затем ExceptionGroup |
| Ссылки на задачи | нужно хранить самому | хранит группа |
| Отмена внешней задачи | отменяет дочерние | отменяет дочерние и ждёт их |
| Результаты | список в порядке аргументов | через task.result() у каждой |
| Частичные результаты при ошибках | return_exceptions=True | нет, нужно ловить внутри задач |
Ограничение числа одновременных задач
TaskGroup не ограничивает число задач: цикл по десяти тысячам идентификаторов создаст десять тысяч задач и десять тысяч одновременных запросов к базе. Ограничение добавляют семафором внутри задачи или очередью с фиксированным числом потребителей:
async def fetch_all(ids: list[int], limit: int = 20) -> list[Item]:
sem = asyncio.Semaphore(limit)
async def one(item_id: int) -> Item:
async with sem:
return await client.get_item(item_id)
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(one(i)) for i in ids]
return [t.result() for t in tasks]
Десять тысяч задач создаются, но одновременно работают двадцать. Для больших потоков дешевле не создавать задачу на элемент, а запустить двадцать потребителей над asyncio.Queue, об этом статья про синхронизацию.
Фоновые задачи сервиса
Потребитель Kafka, опрос очереди, периодическая чистка живут столько же, сколько приложение. Структурный способ запустить их во FastAPI: TaskGroup внутри lifespan:
@asynccontextmanager
async def lifespan(app: FastAPI):
async with asyncio.TaskGroup() as tg:
tg.create_task(consume_orders(), name="kafka-consumer-orders")
tg.create_task(cleanup_loop(), name="cleanup")
yield
tg.cancel_scope.cancel() if hasattr(tg, "cancel_scope") else None
У asyncio.TaskGroup нет метода остановки группы: чтобы завершить фоновые задачи после yield, их отменяют поимённо или через событие остановки, которое задачи проверяют. Поэтому для фоновых задач часто берут anyio, где у группы есть cancel_scope.cancel(), и он же используется внутри Starlette:
import anyio
@asynccontextmanager
async def lifespan(app: FastAPI):
async with anyio.create_task_group() as tg:
tg.start_soon(consume_orders)
tg.start_soon(cleanup_loop)
yield
tg.cancel_scope.cancel()
На выходе из lifespan группа отменяет задачи и ждёт их; падение фоновой задачи во время работы приложения поднимется в группе и обрушит lifespan, что честнее тихо умершего потребителя. Если фоновая задача должна переживать свои ошибки, повтор с паузой пишут внутри неё.
Глубже: вложенные группы и отмена снаружирасширенное
Группы вкладываются: задача группы может открыть свою группу, и дерево отмены повторяет дерево кода. Если отменить внешнюю задачу (сработал внешний asyncio.timeout, сервис останавливается), отмена доходит до каждой задачи во всех вложенных группах, каждая группа дожидается своих детей, и только потом CancelledError поднимается наружу. Это и есть ответ на вопрос «что станет с запросами к базе при отмене обработчика»: они будут отменены, соединения вернутся в пул, и обработчик завершится без побочных эффектов, если не успел сделать commit. Уход клиента такой отменой не является: проверено на Starlette 1.7, что обработчик продолжает работать после разрыва соединения клиентом, и для долгих обработчиков это проверяют руками через request.is_disconnected().
Есть один случай, где структурность мешает: операция, которая должна завершиться даже при отмене обработчика, например, запись события аудита. Внутри TaskGroup её защищают asyncio.shield, но защищённая задача формально остаётся в группе, и группа дождётся её завершения. Поэтому такие операции чаще отдают отдельному долгоживущему компоненту (outbox, очередь), а не защищают на месте.
И про anyio: его TaskGroup старше и богаче (start с ожиданием готовности, cancel_scope, move_on_after), а семантика та же, что у asyncio.TaskGroup, который с него и списан. В приложении на FastAPI смешивать их можно, потому что anyio работает поверх asyncio, но стоит выбрать один стиль на проект.
Коротко
TaskGroup: задачи живут внутриasync with, блок не завершится, пока все не закончатся; ссылки хранить не нужно.- Первая ошибка отменяет остальные задачи группы, затем наружу поднимается
ExceptionGroup; ловить черезexcept*, обычныйexceptгруппу не видит. - Глобальный обработчик ошибок должен уметь разворачивать
ExceptionGroup, иначе прикладные ошибки превратятся в 500. gatherостаётся дляreturn_exceptions=Trueи старого кода; в остальномTaskGroupбезопаснее.- Группа не ограничивает число задач: семафор внутри задачи или очередь с фиксированным числом потребителей.
- Фоновые задачи сервиса запускают в группе внутри
lifespan; уanyioестьcancel_scope.cancel()для остановки послеyield. - Отмена внешней задачи доходит до всех вложенных групп и ждёт их; неделимые операции лучше отдавать outbox, чем
shield.
Что почитать дальше
- Синхронизация — Semaphore, Queue и потребители с ограничением.
- Типичные ошибки конкурентности — забытые задачи, проглоченная отмена и другие грабли.
- Глобальная обработка ошибок на Python — куда встроить разворачивание
ExceptionGroup.