Вам нужно вынести вычисление из event loop в отдельный процесс. Какой механизм помешает этому коду передать функцию исполнителю?
import asyncio
from concurrent.futures import ProcessPoolExecutor
async def main():
def square(x):
return x * x
loop = asyncio.get_running_loop()
with ProcessPoolExecutor() as pool:
result = await loop.run_in_executor(pool, square, 3)
print(result)
if __name__ == "__main__":
asyncio.run(main())
Код завершится ошибкой сериализации: локальную функцию square, объявленную внутри main, нельзя передать в ProcessPoolExecutor. Задача, вызываемая в отдельном процессе, вместе с аргументами должна быть передана через межпроцессный канал и обычно сериализуется с помощью pickle.
Вынесение функции на уровень модуля исправляет проблему. Это ограничение относится не к asyncio, а к способу обмена заданиями между процессами.
Процессы применяют для изоляции памяти и обхода ограничения GIL в CPython при CPU-bound вычислениях на чистом Python. В отличие от потоков, процессы не разделяют общую память напрямую, поэтому исполнитель должен передавать им функции и данные через механизм межпроцессного обмена.
ProcessPoolExecutor предоставляет единый интерфейс для запуска таких задач, а asyncio умеет ожидать результат исполнителя без блокировки event loop. Цена этого подхода — сериализация данных, создание процессов и дополнительные накладные расходы.
run_in_executor принимает вызываемый объект и аргументы, после чего передаёт их выбранному executor. Для ProcessPoolExecutor они должны быть пригодны для сериализации и восстановления в другом процессе.
Локальная функция зависит от области видимости конкретного вызова main и не является доступным для импорта объектом верхнего уровня модуля. Поэтому worker-процесс не может корректно получить ссылку на square.
Функцию следует определить на верхнем уровне модуля:
Теперь worker может импортировать функцию по имени модуля и выполнить её. Аргумент 3 и возвращаемый результат также должны быть сериализуемыми; например, обычные числа, строки, списки и словари обычно подходят, а открытые файлы, сокеты, event loop и многие объекты, связанные с потоками, — нет.
await loop.run_in_executor(...) приостанавливает только текущую корутину. Event loop продолжает обслуживать другие задачи, пока процесс выполняет вычисление. Само CPU-bound вычисление при этом происходит вне основного процесса Python.
Для платформ и режимов запуска, использующих spawn, особенно важна защита if __name__ == "__main__". Она предотвращает повторное создание пула при импорте основного модуля worker-процессом. На практике также важно не передавать в процессные задачи большие объёмы данных без необходимости: копирование и сериализация могут съесть выигрыш от параллелизма.
В сервисе нужно считать хэши больших файлов. Прямой вызов функции внутри корутины блокирует event loop. ThreadPoolExecutor проще применить и дешевле создать, но для чистого Python-кода вычисление может упереться в GIL. ProcessPoolExecutor лучше использует несколько ядер, однако требует сериализуемых аргументов и создаёт заметные накладные расходы.
Передавать в процесс целиком открытый файловый объект нельзя или нецелесообразно. Рассматривался вариант передавать содержимое файла, но это увеличивало потребление памяти и стоимость копирования. Выбрали передачу пути к файлу и диапазона байтов, а открытие файла выполняли уже внутри worker-процесса.
Такой вариант сохранил неблокирующий event loop и уменьшил объём межпроцессного обмена. Для маленьких файлов оставили обработку в потоке или обычном коде, поскольку запуск процесса был бы дороже самого вычисления.
Вопрос: Достаточно ли сделать функцию верхнеуровневой, если её аргументом является объект asyncio.Lock?
Ответ: Нет. Верхнеуровневость функции решает только проблему поиска вызываемого объекта. asyncio.Lock связан с асинхронным исполнением и не предназначен для передачи в другой процесс через pickle. Межпроцессную синхронизацию нужно строить через подходящие средства multiprocessing, IPC или внешнее хранилище, а не через asyncio-примитив.
Вопрос: Остановит ли task.cancel() уже выполняющееся вычисление в worker-процессе?
Ответ: Не обязательно. Отмена корутины отменяет ожидание результата со стороны event loop, но не гарантирует принудительное завершение уже исполняющейся функции в другом процессе. Вычисление может продолжиться и потреблять CPU. Если требуется управляемая остановка, её проектируют отдельно: например, через кооперативный флаг, разбивку работы на части или контролируемое завершение worker-процесса.
Вопрос: Почему изменение глобальной переменной внутри функции не является способом вернуть состояние из worker-процесса?
Ответ: У каждого процесса собственное адресное пространство. Worker получает сериализованную копию задания и изменяет свою копию глобального состояния; глобальная переменная родительского процесса от этого не меняется. Состояние нужно явно вернуть результатом задачи либо передать через IPC, очередь, общую память или внешнее хранилище — с учётом соответствующих затрат и синхронизации.