Разберите, почему этот способ обработки большого набора данных может резко увеличить пиковое потребление памяти:
import asyncio
async def fetch(item):
await asyncio.sleep(0)
return item
async def load(items):
return await asyncio.gather(*(fetch(item) for item in items))
asyncio.gather в этом примере создаёт и планирует задачу для каждого элемента, а затем удерживает все результаты до завершения всей группы. Поэтому пиковая память растёт примерно с числом одновременно созданных корутин и размером накопленных результатов, даже если каждая отдельная операция использует мало памяти.
Ограничение числа активных запросов через Semaphore уменьшит нагрузку на внешнюю систему, но само по себе не ограничит число созданных задач. Для ограничения памяти нужен конечный пул workers или другая форма обработки с ограниченным количеством задач и потоковой выдачей результатов.
Асинхронная модель Python предназначена для эффективного ожидания большого числа операций ввода-вывода без создания отдельного потока на каждую операцию. asyncio.gather решает задачу группового ожидания: запускает переданные awaitable-объекты и возвращает результаты в исходном порядке.
Такой интерфейс удобен, когда набор задач невелик или его размер заранее ограничен. При обработке миллионов элементов удобство группового ожидания конфликтует с требованием ограничить количество объектов, находящихся в памяти.
Выражение *(fetch(item) for item in items) разворачивается в полный набор аргументов до вызова gather. Затем gather отслеживает все переданные операции, а итоговый список результатов живёт до возврата из load.
Ошибочно считать, что асинхронность автоматически делает потребление памяти постоянным. При большом items в памяти одновременно находятся корутины или задачи, служебные структуры планировщика, исходные элементы и результаты. Если результаты крупные, их накопление становится основным источником пикового потребления.
gather планирует переданные корутины и ждёт их все. Каждая корутина сохраняет своё состояние между точками await, включая локальные переменные и служебные данные. Пока группа не завершена, эти состояния доступны планировщику и не могут быть освобождены.
Кроме того, gather формирует результирующую последовательность. Даже если результат первой операции уже готов, он обычно остаётся частью итогового набора до завершения остальных операций. Следовательно, gather не является потоковым конвейером.
Надёжный вариант — создать ограниченное число workers, передавать им элементы через asyncio.Queue и обрабатывать результаты по мере готовности:
Размер очереди и число workers ограничивают количество элементов, находящихся в обработке. Однако приведённый вариант всё ещё накапливает все результаты в output; если их можно сразу записывать в файл, базу данных или отправлять потребителю, это следует делать вместо списка.
Semaphore полезен для ограничения числа одновременных сетевых операций, но конструкция с миллионами уже созданных задач всё равно может потреблять много памяти. Для настоящего ограничения памяти нужно ограничивать не только активные операции, но и количество существующих задач и буферизованных результатов.
Сервис загружал документы из внешнего API. Первый вариант создавал через gather задачу для каждого документа: он был простым и сохранял порядок результатов, но при большой партии удерживал в памяти все задачи и ответы.
Рассматривались два решения. Semaphore ограничивал число запросов и защищал API, но не устранял память, занятую уже созданными задачами. Разбиение входа на фиксированные батчи снижало пик памяти, однако требовало выбирать размер батча и задерживало обработку результатов до завершения каждого батча.
Выбрали пул workers с ограниченной очередью. Число одновременно выполняемых запросов стало контролируемым, элементы начали обрабатываться по мере готовности, а результаты стали сразу передаваться следующему этапу. Если требовалось сохранить порядок, каждому результату присваивали индекс и использовали ограниченное буферизованное хранилище; если порядок не имел значения, результаты можно было отправлять немедленно.
1. Достаточно ли заменить gather на Semaphore?
Нет. Такой код ограничивает число операций внутри критической секции, но все корутины могут быть созданы заранее:
В памяти по-прежнему находятся задачи для всех items, а gather удерживает все результаты. Semaphore регулирует параллелизм, но не задаёт bounded-модель потребления памяти.
2. Освобождаются ли завершившиеся задачи до окончания gather?
Внутренние состояния завершившихся операций могут стать не нужны для продолжения их выполнения, но их результаты должны сохраняться, чтобы gather вернул полный список в правильном порядке. Поэтому завершение части задач не означает пропорционального освобождения памяти результатов.
Для обработки по мере готовности применяют потоковую архитектуру с workers и очередью либо специализированный шаблон, выдающий результаты без ожидания всей группы. При этом нужно отдельно контролировать буфер результатов, иначе память снова будет расти у потребителя.
3. Почему asyncio не делает память пропорциональной только числу активных запросов?
Потому что активность запроса и существование задачи — разные понятия. Задача, ожидающая места в семафоре или данных из очереди, всё ещё является объектом Python со своим состоянием и занимает память.
Поэтому оценивать систему нужно по нескольким величинам: числу созданных задач, размеру очередей, объёму результатов, размеру локальных данных корутин и числу одновременно выполняемых операций. Ограничение каждой из этих величин должно соответствовать архитектуре конкретного конвейера.