Асинхронность Многопоточность Многопроцессорность #
1. Что такое асинхронность и как она устроена? #
Асинхронный код - это подход к написанию кода, который позволяет выполнять несколько задач одновременно в рамках одного потока. Это достигается за счет использования асинхронных функций и корутин. В отличие от синхронного кода, который выполняет каждую задачу последовательно, асинхронный код может запустить несколько задач конкурентно.
Ключевой механизм: когда одна задача выполняет операцию ввода-вывода (например, ожидает ответ от сети, чтение файла или запрос к базе данных) она освобождает управление, позволяя работать другой задаче. Таким образом, сокращается время, которое поток мог бы проводить в ожидании результатов операции ввода-вывода.
Основное отличие от многопоточности: хотя в один момент времени выполняется только одна задача (как и при работе с GIL), переключение между задачами происходит не произвольно по решению операционной системы, а предсказуемо — только в местах, помеченных await, то есть в точках ожидания результатов IO операций, специально отмечаемых программистом.
Это позволяет одному потоку обрабатывать тысячи одновременных соединений, делая программу чрезвычайно эффективной для I/O-bound операций (задач, связанных с операциями ввода-вывода).
Механизм возврата к выполнению задачи:
Задача встречает await и приостанавливается, сообщая event loop: “Я жду I/O-операцию, займись другими задачами” Event loop продолжает выполнять другие готовые задачи Когда I/O-операция завершается, система оповещает event loop, что результат готов Event loop помещает приостановленную задачу обратно в очередь готовых к выполнению Когда доходит очередь, задача возобновляется именно с того места, где остановилась, получая результат операции Ключевые особенности:
Запомнилось состояние стека и локальные переменные Возобновление происходит точно после await Не нужно заново запускать функцию с начала Аналогия: Как закладка в книге — закрыли на интересном месте, потом открыли и продолжили читать с того же момента.
Примером использования асинхронного кода является библиотека asyncio в Python. Например, вот простой пример кода, который использует asyncio для запуска нескольких задач одновременно и ожидания их завершения: import asyncio
async def hello(): await asyncio.sleep(1) print(“Hello”)
async def world(): await asyncio.sleep(2) print(“World”)
async def main(): await asyncio.gather(hello(), world())
if name == ‘main’: asyncio.run(main()) В этом примере мы определяем 3 асинхронные функции: hello(), world() и main(). Функции hello() и world() печатают соответствующие сообщения и ждут 1 и 2 секунды соответственно. Функция main() запускает эти две функции одновременно с помощью asyncio.gather() и ждет, пока они завершат свою работу. Затем мы запускаем функцию main() с помощью asyncio.run(). В результате мы получим сообщения “Hello” и “World”, каждое через 1 и 2 секунды соответственно, при этом результаты двух задач были получены почти одновременно.
Основные преимущества асинхронного программирования:
Улучшенная отзывчивость: Позволяет обрабатывать множество задач конкурентно без блокировки основного потока исполнения (для I/O-bound операций). Параллелизм — одновременное выполнение нескольких задач (на разных ядрах/процессорах). Конкурентность — управление несколькими задачами, которые выполняются в перекрывающиеся промежутки времени, но не обязательно одновременно (Asyncio). Эффективное использование ресурсов: Позволяет эффективно использовать процессорное время, так как задачи выполняются в моменты ожидания операций ввода/вывода или других блокирующих операций. Простота масштабирования: Позволяет легко создавать множество параллельных (конкурентных) задач без создания большого количества потоков или процессов.
2. Что такое конкурентность и параллельность? #
Конкурентность — это когда программа умеет вести несколько задач “одновременно” по времени, переключаясь между ними.
Параллельность — это когда несколько задач реально выполняются в один и тот же момент, например на разных ядрах CPU.
Конкурентность:
один исполнитель переключается между задачами
Задача A: ███ wait ███ wait ███
Задача B: ███ wait ███ wait ███
CPU: A → B → A → B → A
Параллельность:
несколько исполнителей работают реально одновременно
CPU 1: █████ Задача A █████
CPU 2: █████ Задача B █████
Главное различие #
| Понятие | Смысл | Реально одновременно? |
|---|---|---|
| Конкурентность | Несколько задач продвигаются вперёд, но могут чередоваться | Не обязательно |
| Параллельность | Несколько задач выполняются физически одновременно | Да |
То есть:
параллельность ⊂ конкурентность
Любая параллельность — это конкурентность, но не всякая конкурентность — параллельность.
Пример из жизни #
Конкурентность:
Ты готовишь суп и жаришь мясо один.
Поставил суп вариться → пока ждёшь, жаришь мясо → вернулся к супу.
Ты один, но продвигаешь две задачи.
Параллельность:
Один человек готовит суп.
Другой человек жарит мясо.
Оба работают одновременно.
В Python #
В Python конкурентность часто делают через:
asyncio
threading
concurrent.futures.ThreadPoolExecutor
asyncio в официальной документации Python описан как библиотека для написания конкурентного кода через async/await; он особенно полезен, когда программа часто ждёт I/O: сеть, базу данных, файлы, HTTP-запросы. (
Python documentation)
Пример конкурентности через asyncio:
import asyncio
async def download_file(name: str):
print(f"start {name}")
await asyncio.sleep(2) # имитация ожидания сети
print(f"finish {name}")
async def main():
await asyncio.gather(
download_file("file1"),
download_file("file2"),
download_file("file3"),
)
asyncio.run(main())
Здесь задачи не обязательно выполняются физически одновременно. Они переключаются в моменты ожидания.
Параллельность в Python #
Для настоящей параллельности в Python часто используют процессы:
from concurrent.futures import ProcessPoolExecutor
def heavy_calculation(n: int):
return sum(i * i for i in range(n))
if __name__ == "__main__":
with ProcessPoolExecutor() as executor:
results = executor.map(heavy_calculation, [10_000_000, 10_000_000, 10_000_000])
print(list(results))
concurrent.futures поддерживает запуск задач через потоки, процессы и интерпретаторы; ProcessPoolExecutor использует отдельные процессы, поэтому CPU-bound задачи могут реально выполняться параллельно на разных ядрах.
Важный момент про GIL #
В обычном CPython с включённым GIL несколько потоков не выполняют Python-байткод одновременно на разных ядрах. GIL ограничивает одновременное выполнение Python-кода потоками, хотя потоки всё равно полезны для I/O-задач.
Поэтому:
I/O-bound задачи:
asyncio / threads часто подходят
CPU-bound задачи:
processes / multiprocessing / ProcessPoolExecutor чаще подходят
CPU-bound и I/O-bound #
CPU-bound — задача грузит процессор:
математика
сжатие данных
обработка изображений
парсинг большого объёма данных
Для них нужна параллельность.
I/O-bound — задача ждёт внешние ресурсы:
запрос в БД
HTTP-запрос
чтение файла
ответ от API
Для них часто достаточно конкурентности.
Итог #
Конкурентность = уметь работать с несколькими задачами в один период времени.
Параллельность = реально выполнять несколько задач одновременно.
В Python:
asyncio → конкурентность
threading → конкурентность, хорошо для I/O
multiprocessing / ProcessPoolExecutor → параллельность
3. Чем отличаются асинхронность, многопроцессорность и многопоточность? | Виды конкурентности #
Асинхронность, многопоточность и многопроцессность — это разные способы организовать конкурентность, то есть выполнение нескольких задач в один период времени.
Конкурентность
├── Асинхронность
├── Многопоточность
└── Многопроцессность
Важно: правильнее говорить многопроцессность, если речь про multiprocessing.
Многопроцессорность — это скорее про железо: несколько процессоров/ядер.
1. Асинхронность #
Асинхронность — это когда один поток выполнения переключается между задачами в моменты ожидания.
Например:
Задача A делает HTTP-запрос и ждёт ответ
↓
Пока A ждёт, программа выполняет задачу B
↓
Потом возвращается к A
В Python это обычно:
async def func():
await something()
asyncio в официальной документации Python описан как библиотека для написания конкурентного кода через async/await; он часто используется в сетевых серверах, клиентах, библиотеках для БД и других I/O-задачах.
Пример:
import asyncio
async def request(name: str):
print(f"start {name}")
await asyncio.sleep(2) # имитация ожидания сети
print(f"finish {name}")
async def main():
await asyncio.gather(
request("A"),
request("B"),
request("C"),
)
asyncio.run(main())
Здесь задачи не выполняют Python-код одновременно. Они кооперативно уступают управление через await.
Асинхронность:
один поток
один event loop
много задач
переключение в await
Хорошо подходит для:
HTTP-запросов
запросов в БД
WebSocket
сетевых сервисов
очередей
ожидания файлов/сокетов
Плохо подходит для тяжёлых CPU-вычислений, потому что одна тяжёлая функция может заблокировать event loop.
2. Многопоточность #
Многопоточность — это когда внутри одного процесса создаётся несколько потоков.
Один процесс
├── Thread 1
├── Thread 2
└── Thread 3
Потоки имеют общую память процесса:
import threading
def worker():
print("work")
t1 = threading.Thread(target=worker)
t2 = threading.Thread(target=worker)
t1.start()
t2.start()
t1.join()
t2.join()
В Python модуль threading используется для потокового выполнения, но в обычном CPython с GIL только один поток за раз выполняет Python-байткод, поэтому прирост на CPU-bound задачах ограничен.
То есть потоки полезны, когда задача часто ждёт:
поток 1 ждёт ответ от сети
поток 2 в это время может работать
поток 3 ждёт файл
Но для тяжёлой математики:
Thread 1 считает
Thread 2 считает
Thread 3 считает
в обычном CPython это часто не даёт настоящего ускорения из-за GIL.
Многопоточность:
один процесс
несколько потоков
общая память
переключение потоков делает ОС
в CPython GIL мешает CPU-параллелизму
Хорошо подходит для:
блокирующего I/O
работы с библиотеками без async API
фоновых задач
простого параллельного ожидания
Плохо подходит для:
тяжёлых CPU-вычислений в чистом Python
3. Многопроцессность #
Многопроцессность — это когда создаётся несколько отдельных процессов.
Process 1
Process 2
Process 3
У каждого процесса своя память и свой интерпретатор Python.
from multiprocessing import Process
def worker():
print("work")
if __name__ == "__main__":
p1 = Process(target=worker)
p2 = Process(target=worker)
p1.start()
p2.start()
p1.join()
p2.join()
multiprocessing использует подпроцессы вместо потоков и тем самым обходит GIL, позволяя задействовать несколько процессоров/ядер.
Поэтому для CPU-bound задач обычно лучше:
from concurrent.futures import ProcessPoolExecutor
def heavy(n: int) -> int:
return sum(i * i for i in range(n))
if __name__ == "__main__":
with ProcessPoolExecutor() as executor:
results = executor.map(heavy, [10_000_000, 10_000_000, 10_000_000])
print(list(results))
concurrent.futures даёт единый интерфейс для асинхронного запуска задач через потоки, процессы или интерпретаторы: ThreadPoolExecutor, ProcessPoolExecutor, InterpreterPoolExecutor.
Многопроцессность:
несколько процессов
память изолирована
GIL не мешает между процессами
есть реальный CPU-параллелизм
обмен данными дороже
Хорошо подходит для:
математики
обработки изображений
сжатия данных
парсинга больших объёмов
ML/аналитики
CPU-bound задач
Минусы:
процессы тяжелее потоков
данные нужно передавать между процессами
сложнее общий state
больше расход памяти
Главное отличие #
| Подход | Что создаётся | Память | Реальный параллелизм в обычном CPython | Лучше для |
|---|---|---|---|---|
| Асинхронность | задачи/coroutines | общая | нет, обычно один поток | I/O-bound |
| Многопоточность | потоки | общая | ограничен GIL для Python-кода | I/O-bound |
| Многопроцессность | процессы | раздельная | да | CPU-bound |
CPU-bound и I/O-bound #
CPU-bound — задача упирается в процессор:
считать хэши
обрабатывать изображения
сжимать видео
делать тяжёлые вычисления
парсить огромный JSON
Для этого обычно нужна многопроцессность.
I/O-bound — задача в основном ждёт внешний ресурс:
запрос в БД
HTTP-запрос
чтение файла
ответ от API
сообщение из очереди
Для этого обычно хватает асинхронности или потоков.
Схема выбора #
Много сетевого ожидания?
↓
asyncio или threads
Тяжёлые вычисления на CPU?
↓
multiprocessing / ProcessPoolExecutor
Есть async-библиотеки?
↓
asyncio
Библиотека блокирующая и async API нет?
↓
ThreadPoolExecutor / threading
Нужно использовать несколько ядер CPU?
↓
multiprocessing / ProcessPoolExecutor
Итог #
Асинхронность — много задач в одном потоке, переключение через await.
Многопоточность — много потоков в одном процессе, общая память, удобно для I/O.
Многопроцессность — много процессов, отдельная память, хорошо для CPU и настоящего параллелизма.
В Python чаще всего:
asyncio → конкурентность для I/O
threading → конкурентность для I/O
multiprocessing → параллельность для CPU
ProcessPoolExecutor → удобная обёртка над процессами
ThreadPoolExecutor → удобная обёртка над потоками
4. Многопоточность vs Многопроцессорность | threading vs multiprocessing #
Главное различие #
threading — несколько потоков внутри одного процесса.multiprocessing — несколько отдельных процессов.
Точнее: в Python чаще говорят не “многопроцессорность”, а “многопроцессность”. Многопроцессорность — это скорее про железо: несколько CPU/ядер. multiprocessing — про запуск нескольких процессов, которые могут использовать эти ядра.
Коротко #
| Критерий | threading | multiprocessing |
|---|---|---|
| Единица выполнения | Поток | Процесс |
| Память | Общая память процесса | У каждого процесса своя память |
| GIL в обычном CPython | Один GIL на процесс | У каждого процесса свой GIL |
| CPU-bound задачи | Обычно плохо ускоряются | Хорошо подходят |
| I/O-bound задачи | Хорошо подходят | Обычно избыточно |
| Обмен данными | Проще, но нужны lock-и | Сложнее: queue, pipe, shared memory |
| Расход памяти | Меньше | Больше |
| Создание/переключение | Дешевле | Дороже |
| Риск race condition | Высокий из-за общей памяти | Ниже, потому что память изолирована |
threading
#
threading создаёт потоки внутри одного процесса:
Process
├── Thread 1
├── Thread 2
└── Thread 3
Потоки разделяют одну память:
import threading
counter = 0
def work():
global counter
for _ in range(100_000):
counter += 1
threads = [threading.Thread(target=work) for _ in range(4)]
for t in threads:
t.start()
for t in threads:
t.join()
print(counter)
Проблема: потоки работают с общими объектами. Поэтому возможны race condition, и для защиты общих данных нужны Lock, RLock, Semaphore, Queue и другие механизмы синхронизации.
В обычной сборке CPython GIL ограничивает выполнение Python-байткода: в один момент времени только один поток выполняет Python-код. Поэтому threading обычно не даёт сильного ускорения для CPU-bound задач. Но он полезен для I/O-bound задач: сеть, файлы, ожидание ответа БД, HTTP-запросы и т.д. Документация Python прямо указывает, что GIL ограничивает выгоду от threading для CPU-bound задач, но потоки всё равно полезны для конкурентного выполнения, особенно когда поток часто ждёт I/O.
multiprocessing
#
multiprocessing запускает отдельные процессы:
Main Process
├── Worker Process 1
├── Worker Process 2
└── Worker Process 3
У каждого процесса своя память и свой интерпретатор Python:
from multiprocessing import Process
def work():
total = 0
for i in range(10_000_000):
total += i
print(total)
processes = [Process(target=work) for _ in range(4)]
for p in processes:
p.start()
for p in processes:
p.join()
Так как процессы отдельные, они обходят ограничение одного GIL на процесс. Поэтому multiprocessing позволяет реально задействовать несколько CPU-ядер для CPU-bound задач. В документации Python сказано, что multiprocessing использует subprocesses вместо threads и таким образом обходит GIL, позволяя использовать несколько процессоров/ядер.
Минус: процессы тяжелее потоков. У них отдельная память, дороже запуск, сложнее обмен данными. Чтобы передать данные между процессами, обычно используют Queue, Pipe, Manager, shared_memory или сериализацию через pickle.
CPU-bound vs I/O-bound #
CPU-bound — задача упирается в вычисления процессора:
# CPU-bound
for i in range(100_000_000):
result += i * i
Для таких задач лучше:
multiprocessing
ProcessPoolExecutor
I/O-bound — задача большую часть времени ждёт внешний ресурс:
# I/O-bound
response = requests.get("https://example.com")
data = db.query(...)
file.read()
Для таких задач обычно подходят:
threading
ThreadPoolExecutor
asyncio
Python-документация по конкурентному выполнению прямо разделяет выбор инструмента по типу задачи: CPU-bound или I/O-bound.
Через concurrent.futures
#
На практике часто используют не Thread и Process напрямую, а высокоуровневые пулы:
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
Для I/O-bound:
from concurrent.futures import ThreadPoolExecutor
import requests
urls = [
"https://example.com",
"https://python.org",
]
def fetch(url):
return requests.get(url).status_code
with ThreadPoolExecutor(max_workers=4) as executor:
results = list(executor.map(fetch, urls))
print(results)
Для CPU-bound:
from concurrent.futures import ProcessPoolExecutor
def calc(n):
return sum(i * i for i in range(n))
with ProcessPoolExecutor(max_workers=4) as executor:
results = list(executor.map(calc, [10_000_000] * 4))
print(results)
concurrent.futures даёт единый интерфейс: ThreadPoolExecutor для потоков и ProcessPoolExecutor для процессов. ProcessPoolExecutor использует multiprocessing, что позволяет обходить GIL, но требует, чтобы передаваемые функции и данные могли сериализоваться.
Почему threading не ускоряет CPU-bound код в обычном CPython
#
Пример:
def cpu_task():
total = 0
for i in range(50_000_000):
total += i
return total
Если запустить такую функцию в нескольких потоках, потоки будут конкурировать за GIL. Они могут переключаться между собой, но не будут одновременно выполнять Python-байткод на разных ядрах в рамках одного процесса.
То есть:
Thread 1: работает
Thread 2: ждёт GIL
Thread 3: ждёт GIL
Thread 4: ждёт GIL
А при multiprocessing:
Process 1: работает на CPU core 1
Process 2: работает на CPU core 2
Process 3: работает на CPU core 3
Process 4: работает на CPU core 4
Поэтому для тяжёлых вычислений чаще выбирают процессы.
Когда выбирать threading
#
threading подходит, когда задача часто ждёт:
HTTP-запросы
работа с файлами
запросы к БД
сетевые операции
ожидание внешнего API
Пример:
# Много сетевых запросов
# Поток ждёт ответ сервера, в это время другой поток может работать
Плюсы:
+ меньше расход памяти
+ проще обмениваться объектами
+ быстрее запуск, чем у процессов
+ хорошо для I/O-bound
Минусы:
- GIL мешает CPU-bound ускорению
- общая память создаёт race condition
- нужны lock-и
- сложнее отлаживать ошибки конкуренции
Когда выбирать multiprocessing
#
multiprocessing подходит, когда задача реально грузит CPU:
математические вычисления
обработка изображений
парсинг больших файлов
сжатие/шифрование
тяжёлая обработка данных
ML/численные расчёты, если код не уходит в C-библиотеки
Плюсы:
+ реальный параллелизм на нескольких ядрах
+ обход GIL
+ изоляция памяти
+ падение одного процесса меньше портит состояние других
Минусы:
- больше расход RAM
- дороже запуск
- сложнее обмен данными
- данные часто нужно сериализовать
- на Windows особенно важно использовать if __name__ == "__main__"
Пример для Windows и кроссплатформенного кода:
from multiprocessing import Process
def work():
print("worker")
if __name__ == "__main__":
p = Process(target=work)
p.start()
p.join()
Важный нюанс про новые версии Python #
Начиная с Python 3.13, в CPython появилась поддержка экспериментальной free-threaded сборки, где GIL может быть отключён. В такой сборке потоки могут использовать несколько CPU-ядер полноценнее. Но это не отменяет стандартное правило для обычного CPython: в типичной сборке GIL всё ещё важен, и для CPU-bound задач чаще используют процессы.
Итоговая формула #
Много ожидания I/O?
threading / ThreadPoolExecutor / asyncio
Много чистых CPU-вычислений?
multiprocessing / ProcessPoolExecutor
Нужно простое конкурентное выполнение с общей памятью?
threading
Нужно реально загрузить несколько ядер CPU?
multiprocessing
Главная мысль:
threading = конкурентность внутри одного процесса
multiprocessing = параллельность через несколько процессов
В Python threading чаще используют для I/O-bound задач, а multiprocessing — для CPU-bound задач.
5. Что такое многопроцессорность, процесс и зачем нужен multiprocessing #
multiprocessing нужен, чтобы запускать код Python в нескольких отдельных процессах, а не в нескольких потоках.
Главная причина: в обычном CPython есть GIL, из-за которого несколько потоков не могут одновременно выполнять Python-байткод в одном процессе. multiprocessing обходит это ограничение, потому что каждый процесс имеет свой интерпретатор Python и свой GIL. Официальная документация Python прямо описывает multiprocessing как process-based parallelism и указывает, что он использует subprocesses вместо threads, позволяя задействовать несколько процессоров/ядер.
Что такое многопроцессорность #
Строго говоря, есть два похожих термина:
Многопроцессорность / multicore / multiprocessor
= несколько CPU или ядер могут выполнять работу параллельно
Многопроцессность / multiprocessing
= программа запускает несколько отдельных процессов
В Python, когда говорят про multiprocessing, обычно имеют в виду именно многопроцессность:
одна программа
├── процесс 1
├── процесс 2
├── процесс 3
└── процесс 4
Эти процессы операционная система может распределить по разным ядрам CPU.
Что такое процесс #
Процесс — это отдельный экземпляр выполняемой программы.
У процесса есть:
своё адресное пространство памяти
свои переменные
свой Python-интерпретатор
свой GIL
свои ресурсы ОС
один или несколько потоков внутри
Пример:
python main.py
Когда ты запускаешь этот файл, ОС создаёт процесс Python.
Упрощённо:
Process
├── память процесса
├── глобальные переменные
├── открытые файлы/сокеты
├── Python-интерпретатор
└── main thread
Поток — это более мелкая единица выполнения внутри процесса. Документация Python описывает threading как способ запускать несколько потоков, то есть smaller units of a process, внутри одного процесса.
Что делает multiprocessing
#
Модуль multiprocessing позволяет из Python-кода создавать и управлять отдельными процессами.
Пример:
from multiprocessing import Process
def work():
print("Работаю в отдельном процессе")
if __name__ == "__main__":
p = Process(target=work)
p.start()
p.join()
Что здесь происходит:
main process
└── создаёт child process
└── child process выполняет work()
p.start() запускает новый процесс.p.join() заставляет главный процесс дождаться завершения дочернего процесса.
Зачем нужен multiprocessing
#
Главная задача multiprocessing — выполнять CPU-bound задачи параллельно.
CPU-bound — это задачи, которые сильно грузят процессор:
математические вычисления
обработка больших файлов
обработка изображений
сжатие данных
криптографические вычисления
тяжёлый парсинг
расчёты
Пример CPU-bound задачи:
def calc(n):
total = 0
for i in range(n):
total += i * i
return total
Для такой задачи threading в обычном CPython обычно не даст нормального ускорения, потому что потоки будут конкурировать за GIL.
Один процесс + несколько потоков:
Process
├── Thread 1 ждёт GIL
├── Thread 2 выполняет Python-код
├── Thread 3 ждёт GIL
└── Thread 4 ждёт GIL
А multiprocessing создаёт несколько процессов:
Process 1 → CPU core 1
Process 2 → CPU core 2
Process 3 → CPU core 3
Process 4 → CPU core 4
Поэтому процессы могут реально выполнять Python-код параллельно на разных ядрах.
Почему не всегда использовать multiprocessing
#
Потому что процессы тяжелее потоков.
Минусы:
процессы занимают больше RAM
процесс создать дороже, чем поток
данные между процессами передавать сложнее
обычные Python-объекты не разделяются напрямую
часто нужна сериализация через pickle
Документация ProcessPoolExecutor отдельно указывает, что он использует multiprocessing, обходит GIL, но требует, чтобы вызываемые объекты и возвращаемые значения были picklable.
Почему у процессов отдельная память #
При threading потоки живут внутри одного процесса и видят общие объекты:
Process
├── Thread 1
├── Thread 2
└── shared memory
При multiprocessing каждый процесс изолирован:
Process 1 → memory 1
Process 2 → memory 2
Process 3 → memory 3
Поэтому такой код не работает так, как может показаться:
from multiprocessing import Process
counter = 0
def increment():
global counter
counter += 1
if __name__ == "__main__":
p1 = Process(target=increment)
p2 = Process(target=increment)
p1.start()
p2.start()
p1.join()
p2.join()
print(counter)
counter в главном процессе останется 0, потому что дочерние процессы меняли свои копии переменной, а не оригинальную переменную главного процесса.
Как процессы обмениваются данными #
Для обмена данными используют специальные механизмы:
Queue
Pipe
Manager
Value / Array
shared_memory
Пример с Queue:
from multiprocessing import Process, Queue
def worker(q):
result = 10 + 20
q.put(result)
if __name__ == "__main__":
q = Queue()
p = Process(target=worker, args=(q,))
p.start()
p.join()
result = q.get()
print(result)
Здесь дочерний процесс не меняет переменную главного процесса напрямую. Он кладёт результат в очередь.
Pool и ProcessPoolExecutor
#
Обычно вручную создавать Process нужно не всегда. Часто удобнее использовать пул процессов.
Через multiprocessing.Pool:
from multiprocessing import Pool
def calc(n):
return sum(i * i for i in range(n))
if __name__ == "__main__":
with Pool(processes=4) as pool:
results = pool.map(calc, [10_000_000, 10_000_000, 10_000_000, 10_000_000])
print(results)
Через concurrent.futures.ProcessPoolExecutor:
from concurrent.futures import ProcessPoolExecutor
def calc(n):
return sum(i * i for i in range(n))
if __name__ == "__main__":
with ProcessPoolExecutor(max_workers=4) as executor:
results = list(executor.map(calc, [10_000_000] * 4))
print(results)
ProcessPoolExecutor — более высокоуровневый интерфейс для запуска задач в пуле процессов. Документация Python описывает его как executor, который использует пул процессов для асинхронного выполнения вызовов.
Когда использовать multiprocessing
#
Используй multiprocessing, когда:
задача CPU-bound
нужно задействовать несколько ядер CPU
вычисления тяжёлые
потоки не дают ускорения из-за GIL
данные между задачами можно передавать относительно редко
Примеры:
обработать 10 000 изображений
посчитать тяжёлую математику
распарсить большой объём данных
запустить независимые CPU-задачи
обработать файлы параллельно
Когда лучше не использовать multiprocessing
#
Не лучший выбор, когда задача I/O-bound:
HTTP-запросы
запросы в БД
ожидание Redis
чтение файлов
сетевые операции
Для такого чаще подходят:
threading
ThreadPoolExecutor
asyncio
Официальная документация Python по конкурентному выполнению указывает, что выбор инструмента зависит от типа задачи: CPU-bound или I/O-bound.
Итог #
Процесс
= отдельный запущенный экземпляр программы со своей памятью
Многопроцессность
= запуск нескольких процессов
Многопроцессорность
= возможность использовать несколько CPU/ядер
multiprocessing
= Python-модуль для запуска нескольких процессов
Главная польза multiprocessing:
обойти GIL
загрузить несколько ядер CPU
ускорить CPU-bound задачи
изолировать выполнение по процессам
Главная цена:
больше памяти
дороже запуск
сложнее обмен данными
нужна сериализация объектов
6. Что такое и чем отличается вытесняющая многозадачность от кооперативной? #
Вытесняющая и кооперативная многозадачность отличаются тем, кто решает, когда одна задача должна уступить выполнение другой.
Вытесняющая многозадачность
= задачу может остановить планировщик ОС/рантайма
Кооперативная многозадачность
= задача сама добровольно отдаёт управление
Вытесняющая многозадачность #
Вытесняющая многозадачность — это когда планировщик может прервать выполняющуюся задачу без её согласия.
Например:
Thread 1 работает
↓
ОС решает: "хватит, теперь работает Thread 2"
↓
Thread 2 работает
↓
ОС снова переключает выполнение
То есть задача не обязана сама писать:
yield_control()
Её могут остановить извне.
Обычно так работают потоки ОС:
Process
├── Thread 1
├── Thread 2
└── Thread 3
Операционная система выдаёт каждому потоку немного CPU-времени. Когда квант времени закончился или появилась более приоритетная задача, планировщик может переключить CPU на другой поток. В документации Linux это описывается как preemption: задача может быть вытеснена, когда планировщик считает, что она уже использовала свою долю CPU-времени.
Пример вытесняющей многозадачности #
import threading
def task(name):
for i in range(5):
print(name, i)
t1 = threading.Thread(target=task, args=("A",))
t2 = threading.Thread(target=task, args=("B",))
t1.start()
t2.start()
t1.join()
t2.join()
Ты не указываешь вручную, где именно поток A должен уступить место потоку B.
Переключение происходит неявно:
A 0
A 1
B 0
A 2
B 1
B 2
...
Порядок может быть разным при разных запусках.
Главный плюс вытесняющей модели #
Плюс в том, что одна задача не обязана быть “вежливой”.
Даже если поток выполняет долгий цикл:
while True:
do_work()
ОС всё равно может переключить CPU на другие потоки или процессы.
То есть система остаётся более отзывчивой.
Главный минус вытесняющей модели #
Минус — задача может быть прервана почти в любой момент.
Например:
counter += 1
На уровне реального выполнения это не обязательно одна атомарная операция:
прочитать counter
увеличить значение
записать обратно
Если поток прервали между этими шагами, другой поток может изменить те же данные. Поэтому в вытесняющей модели часто нужны блокировки:
import threading
lock = threading.Lock()
counter = 0
def increment():
global counter
with lock:
counter += 1
Иначе возможны race condition.
Кооперативная многозадачность #
Кооперативная многозадачность — это когда задача сама решает, когда отдать управление.
Например:
Task 1 работает
↓
Task 1 дошла до await
↓
event loop запускает Task 2
↓
Task 2 дошла до await
↓
event loop возвращается к Task 1
В Python это типичная модель asyncio.
Официальная документация Python прямо говорит, что event loop использует cooperative scheduling: event loop выполняет одну Task за раз, а когда задача делает await и ждёт Future, event loop запускает другие задачи.
Пример кооперативной многозадачности #
import asyncio
async def task(name):
for i in range(3):
print(name, i)
await asyncio.sleep(1)
async def main():
await asyncio.gather(
task("A"),
task("B"),
)
asyncio.run(main())
Здесь переключение происходит на await:
A 0
await asyncio.sleep(1)
↓
B 0
await asyncio.sleep(1)
↓
A 1
...
await — это место, где корутина добровольно отдаёт управление event loop.
Важный пример проблемы #
import asyncio
async def bad_task():
while True:
pass
async def normal_task():
while True:
print("Я тоже хочу работать")
await asyncio.sleep(1)
async def main():
await asyncio.gather(
bad_task(),
normal_task(),
)
asyncio.run(main())
normal_task() почти не получит управление, потому что bad_task() никогда не делает await.
В кооперативной модели задача обязана сама уступать управление:
async def better_task():
while True:
# какая-то работа
await asyncio.sleep(0)
await asyncio.sleep(0) здесь означает: “отдай управление event loop, пусть другие задачи тоже выполнятся”.
Ключевое отличие #
Вытесняющая:
планировщик может остановить задачу сам
Кооперативная:
задача должна сама уступить управление
Сравнение:
| Критерий | Вытесняющая | Кооперативная |
|---|---|---|
| Кто переключает задачи | ОС/планировщик | Сама задача через await/yield |
| Можно ли остановить задачу без её согласия | Да | Обычно нет |
| Пример в Python | threading | asyncio |
| Риск race condition | Выше | Ниже внутри одного event loop |
| Риск зависания всей системы задач | Ниже | Выше, если задача не отдаёт управление |
| Подходит для | Потоки, процессы, ОС | Асинхронный I/O, корутины |
| Контроль над точками переключения | Меньше | Больше |
Как это связано с threading
#
threading в Python работает с потоками ОС. Потоки — это единицы выполнения внутри процесса. Документация Python описывает threading как модуль для запуска нескольких потоков внутри одного процесса.
Схема:
Process
├── Thread 1
├── Thread 2
└── Thread 3
Потоки могут быть переключены планировщиком в неожиданный момент:
Thread 1 начал менять общий список
↓
его вытеснили
↓
Thread 2 тоже меняет этот список
↓
получили некорректное состояние
Поэтому для общего состояния нужны синхронизация и осторожность.
Как это связано с asyncio
#
asyncio работает иначе:
Event loop
├── Task 1
├── Task 2
└── Task 3
Но event loop выполняет только одну задачу в конкретный момент времени:
Task 1 работает до await
Task 2 работает до await
Task 3 работает до await
Поэтому в asyncio меньше неожиданных переключений между строками кода. Но есть другое правило: нельзя надолго занимать event loop синхронным кодом.
Плохой вариант:
async def handler():
heavy_cpu_calculation()
Лучше:
async def handler():
result = await asyncio.to_thread(heavy_cpu_calculation)
Или для тяжёлых CPU-bound задач — вынести работу в процессы.
Простая аналогия #
Вытесняющая модель:
Учитель сам забирает слово у ученика и передаёт другому.
Кооперативная модель:
Ученик сам должен сказать: "Я закончил, пусть говорит другой".
Если ученик в кооперативной модели никогда не замолчит, остальные не смогут говорить.
Главное для Python #
threading
= ближе к вытесняющей многозадачности
asyncio
= кооперативная многозадачность
multiprocessing
= отдельные процессы, которые ОС тоже планирует вытесняюще
Итог #
Вытесняющая многозадачность удобна тем, что планировщик сам распределяет CPU между задачами, но из-за неожиданных переключений появляются race condition и нужны lock-и.
Кооперативная многозадачность удобна тем, что переключения происходят в понятных местах — обычно на await, но задача может заблокировать весь event loop, если долго работает без отдачи управления.
7. Что такое Гринлеты (greenlet) и чем они отличаются от Асинхронности #
greenlet — это лёгкая корутина, которая умеет вручную переключать выполнение с одного участка кода на другой внутри одного процесса и обычно внутри одного потока.
asyncio / асинхронность — это модель, где корутины переключаются через await, а их выполнением управляет event loop.
greenlet
= ручное переключение: g1.switch()
asyncio
= переключение через await под управлением event loop
Официальная документация greenlet описывает greenlets как lightweight coroutines для in-process concurrent programming. Документация asyncio описывает asyncio как библиотеку для concurrent code через синтаксис async/await.
Что такое greenlet #
greenlet — это объект, который хранит собственный стек выполнения и может быть приостановлен, а потом продолжен с того же места.
Пример:
from greenlet import greenlet
def task1():
print("task1: start")
g2.switch()
print("task1: end")
def task2():
print("task2: start")
g1.switch()
print("task2: end")
g1 = greenlet(task1)
g2 = greenlet(task2)
g1.switch()
Результат будет примерно такой:
task1: start
task2: start
task1: end
Что произошло:
g1 начал выполняться
↓
g1 вручную переключился на g2
↓
g2 начал выполняться
↓
g2 вручную переключился обратно на g1
↓
g1 продолжил выполнение
В greenlet переключение происходит через .switch(). Документация прямо говорит, что переключение между greenlet-ами происходит при вызове greenlet.switch() или greenlet.throw().
Главное свойство greenlet #
Greenlet может переключиться из глубоко вложенного вызова без await.
Например:
from greenlet import greenlet
def deep_function():
print("inside deep function")
main.switch()
print("back inside deep function")
def worker():
print("worker start")
deep_function()
print("worker end")
main = greenlet.getcurrent()
g = greenlet(worker)
g.switch()
print("main continues")
g.switch()
Здесь deep_function() может отдать управление наружу через main.switch(), хотя сама функция не объявлена как async def и внутри нет await.
Это важное отличие от asyncio.
Что такое асинхронность в Python #
В современном Python под асинхронностью обычно имеют в виду asyncio:
import asyncio
async def task1():
print("task1: start")
await asyncio.sleep(1)
print("task1: end")
async def task2():
print("task2: start")
await asyncio.sleep(1)
print("task2: end")
async def main():
await asyncio.gather(task1(), task2())
asyncio.run(main())
Здесь переключение происходит не через .switch(), а через await.
task1 работает до await
↓
event loop запускает task2
↓
task2 работает до await
↓
event loop возвращается к task1
Документация Python указывает, что event loop выполняет одну Task за раз; когда задача ждёт Future, event loop запускает другие задачи.
Основное отличие #
| Критерий | greenlet | asyncio |
|---|---|---|
| Как переключается выполнение | Через .switch() | Через await |
| Кто управляет задачами | Сам код / библиотека поверх greenlet | Event loop |
Нужно ли писать async def | Нет | Да |
Нужно ли писать await | Нет | Да |
| Уровень абстракции | Ниже | Выше |
| Основное применение | Основа для библиотек вроде gevent, внутренние механизмы | Асинхронные приложения, веб-серверы, клиенты БД, HTTP |
| Видимость переключений | Менее очевидная | Более явная |
| Параллелизм на CPU | Нет сам по себе | Нет сам по себе |
Greenlet — это не поток #
Greenlet не равен thread.
Thread
= поток ОС
Greenlet
= лёгкая корутина внутри процесса/потока
Обычно несколько greenlet-ов выполняются внутри одного OS thread:
Process
└── Thread
├── greenlet 1
├── greenlet 2
└── greenlet 3
Они не выполняют Python-код одновременно на разных ядрах. Они просто переключаются между собой.
Документация greenlet отдельно отмечает, что greenlet-ы можно сочетать с Python threads, но у каждого потока своё дерево greenlet-ов, и нельзя переключаться между greenlet-ами разных потоков.
Greenlet — это низкоуровневая корутина #
Можно сказать так:
greenlet
= низкоуровневый механизм переключения контекста
asyncio
= полноценная асинхронная модель с event loop, Task, Future, await
Сам по себе greenlet не делает сетевой код асинхронным.
Например:
import requests
from greenlet import greenlet
def task():
requests.get("https://example.com")
Если requests.get() блокирует поток, то greenlet сам по себе не спасёт. Он не умеет автоматически превращать блокирующий I/O в неблокирующий.
Для этого нужны библиотеки поверх greenlet, например gevent. Документация gevent описывает его как coroutine-based networking library, которая использует greenlet и event loop на базе libev/libuv.
Зачем тогда greenlet нужен #
greenlet полезен как строительный блок для библиотек, которые хотят дать “синхронный на вид” код, но внутри выполнять его конкурентно.
Например, вместо такого async-кода:
async def handler():
result = await db.fetch(...)
библиотека может дать более обычный стиль:
def handler():
result = db.fetch(...)
А внутри библиотека сама переключает greenlet-ы, когда операция ждёт I/O.
Идея:
код выглядит синхронно
↓
но библиотека внутри переключает greenlet-ы
↓
другие задачи могут выполняться во время ожидания I/O
Такой подход используют, например, gevent и некоторые внутренние механизмы библиотек.
В чём отличие от async/await
#
В asyncio точки переключения явно видны:
await something()
Ты сразу видишь: здесь функция может приостановиться.
В greenlet переключение может быть спрятано внутри обычного вызова:
result = db.fetch()
Снаружи это выглядит как обычная синхронная функция, но внутри она может переключить greenlet.
Сравнение:
asyncio:
явная асинхронность
greenlet:
скрытое переключение контекста
Пример на уровне идеи #
asyncio:
async def get_data():
data = await fetch()
return data
Точка переключения видна:
await fetch()
greenlet:
def get_data():
data = fetch()
return data
Точка переключения может быть внутри fetch():
fetch()
└── внутри где-то вызван greenlet.switch()
Связь с кооперативной многозадачностью #
И greenlet, и asyncio относятся к кооперативной модели.
То есть задача сама должна уступить управление:
greenlet → через switch()
asyncio → через await
Но разница в том, что asyncio встроен в язык и стандартную библиотеку как понятная модель async/await, а greenlet — сторонний низкоуровневый механизм переключения стеков.
Важное ограничение #
greenlet не решает проблему GIL.
greenlet ≠ multiprocessing
greenlet ≠ параллельное выполнение на нескольких ядрах
greenlet ≠ ускорение CPU-bound задач
Если задача CPU-bound:
def cpu_heavy():
for i in range(100_000_000):
...
то greenlet не сделает её параллельной. Для CPU-bound обычно нужны:
multiprocessing
ProcessPoolExecutor
нативные расширения, которые отпускают GIL
Итог #
greenlet
= низкоуровневая лёгкая корутина
= ручное переключение через .switch()
= может прятать переключение внутри обычных функций
= часто используется библиотеками вроде gevent
asyncio
= стандартная async-модель Python
= event loop + Task + Future + async/await
= переключение явно через await
= основной современный способ писать асинхронный I/O в Python
Главная разница:
asyncio заставляет явно писать async/await
greenlet позволяет переключать выполнение без async/await,
но сам по себе не даёт полноценный асинхронный I/O
Практически: для обычного современного Python-кода чаще выбирают asyncio. greenlet чаще встречается внутри библиотек или фреймворков, которым нужно дать синхронный API поверх конкурентного выполнения.
8. Какой тип многозадачности используют greenlet? #
greenlet используют кооперативную многозадачность.
То есть greenlet не вытесняется планировщиком ОС сам по себе. Он продолжает выполняться, пока сам не отдаст управление другому greenlet через switch().
Официальная документация greenlet описывает greenlet-ы как лёгкие корутины для конкурентного выполнения внутри процесса. В разделе про переключение сказано, что переход между greenlet-ами выполняется через g.switch(...). То есть переключение не происходит автоматически, его явно инициирует сам код.
Почему это кооперативная модель #
В кооперативной многозадачности задача сама уступает управление:
Greenlet A работает
↓
Greenlet A вызывает g2.switch()
↓
Greenlet B продолжает выполнение
↓
Greenlet B вызывает g1.switch()
↓
Greenlet A продолжает выполнение
Пример:
from greenlet import greenlet
def task1():
print("task1 start")
g2.switch()
print("task1 end")
def task2():
print("task2 start")
g1.switch()
print("task2 end")
g1 = greenlet(task1)
g2 = greenlet(task2)
g1.switch()
Здесь task1 не вытесняется извне. Она сама вызывает:
g2.switch()
и добровольно передаёт управление task2.
Важный момент #
Сам greenlet — это низкоуровневый механизм переключения контекста. Он не является полноценным планировщиком задач.
То есть greenlet сам по себе не делает так:
каждому greenlet дать по 10 мс CPU-времени
потом автоматически переключить на следующий
Так работает вытесняющая многозадачность у потоков ОС.
У greenlet иначе:
нет switch() → нет переключения
Документация greenlet прямо формулирует это так: переходы между greenlet-ами не являются неявными; greenlet должен сам выбрать переход к другому greenlet-у.
Сравнение с потоками #
threading
= вытесняющая многозадачность на уровне потоков ОС
greenlet
= кооперативная многозадачность внутри процесса/потока
Схема:
Process
└── OS Thread
├── greenlet 1
├── greenlet 2
└── greenlet 3
Обычно greenlet-ы живут внутри одного потока ОС. Поэтому они не выполняют Python-код параллельно на разных ядрах. Они просто переключаются между собой.
Если использовать gevent #
gevent построен поверх greenlet и добавляет event loop. В документации gevent прямо сказано, что greenlet-ы выполняются в одном OS thread и планируются кооперативно: пока конкретный greenlet не отдаст управление, остальные не получают шанс выполниться.
То есть:
greenlet сам по себе
= ручное кооперативное переключение через switch()
gevent
= кооперативная многозадачность поверх greenlet + event loop
Итог #
greenlet используют кооперативную многозадачность
Главная причина:
greenlet не прерывается автоматически;
он сам должен передать управление другому greenlet-у
Ключевая формула:
threading → вытесняющая модель
asyncio → кооперативная модель через await
greenlet → кооперативная модель через switch()
9. Для каких типов задач целесообразно использовать асинхронность, многопоточность, мультипроцессность (I/O-bound vs CPU-bound)? #
I/O-bound задача
= программа в основном ждёт внешний ресурс:
сеть, БД, файл, Redis, HTTP API
CPU-bound задача
= программа в основном грузит процессор:
вычисления, парсинг, обработка изображений, сжатие, криптография
Выбор в Python обычно такой:
I/O-bound:
asyncio
threading
ThreadPoolExecutor
CPU-bound:
multiprocessing
ProcessPoolExecutor
Основная таблица #
| Тип задачи | Что лучше использовать | Почему |
|---|---|---|
| Много HTTP-запросов | asyncio или threading | Основное время уходит на ожидание сети |
| Много запросов в БД | asyncio, если драйвер async; иначе threading | Пока один запрос ждёт БД, можно выполнять другой |
| Работа с файлами | threading или async-обёртки | Часто есть ожидание диска |
| WebSocket / long polling | asyncio | Много долгоживущих соединений |
| Web API на FastAPI | asyncio | Хорошо подходит под сетевой I/O |
| Тяжёлые вычисления | multiprocessing | Можно задействовать несколько CPU-ядер |
| Обработка изображений | multiprocessing, если чистый Python/CPU | Задача грузит CPU |
| Парсинг большого файла с тяжёлой логикой | multiprocessing | Много вычислений |
| Несколько блокирующих библиотек | threading | Можно вынести блокирующие вызовы в потоки |
Асинхронность #
Асинхронность в Python обычно означает asyncio: async def, await, event loop, Task, Future.
asyncio хорошо подходит для I/O-bound и сетевого кода. В официальной документации Python сказано, что asyncio используется для конкурентного кода через async/await и часто хорошо подходит для I/O-bound и высокоуровневого сетевого кода.
Пример задачи:
import asyncio
import httpx
async def fetch(url):
async with httpx.AsyncClient() as client:
response = await client.get(url)
return response.status_code
async def main():
results = await asyncio.gather(
fetch("https://example.com"),
fetch("https://python.org"),
)
print(results)
asyncio.run(main())
Смысл:
одна корутина ждёт ответ от сети
↓
event loop переключается на другую корутину
↓
CPU не простаивает зря
Асинхронность не ускоряет тяжёлые вычисления сама по себе. Если внутри async def написать долгий CPU-bound цикл без await, он заблокирует event loop.
Плохой вариант:
async def handler():
result = heavy_cpu_calculation()
return result
Лучше вынести тяжёлую CPU-задачу в процесс, а блокирующую I/O-задачу — в поток.
Многопоточность #
threading подходит для I/O-bound задач, особенно когда используешь синхронные блокирующие библиотеки.
Например:
requests
psycopg2
обычный boto3
синхронный SDK внешнего сервиса
Официальная документация Python указывает, что threading полезен для I/O-bound задач, например файловых операций или сетевых запросов, где большая часть времени уходит на ожидание внешних ресурсов.
Пример:
from concurrent.futures import ThreadPoolExecutor
import requests
def fetch(url):
response = requests.get(url)
return response.status_code
urls = [
"https://example.com",
"https://python.org",
]
with ThreadPoolExecutor(max_workers=10) as executor:
results = list(executor.map(fetch, urls))
print(results)
Смысл:
Thread 1 ждёт HTTP-ответ
Thread 2 ждёт HTTP-ответ
Thread 3 ждёт HTTP-ответ
...
Пока один поток ждёт сеть, другой может тоже выполнять полезную работу или ждать свой I/O.
Но для CPU-bound задач обычный threading в CPython обычно не помогает: из-за GIL только один поток в процессе выполняет Python-байткод одновременно. Документация Python прямо указывает, что для лучшего использования многоядерных машин стоит использовать multiprocessing или ProcessPoolExecutor, а threading остаётся подходящей моделью для одновременного выполнения I/O-bound задач.
Мультипроцессность #
multiprocessing подходит для CPU-bound задач.
Пример:
from concurrent.futures import ProcessPoolExecutor
def calc(n):
total = 0
for i in range(n):
total += i * i
return total
if __name__ == "__main__":
with ProcessPoolExecutor(max_workers=4) as executor:
results = list(executor.map(calc, [10_000_000] * 4))
print(results)
Смысл:
Process 1 → CPU core 1
Process 2 → CPU core 2
Process 3 → CPU core 3
Process 4 → CPU core 4
В Python с обычным GIL использование нескольких CPU-ядер обычно требует нескольких процессов или нативных расширений, которые отпускают GIL. Это указано в официальном глоссарии Python.
multiprocessing даёт реальный параллелизм для Python-кода, но цена выше:
больше расход RAM
дороже запуск процессов
сложнее обмен данными
нужна сериализация данных
Документация Python показывает, что multiprocessing.Pool используется для распределения работы по процессам, а ProcessPoolExecutor даёт более высокоуровневый интерфейс для выполнения задач в фоновых процессах.
Как выбирать на практике #
Задача ждёт сеть/БД/файл?
Да → I/O-bound
async-библиотеки есть?
Да → asyncio
Нет → threading / ThreadPoolExecutor
Задача грузит CPU?
Да → CPU-bound
multiprocessing / ProcessPoolExecutor
Задача смешанная?
I/O часть → asyncio/threading
CPU часть → вынести в процессы
Пример смешанной задачи #
Допустим, надо:
1. скачать 1000 страниц
2. распарсить HTML
3. посчитать тяжёлую статистику
4. сохранить результат в БД
Разделение:
скачать страницы → asyncio
тяжёлый парсинг/расчёт → ProcessPoolExecutor
сохранить в БД → asyncio, если async-драйвер
Упрощённая схема:
asyncio event loop
├── скачивает данные
├── ждёт сеть/БД
└── отправляет тяжёлую CPU-работу в process pool
Что выбрать для backend #
Для обычного backend на Python:
FastAPI + async DB driver + httpx.AsyncClient
→ asyncio
Django/DRF с синхронными ORM-вызовами
→ обычный sync-код + workers/threads на уровне сервера
Тяжёлая обработка файлов/изображений
→ отдельный worker/process pool/Celery workers
Много внешних API через синхронные SDK
→ ThreadPoolExecutor или отдельные workers
Ключевая формула #
asyncio
= много ожидания, мало CPU, async-библиотеки
threading
= много ожидания, но библиотеки синхронные/блокирующие
multiprocessing
= много CPU, нужно использовать несколько ядер
Самая частая ошибка:
CPU-bound задачу класть в asyncio
asyncio не делает вычисления параллельными. Оно эффективно тогда, когда задачи часто отдают управление на await, то есть ждут I/O.
10. Что такое и зачем нужны потоки в Python? #
Что такое поток #
Поток в Python — это отдельная линия выполнения внутри одного процесса.
Упрощённо:
Process
├── Main Thread
├── Thread 1
├── Thread 2
└── Thread 3
Один процесс может иметь несколько потоков. Эти потоки выполняют разные функции, но живут внутри одного процесса и используют общую память.
Документация Python описывает threading как модуль для запуска нескольких потоков внутри одного процесса; потоки особенно полезны для I/O-bound задач — например, сетевых запросов и файловых операций.
Зачем нужны потоки #
Потоки нужны, чтобы программа могла выполнять несколько задач конкурентно.
Например, без потоков:
скачать файл 1
↓
дождаться завершения
↓
скачать файл 2
↓
дождаться завершения
↓
скачать файл 3
С потоками:
Thread 1 → скачивает файл 1
Thread 2 → скачивает файл 2
Thread 3 → скачивает файл 3
Пока один поток ждёт ответ от сети, другой поток тоже может выполнять работу или ждать свой ответ.
Пример #
import threading
import time
def task(name):
print(f"{name}: start")
time.sleep(2)
print(f"{name}: end")
t1 = threading.Thread(target=task, args=("Thread 1",))
t2 = threading.Thread(target=task, args=("Thread 2",))
t1.start()
t2.start()
t1.join()
t2.join()
print("Done")
Что происходит:
t1.start() → запускает первый поток
t2.start() → запускает второй поток
t1.join() → ждём завершения первого потока
t2.join() → ждём завершения второго потока
Оба потока могут “спать” одновременно, поэтому программа завершится примерно за 2 секунды, а не за 4.
Где потоки полезны #
Потоки хорошо подходят для I/O-bound задач:
HTTP-запросы
запросы к БД
чтение/запись файлов
работа с сокетами
ожидание внешнего API
фоновые задачи в приложении
Пример с HTTP:
from concurrent.futures import ThreadPoolExecutor
import requests
def fetch(url):
response = requests.get(url)
return response.status_code
urls = [
"https://example.com",
"https://python.org",
]
with ThreadPoolExecutor(max_workers=4) as executor:
results = list(executor.map(fetch, urls))
print(results)
ThreadPoolExecutor использует пул потоков для асинхронного выполнения вызовов, то есть позволяет отправлять задачи в фоновые потоки через более удобный интерфейс, чем ручное создание Thread.
Почему потоки не всегда ускоряют код #
В обычном CPython есть GIL — Global Interpreter Lock.
Из-за GIL в одном процессе только один поток одновременно выполняет Python-байткод. Поэтому потоки обычно не дают сильного ускорения для CPU-bound задач. Это прямо указано в документации Python: GIL ограничивает выгоду threading для CPU-bound задач, хотя потоки остаются полезными для I/O-bound сценариев.
CPU-bound пример:
def calc():
total = 0
for i in range(100_000_000):
total += i * i
return total
Для такой задачи threading обычно не лучший выбор.
Схема:
Thread 1 хочет выполнять Python-код
Thread 2 хочет выполнять Python-код
Thread 3 хочет выполнять Python-код
Но GIL разрешает выполнять Python-байткод только одному потоку за раз
Для CPU-bound задач чаще используют:
multiprocessing
ProcessPoolExecutor
Важное свойство потоков: общая память #
Потоки внутри одного процесса видят одни и те же объекты:
import threading
counter = 0
def increment():
global counter
for _ in range(100_000):
counter += 1
threads = [
threading.Thread(target=increment),
threading.Thread(target=increment),
]
for thread in threads:
thread.start()
for thread in threads:
thread.join()
print(counter)
На первый взгляд кажется, что результат всегда должен быть:
200000
Но при работе с общими изменяемыми данными могут возникать race condition.
Проблема в том, что операция:
counter += 1
логически состоит из нескольких шагов:
прочитать counter
увеличить значение
записать обратно
Если потоки вмешаются друг в друга между этими шагами, результат может стать некорректным.
Для защиты используют Lock #
import threading
counter = 0
lock = threading.Lock()
def increment():
global counter
for _ in range(100_000):
with lock:
counter += 1
threads = [
threading.Thread(target=increment),
threading.Thread(target=increment),
]
for thread in threads:
thread.start()
for thread in threads:
thread.join()
print(counter)
Lock защищает критическую секцию:
один поток вошёл в with lock
↓
другие ждут
↓
первый поток закончил
↓
следующий поток входит
Потоки vs процессы #
threading
= несколько потоков внутри одного процесса
= общая память
= хорошо для I/O-bound
= хуже для CPU-bound из-за GIL
multiprocessing
= несколько отдельных процессов
= отдельная память
= лучше для CPU-bound
= дороже по памяти и запуску
Потоки vs asyncio #
threading
= вытесняющая многозадачность потоков ОС
asyncio
= кооперативная многозадачность через await
В потоках переключение может произойти неявно:
Thread 1 работает
↓
ОС/интерпретатор переключил выполнение
↓
Thread 2 работает
В asyncio переключение обычно видно явно:
await something()
Когда использовать потоки #
Потоки стоит использовать, когда:
задача часто ждёт внешний ресурс
используются синхронные библиотеки
нужно не блокировать основной поток
нужно выполнить несколько I/O-операций конкурентно
Примеры:
скачать много URL через requests
сделать много запросов к синхронной БД
работать с несколькими файлами
вынести блокирующую операцию из основного потока
Когда не стоит использовать потоки #
Потоки не лучший выбор, когда:
задача тяжёлая по CPU
много общего изменяемого состояния
нужен простой и предсказуемый порядок выполнения
можно использовать нормальные async-библиотеки
Для CPU-bound:
ProcessPoolExecutor / multiprocessing
Для async I/O:
asyncio
Итог #
Потоки в Python нужны для конкурентного выполнения задач внутри одного процесса.
Главная польза:
не простаивать во время ожидания I/O
выполнять блокирующие операции в фоне
обрабатывать несколько внешних операций одновременно
Главное ограничение:
из-за GIL потоки обычно не ускоряют CPU-bound Python-код
Формула выбора:
много ожидания I/O + синхронные библиотеки → threading
много ожидания I/O + async-библиотеки → asyncio
много CPU-вычислений → multiprocessing
11. Как реализована модель асинхронного выполнения? #
Асинхронное выполнение в Python обычно реализуется через asyncio: это модель конкурентного выполнения, где один поток может обслуживать много задач, переключаясь между ними в моменты ожидания I/O, например сети, БД, файловых операций через async-библиотеки. asyncio описывается в документации как библиотека для конкурентного кода с синтаксисом async / await.
Главная идея:
Одна задача ждёт I/O
↓
она отдаёт управление event loop
↓
event loop запускает другую готовую задачу
↓
когда I/O завершилось, первая задача продолжается
Это не означает, что Python-код внутри одного event loop выполняется параллельно на нескольких ядрах. Это кооперативная многозадачность: задача сама отдаёт управление через await.
Основные компоненты #
1. Coroutine #
Корутина — это объект, который получается при вызове async def функции:
async def get_data():
return 123
coro = get_data()
Важный момент: сам вызов get_data() ещё не выполняет тело функции. Корутина начнёт выполняться только когда её await-нут или запланируют как задачу. В документации Python прямо указано, что корутины и задачи могут выполняться только при работающем event loop.
result = await get_data()
2. Event loop #
event loop — центральный механизм asyncio. Он запускает асинхронные задачи, callback-функции, выполняет сетевой I/O и работает с подпроцессами.
Упрощённо:
event loop:
1. берёт готовую задачу
2. выполняет её до ближайшего await
3. если задача ждёт I/O — откладывает её
4. запускает другую готовую задачу
5. возвращается к первой, когда её ожидание завершилось
3. Task #
Task — это обёртка над корутиной, которая позволяет запланировать её выполнение в event loop. Например, asyncio.create_task() запускает корутину конкурентно с другими задачами.
task = asyncio.create_task(get_data())
То есть:
coroutine
↓ оборачивается в Task
Task попадает в event loop
↓
event loop решает, когда продолжить её выполнение
4. Future #
Future — объект, который представляет будущий результат асинхронной операции. Корутина может ждать Future, пока в нём не появится результат, исключение или отмена.
Примерно:
Future = "результат ещё не готов, но будет позже"
Когда результат готов, event loop возобновляет задачу, которая ждала этот Future.
Что делает await
#
await — это точка добровольной остановки корутины.
async def main():
data = await fetch_data()
print(data)
Когда выполнение доходит до await fetch_data():
main() временно приостанавливается
↓
управление возвращается event loop
↓
event loop запускает другие задачи
↓
fetch_data() завершилась
↓
main() продолжается с места await
Важно: await не создаёт новый поток. Он только говорит:
«я сейчас жду, можешь пока выполнять другие задачи».
Пример #
import asyncio
async def task(name, delay):
print(f"{name}: start")
await asyncio.sleep(delay)
print(f"{name}: end")
async def main():
t1 = asyncio.create_task(task("A", 2))
t2 = asyncio.create_task(task("B", 1))
await t1
await t2
asyncio.run(main())
Выполнение будет примерно таким:
A: start
B: start
B: end
A: end
Почему?
task A дошла до await sleep(2)
↓
отдала управление event loop
task B дошла до await sleep(1)
↓
отдала управление event loop
через 1 секунду B готова
↓
event loop продолжает B
через 2 секунды A готова
↓
event loop продолжает A
Важная деталь: внутри event loop задача выполняется одна за раз #
Обычно event loop работает в одном потоке. Документация Python указывает, что event loop выполняет callback-и и Tasks в своём потоке; пока одна Task выполняется, другие Tasks в этом же loop не выполняются.
То есть асинхронность в Python — это не так:
Task A реально выполняется на CPU
Task B реально выполняется на CPU
Task C реально выполняется на CPU
А так:
Task A выполняется до await
Task B выполняется до await
Task C выполняется до await
Task A продолжает после I/O
Task B продолжает после I/O
Почему это эффективно для I/O-bound задач #
Асинхронность полезна, когда программа много ждёт:
запрос в БД
HTTP-запрос
чтение из сокета
ожидание Redis
ожидание ответа внешнего API
В синхронной модели поток часто простаивает:
отправил запрос → ждёт ответ → ничего не делает
В асинхронной модели:
отправил запрос → переключился на другую задачу → вернулся, когда ответ готов
Поэтому asyncio часто используется как основа для сетевых серверов, асинхронных web-фреймворков, клиентов БД и очередей. Это также указано в официальной документации Python.
Почему это плохо для CPU-bound задач #
Если задача долго считает на CPU и не делает await, она блокирует event loop:
async def bad():
while True:
pass
Пока такой код крутится, другие async-задачи не получают управление.
Плохой вариант:
async def handler():
result = heavy_calculation() # долго считает CPU
return result
Лучше:
result = await asyncio.to_thread(heavy_calculation)
или использовать multiprocessing, если задача реально CPU-bound и нужно задействовать несколько ядер.
Общая схема #
async def
↓
coroutine object
↓
await / create_task()
↓
Task
↓
event loop
↓
выполнение до await
↓
ожидание Future / I/O
↓
другая Task выполняется
↓
I/O завершилось
↓
первая Task продолжается
Коротко по сути #
Асинхронная модель выполнения в Python устроена так:
1. async def создаёт корутину
2. корутина запускается через await или Task
3. event loop управляет задачами
4. await приостанавливает текущую задачу
5. пока одна задача ждёт I/O, выполняются другие
6. после завершения ожидания задача продолжается
Главное отличие от потоков: переключение происходит не принудительно ОС, а в контролируемых местах — на await. Поэтому это кооперативная конкурентность, а не автоматический параллелизм.
12. Что обозначают ключевые слова async и await в механизме корутин? #
async/await - это ключевые слова, представленные в Python 3.5 для определения асинхронных функций и вызовов. Они позволяют объявлять асинхронные функции и делать вызовы асинхронных функций в местах, где обычно используется блокирующий вызов.
Кратко:
- async — объявляет асинхронную функцию (корутину)
- await — приостанавливает выполнение корутины до завершения асинхронной операции
- event loop — управляет выполнением асинхронных задач
Корутина — это специальный тип функции, которая может приостанавливать свое выполнение и возобновлять его позже.
В Python корутины создаются с помощью ключевого слова async def
async def my_coroutine():
print("Начало корутины")
await asyncio.sleep(1)
print("Корутина завершена")
Особенности корутин:
- При вызове возвращают объект корутины, а не результат
- Для выполнения требуют наличия event loop
- Могут содержать выражения await
- Не выполняются до тех пор, пока не будут запланированы в event loop
Event Loop (Цикл событий) — это ядро асинхронного программирования в Python.
Он отвечает за:
- Выполнение корутин и обратных вызовов
- Выполнение сетевых операций ввода-вывода
- Запуск подпроцессов
# Создание и управление event loop вручную
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
try:
loop.run_until_complete(my_coroutine())
finally:
loop.close()
Ключевое слово Await
Выражение await приостанавливает выполнение корутины до тех пор, пока ожидаемый объект (awaitable) не вернет результат. Во время этой паузы event loop может выполнять другие задачи.
Что можно использовать с await:
- Другие корутины
- Задачи (Tasks)
- Футуры (Futures)
- Объекты, реализующие метод await
🚀 Жизненный цикл асинхронной операции:
- Создание корутины — функция объявлена как async def
- Вызов корутины — создается объект корутины, но выполнение не начинается
- Планирование в event loop — корутина помещается в очередь выполнения
- Выполнение до первого await — выполняется код до первого выражения await
- Приостановка и переключение — при встрече await корутина приостанавливается
- Возобновление выполнения — когда ожидаемая операция завершена
- Завершение — возврат результата или исключение
Пример детального выполнения:
async def complex_operation():
print("Шаг 1: Начало")
await asyncio.sleep(1) # Приостановка здесь
print("Шаг 2: После сна")
result = await fetch_data() # Еще одна приостановка
print("Шаг 3: Получены данные")
return result
Объявление асинхронных функций:
async def my_async_function():
# Асинхронная функция (корутина)
return "Результат"
Вызов асинхронных функций:
async def main():
# Правильно: используем await
result = await my_async_function()
# Неправильно: my_async_function() без await вернет корутину, а не результат
wrong_result = my_async_function() # Это объект корутины!
Запуск асинхронного кода:
# Python 3.7+
asyncio.run(main())
# Старые версии Python
loop = asyncio.get_event_loop()
loop.run_until_complete(main())
13. Какими способами реализовать конкурентность без async/await и как на это влияет GIL? | Можно ли реализовать асинхронный подход без asyncio? #
Да, конкурентность можно реализовать без async/await.
Основные варианты:
Конкурентность без async/await
├── threading
├── concurrent.futures.ThreadPoolExecutor
├── multiprocessing / ProcessPoolExecutor
├── InterpreterPoolExecutor
├── selectors / select / epoll / kqueue
├── callback-based event loop
├── generators / yield-based scheduler
└── greenlet / gevent
asyncio — это не единственный способ писать конкурентный код. Это стандартная библиотека для асинхронного I/O с async/await, но сама идея конкурентности шире.
1. Потоки: threading #
Самый прямой способ без async/await — создать несколько потоков:
import threading
import requests
def load_url(url):
response = requests.get(url)
print(url, response.status_code)
threads = [
threading.Thread(target=load_url, args=("https://example.com",)),
threading.Thread(target=load_url, args=("https://python.org",)),
]
for thread in threads:
thread.start()
for thread in threads:
thread.join()
Потоки хорошо подходят для I/O-bound задач: сетевые запросы, работа с файлами, ожидание БД, API, Redis и т.д. Документация Python прямо указывает, что threading полезен для задач, где много ожидания внешних ресурсов.
Но для CPU-bound задач в обычном CPython потоки часто не дают настоящего параллельного выполнения Python-кода из-за GIL.
2. ThreadPoolExecutor #
Более удобный способ — не создавать потоки вручную, а использовать пул потоков:
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests
def load_url(url):
response = requests.get(url)
return url, response.status_code
urls = [
"https://example.com",
"https://python.org",
"https://docs.python.org",
]
with ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(load_url, url) for url in urls]
for future in as_completed(futures):
print(future.result())
Здесь нет async def, await, asyncio, но выполнение всё равно асинхронное в широком смысле: ты отправляешь задачи в пул и забираешь результат позже через Future. В документации concurrent.futures прямо сказано, что модуль предоставляет высокоуровневый интерфейс для асинхронного выполнения callable-объектов через потоки, интерпретаторы или процессы.
3. multiprocessing / ProcessPoolExecutor #
Для CPU-bound задач лучше использовать процессы:
from concurrent.futures import ProcessPoolExecutor
def heavy_calculation(n):
total = 0
for i in range(n):
total += i * i
return total
with ProcessPoolExecutor() as executor:
results = executor.map(heavy_calculation, [10_000_000, 20_000_000, 30_000_000])
print(list(results))
Процессы обходят GIL, потому что каждый процесс имеет свой интерпретатор Python и свою память. Документация multiprocessing прямо говорит, что модуль обходит GIL за счёт subprocesses и позволяет использовать несколько процессоров.
Минусы процессов:
+ настоящая параллельность для CPU-bound
- дороже создание процесса
- данные нужно сериализовать
- память не общая, как у потоков
- обмен данными сложнее
4. InterpreterPoolExecutor #
В Python 3.14 появился InterpreterPoolExecutor: это пул потоков, но каждый worker работает в своём интерпретаторе. У каждого такого интерпретатора свой GIL, поэтому код может выполняться на разных CPU-ядрах параллельно.
Упрощённо:
ThreadPoolExecutor
├── один процесс
├── один интерпретатор
└── один GIL
InterpreterPoolExecutor
├── один процесс
├── несколько интерпретаторов
└── у каждого интерпретатора свой GIL
Но есть ограничение: интерпретаторы изолированы. Нельзя просто так разделять между ними изменяемые объекты. Данные нужно копировать, сериализовать или передавать через специальные механизмы.
5. selectors / select: свой event loop без asyncio #
Можно реализовать асинхронный подход вообще без asyncio, используя неблокирующие сокеты и selectors.
Идея:
1. Сокеты переводятся в non-blocking режим
2. selectors ждёт готовности сокета к чтению/записи
3. Когда сокет готов — вызывается обработчик
4. Один поток обслуживает много соединений
selectors — это высокоуровневая обёртка над select, poll, epoll, kqueue и другими механизмами ожидания I/O-событий. Документация рекомендует использовать selectors, если не нужен точный контроль над низкоуровневыми OS-примитивами.
Схематично:
import selectors
import socket
selector = selectors.DefaultSelector()
server = socket.socket()
server.bind(("localhost", 8000))
server.listen()
server.setblocking(False)
selector.register(server, selectors.EVENT_READ)
while True:
events = selector.select()
for key, mask in events:
sock = key.fileobj
if sock is server:
conn, addr = server.accept()
conn.setblocking(False)
selector.register(conn, selectors.EVENT_READ)
else:
data = sock.recv(1024)
if data:
sock.sendall(data)
else:
selector.unregister(sock)
sock.close()
Это уже асинхронный I/O-подход, но без asyncio и без async/await.
6. Генераторы и yield-based scheduler #
До появления современного async/await корутины в Python можно было строить через генераторы. yield приостанавливает выполнение функции и сохраняет её состояние, а потом выполнение можно продолжить.
Простейшая идея:
def task(name):
for i in range(3):
print(name, i)
yield
tasks = [task("A"), task("B")]
while tasks:
current = tasks.pop(0)
try:
next(current)
tasks.append(current)
except StopIteration:
pass
Вывод будет чередоваться:
A 0
B 0
A 1
B 1
A 2
B 2
Это кооперативная конкурентность:
Задача сама отдаёт управление через yield.
Пока она не отдаст управление — другая задача не выполняется.
Так можно написать собственный мини-планировщик. Но для реального сетевого I/O нужно связать генераторы с selectors, таймерами, очередями событий и обработкой ошибок.
7. greenlet / gevent #
greenlet — это легковесные корутины для конкурентного программирования внутри процесса. Они похожи на кооперативные потоки: переключение происходит явно или через библиотеку, а не вытесняется ОС. (
Greenlet)
gevent строится поверх greenlet и даёт синхронный на вид API поверх event loop libev или libuv.
Примерно так:
import gevent
from gevent import monkey
monkey.patch_all()
import requests
def load_url(url):
response = requests.get(url)
print(url, response.status_code)
jobs = [
gevent.spawn(load_url, "https://example.com"),
gevent.spawn(load_url, "https://python.org"),
]
gevent.joinall(jobs)
Код выглядит синхронным, но I/O может выполняться конкурентно. Минус: если где-то используется блокирующий код, который не был “пропатчен” или не умеет отдавать управление, он может заблокировать весь поток.
Как на всё это влияет GIL #
GIL — это глобальная блокировка интерпретатора CPython. В обычной сборке CPython поток должен владеть GIL, чтобы работать с Python-объектами и выполнять Python C API.
Главное:
I/O-bound + threads
→ обычно нормально
CPU-bound + threads
→ плохо масштабируется из-за GIL
CPU-bound + processes
→ хорошо, потому что процессы обходят GIL
CPU-bound + C extensions
→ может быть хорошо, если расширение отпускает GIL
asyncio / selectors / gevent
→ GIL не основная проблема, потому что одна задача ждёт I/O, пока другая работает
Документация Python отдельно указывает, что GIL ограничивает выигрыш от threading для CPU-bound задач, потому что только один поток может выполнять Python bytecode одновременно.
При этом GIL освобождается вокруг блокирующих I/O-операций, например чтения или записи файла. Поэтому потоки всё равно полезны для сетевых запросов, файлов, БД и других операций ожидания.
Важный момент: GIL не отменяет race condition #
Иногда ошибочно думают:
Есть GIL → значит потоки безопасны
Это неверно.
GIL не делает твою бизнес-логику потокобезопасной. Интерпретатор может переключать потоки между bytecode-инструкциями, поэтому для общего изменяемого состояния всё равно нужны Lock, RLock, Queue, Condition, Semaphore и другие примитивы синхронизации. Документация Python прямо отмечает, что переключения между bytecode-инструкциями — причина, по которой lock-и всё равно нужны для thread-safety в pure-Python коде.
Пример проблемы:
counter = 0
def increment():
global counter
for _ in range(100_000):
counter += 1
Операция counter += 1 не должна восприниматься как одна неделимая операция на уровне логики.
Можно ли реализовать async-подход без asyncio? #
Да.
Есть несколько вариантов:
1. concurrent.futures
Асинхронно отправлять задачи в пул и получать Future.
2. selectors
Написать свой event loop на неблокирующих сокетах.
3. callbacks
Строить код вокруг функций-обработчиков событий.
4. generators
Сделать кооперативный scheduler через yield.
5. greenlet/gevent
Использовать greenlet-based конкурентность.
6. Trio / Curio / AnyIO
Это async-библиотеки не на asyncio или не только на asyncio,
но обычно они всё равно используют async/await.
То есть:
async/await — синтаксис языка
asyncio — стандартная библиотека/event loop
асинхронность — общий подход
конкурентность — ещё более широкое понятие
asyncio — не единственный способ. В документации asyncio event loop описан как ядро asyncio-приложения, которое запускает async tasks, callbacks, сетевой I/O и subprocesses. Но подобный event loop можно построить и самостоятельно на selectors, либо использовать сторонний loop вроде gevent.
Практический выбор #
Много HTTP/API/БД запросов, код уже синхронный
→ ThreadPoolExecutor
CPU-bound вычисления
→ ProcessPoolExecutor / multiprocessing
Нужно много соединений в одном потоке
→ asyncio / selectors / gevent
Нужен низкоуровневый контроль
→ selectors
Нужно писать современный async-код
→ asyncio / Trio / AnyIO
Нужна настоящая параллельность внутри одного процесса на новых версиях Python
→ InterpreterPoolExecutor, но с учётом изоляции интерпретаторов
Итог:
Конкурентность без async/await — да.
Асинхронный подход без asyncio — да.
Настоящий параллелизм CPU-bound в обычных потоках CPython — обычно нет из-за GIL.
Для CPU-bound лучше процессы, отдельные интерпретаторы или free-threaded Python.
Начиная с Python 3.13, CPython поддерживает отдельную free-threaded сборку, где GIL отключён. Но это не обычный режим по умолчанию для большинства окружений, и часть C-расширений может быть не готова к такому режиму или повторно включать GIL.
14. Какие примитивы синхронизации используются в многопоточном программировании? #
Основные примитивы синхронизации #
В многопоточном программировании примитивы синхронизации нужны, чтобы несколько потоков безопасно работали с общими данными: не портили состояние, не читали данные в промежуточном состоянии и не выполняли критические участки одновременно.
В Python основные примитивы находятся в модуле threading: Lock, RLock, Condition, Semaphore, BoundedSemaphore, Event, Barrier. Часть из них можно использовать через with, чтобы гарантированно освобождать ресурс после выхода из блока.
1. Lock #
Lock — обычная блокировка, или mutex.
Используется, когда только один поток должен иметь доступ к критической секции.
import threading
counter = 0
lock = threading.Lock()
def increment():
global counter
for _ in range(100_000):
with lock:
counter += 1
Смысл:
Поток A вошёл в with lock
→ остальные потоки ждут
→ поток A вышел
→ lock освободился
→ другой поток может войти
Типичные случаи:
- изменение общего счётчика
- запись в общий список/словарь
- работа с общим файлом
- изменение общего состояния объекта
Важно: GIL не заменяет Lock. В CPython GIL ограничивает одновременное выполнение Python bytecode, но не делает бизнес-операции потокобезопасными. Документация threading прямо указывает, что из-за GIL только один поток может выполнять Python-код одновременно, но для корректной синхронизации общего состояния всё равно нужны примитивы синхронизации.
2. RLock #
RLock — reentrant lock, то есть повторно-входимая блокировка.
Обычный Lock нельзя захватить повторно из того же потока. RLock — можно.
import threading
lock = threading.RLock()
def outer():
with lock:
inner()
def inner():
with lock:
print("inner")
С обычным Lock такой код мог бы зависнуть:
outer() захватил lock
→ вызвал inner()
→ inner() снова пытается захватить тот же lock
→ поток ждёт сам себя
RLock полезен, когда несколько методов одного объекта используют одну и ту же блокировку и могут вызывать друг друга.
3. Semaphore #
Semaphore ограничивает количество потоков, которые одновременно могут пройти в участок кода.
Например, максимум 3 потока могут одновременно обращаться к внешнему API:
import threading
import time
semaphore = threading.Semaphore(3)
def request_api(user_id):
with semaphore:
print(f"start {user_id}")
time.sleep(1)
print(f"end {user_id}")
Смысл:
Lock → пускает 1 поток
Semaphore → пускает N потоков
Типичные случаи:
- ограничить количество одновременных запросов
- ограничить доступ к пулу соединений
- ограничить работу с тяжёлым ресурсом
4. BoundedSemaphore #
BoundedSemaphore похож на Semaphore, но дополнительно защищает от ошибки лишнего release().
Обычный Semaphore можно случайно “освободить” больше раз, чем захватили. BoundedSemaphore обнаруживает такую ошибку.
import threading
sem = threading.BoundedSemaphore(2)
sem.acquire()
sem.release()
# sem.release() # ошибка: release больше, чем acquire
Используется там, где важно не превысить исходный лимит ресурса.
5. Event #
Event — флаг-сигнал между потоками.
Один поток может ждать события, другой — установить его.
import threading
import time
event = threading.Event()
def worker():
print("worker ждёт сигнал")
event.wait()
print("worker начал работу")
thread = threading.Thread(target=worker)
thread.start()
time.sleep(2)
event.set()
Смысл:
event.clear() → события нет
event.wait() → поток ждёт
event.set() → событие произошло, ожидающие потоки продолжают работу
Типичные случаи:
- сигнал завершения
- сигнал старта
- ожидание готовности ресурса
- остановка фонового потока
Пример остановки фонового потока:
import threading
import time
stop_event = threading.Event()
def worker():
while not stop_event.is_set():
print("working")
time.sleep(1)
thread = threading.Thread(target=worker)
thread.start()
time.sleep(3)
stop_event.set()
thread.join()
6. Condition #
Condition используется, когда поток должен ждать не просто lock, а конкретное условие.
Например:
consumer ждёт, пока появятся данные
producer добавляет данные
producer уведомляет consumer
Пример:
import threading
condition = threading.Condition()
items = []
def producer():
with condition:
items.append("data")
condition.notify()
def consumer():
with condition:
while not items:
condition.wait()
item = items.pop()
print(item)
Важный момент:
while not items:
condition.wait()
Обычно используют именно while, а не if, потому что поток после пробуждения должен заново проверить условие.
Condition часто лежит внутри более высокоуровневых структур, например очередей.
7. Barrier #
Barrier заставляет несколько потоков дождаться друг друга в одной точке.
import threading
import time
barrier = threading.Barrier(3)
def worker(name):
print(f"{name}: подготовка")
time.sleep(1)
print(f"{name}: ждёт остальных")
barrier.wait()
print(f"{name}: продолжает работу")
for i in range(3):
threading.Thread(target=worker, args=(f"thread-{i}",)).start()
Смысл:
Поток 1 дошёл до barrier.wait() → ждёт
Поток 2 дошёл до barrier.wait() → ждёт
Поток 3 дошёл до barrier.wait() → все продолжают
Типичные случаи:
- синхронный старт нескольких потоков
- этапная обработка данных
- тестирование конкурентного кода
8. Queue #
queue.Queue — потокобезопасная очередь.
Это один из самых практичных способов обмена данными между потоками. Модуль queue реализует очереди для схемы multi-producer / multi-consumer и уже содержит нужную locking-семантику.
Пример producer-consumer:
import queue
import threading
q = queue.Queue()
def producer():
for i in range(5):
q.put(i)
def consumer():
while True:
item = q.get()
try:
print("process", item)
finally:
q.task_done()
threading.Thread(target=producer).start()
threading.Thread(target=consumer, daemon=True).start()
q.join()
Смысл:
producer кладёт задачи в очередь
consumer забирает задачи из очереди
Queue сама синхронизирует доступ между потоками
Варианты:
queue.Queue → FIFO
queue.LifoQueue → LIFO, как стек
queue.PriorityQueue → очередь с приоритетом
queue.SimpleQueue → простая неограниченная очередь
9. Future #
Future — не совсем низкоуровневый примитив синхронизации, но часто используется в многопоточном коде.
Он представляет результат задачи, который появится позже.
from concurrent.futures import ThreadPoolExecutor
def calculate():
return 10 + 20
with ThreadPoolExecutor() as executor:
future = executor.submit(calculate)
result = future.result()
print(result)
future.result() блокирует текущий поток до завершения задачи.
Модуль concurrent.futures предоставляет высокоуровневый интерфейс для асинхронного выполнения функций через потоки, процессы или отдельные интерпретаторы.
Что использовать чаще всего #
На практике чаще всего нужны:
Lock
→ защитить общий mutable state
RLock
→ если один и тот же поток может повторно входить в защищённый код
Semaphore
→ ограничить количество одновременных операций
Event
→ дать сигнал потоку
Condition
→ ждать конкретного состояния
Barrier
→ дождаться группы потоков
Queue
→ безопасно передавать задачи/данные между потоками
Краткая таблица #
| Примитив | Для чего нужен |
|---|---|
Lock | Один поток в критической секции |
RLock | Повторный захват lock тем же потоком |
Semaphore | Ограничение количества одновременных потоков |
BoundedSemaphore | Semaphore с защитой от лишнего release() |
Event | Сигнал между потоками |
Condition | Ожидание конкретного условия |
Barrier | Ожидание группы потоков |
Queue | Потокобезопасный обмен данными |
Future | Получение результата задачи позже |
Главное #
GIL ≠ синхронизация бизнес-логики
Даже в CPython с GIL нужно использовать Lock, Queue, Event, Condition и другие примитивы, если несколько потоков работают с общим изменяемым состоянием. Самый безопасный подход — минимизировать общее состояние и передавать данные между потоками через queue.Queue.
15. Какая стандартная библиотека отвечает за асинхронность (async/await) и событийный цикл? #
Стандартная библиотека — asyncio
#
За асинхронность в стиле async/await и событийный цикл в Python отвечает стандартный модуль:
import asyncio
asyncio предоставляет инфраструктуру для:
- корутин
- задач Task
- Future
- event loop
- неблокирующего сетевого I/O
- таймеров
- async-примитивов синхронизации
- запуска subprocess
Официальная документация описывает asyncio как библиотеку для написания конкурентного кода с помощью синтаксиса async/await. Также в ней есть низкоуровневые API для создания и управления event loop.
Минимальный пример #
import asyncio
async def main():
print("start")
await asyncio.sleep(1)
print("end")
asyncio.run(main())
Что здесь происходит:
async def main()
→ создаёт coroutine function
main()
→ создаёт coroutine object
asyncio.run(main())
→ создаёт event loop
→ запускает корутину
→ закрывает event loop после завершения
asyncio.run() — основной высокоуровневый способ запустить async-программу. В документации указано, что он создаёт event loop, запускает корутину и закрывает цикл после выполнения.
Важно разделять #
async/await
→ синтаксис языка Python
asyncio
→ стандартная библиотека для работы с этим синтаксисом
event loop
→ механизм, который планирует и выполняет async-задачи
Сам async def — это часть синтаксиса Python. Функции, объявленные через async def, всегда являются coroutine functions.
А asyncio — библиотека, которая умеет эти корутины запускать, планировать, превращать в Task, ждать Future, обслуживать сетевой I/O и управлять событийным циклом.
16. Какие ключевые особенности asyncio как модели конкурентности (event loop, корутины, задачи)? #
asyncio как модель конкурентности
#
asyncio — это модель кооперативной конкурентности: несколько задач могут продвигаться “как будто одновременно”, но в рамках одного event loop в один момент времени выполняется только одна задача Python-кода. Переключение происходит не принудительно, как у потоков ОС, а в местах, где корутина сама отдаёт управление через await. Официальная документация описывает asyncio как библиотеку для конкурентного кода с использованием async/await.
asyncio
├── event loop
├── coroutine
├── Task
├── Future
├── await
└── неблокирующий I/O
1. Event loop #
event loop — центральный механизм asyncio.
Он:
- запускает корутины
- планирует задачи
- выполняет callbacks
- следит за сетевым I/O
- обрабатывает таймеры
- запускает subprocess
- возобновляет задачи, когда их ожидание завершилось
Документация прямо называет event loop ядром каждого asyncio-приложения: он выполняет асинхронные задачи и callbacks, сетевые операции и subprocesses.
Упрощённо:
event loop
→ взял готовую задачу
→ выполняет её до await
→ задача ждёт I/O/таймер/Future
→ event loop переключается на другую готовую задачу
→ когда ожидание завершилось, задача продолжается
Пример:
import asyncio
async def main():
print("start")
await asyncio.sleep(1)
print("end")
asyncio.run(main())
asyncio.run() запускает awaitable-объект в event loop, управляет циклом, финализирует async generators и закрывает executor после завершения.
2. Корутины #
Корутина — это функция, объявленная через async def.
async def fetch_data():
await asyncio.sleep(1)
return "data"
Важный момент:
coro = fetch_data()
Этот вызов не запускает код функции сразу. Он создаёт coroutine object. Чтобы корутина начала выполняться, её нужно либо await-нуть, либо обернуть в Task.
async def main():
result = await fetch_data()
print(result)
await означает:
эта корутина сейчас ждёт результат
event loop может пока выполнить другие задачи
потом выполнение этой корутины продолжится
То есть await — это точка добровольного переключения.
3. Task #
Task — это запланированная корутина.
Просто создать coroutine object недостаточно:
coro = fetch_data()
А вот так корутина планируется на выполнение:
task = asyncio.create_task(fetch_data())
Документация указывает, что Task используется для конкурентного планирования корутин; когда корутина оборачивается в Task через asyncio.create_task(), она автоматически планируется для выполнения.
Пример:
import asyncio
async def load(name, delay):
print(f"{name}: start")
await asyncio.sleep(delay)
print(f"{name}: end")
return name
async def main():
task1 = asyncio.create_task(load("A", 2))
task2 = asyncio.create_task(load("B", 1))
result1 = await task1
result2 = await task2
print(result1, result2)
asyncio.run(main())
Здесь A и B выполняются конкурентно:
A стартует
B стартует
B заканчивает ожидание раньше
A заканчивает позже
Но это не значит, что Python-код A и B выполнялся параллельно на разных ядрах. Они по очереди выполнялись в одном event loop.
4. Future #
Future — низкоуровневый объект, который представляет результат, который появится позже.
Future
├── пока результата нет
├── потом будет result
├── или exception
└── или cancellation
Документация описывает asyncio.Future как awaitable-объект, представляющий eventual result асинхронной операции. Корутины могут ждать Future до появления результата, исключения или отмены.
Обычно напрямую Future в пользовательском коде создают редко. Чаще ты работаешь с корутинами и Task, а Future остаётся внутри библиотек, event loop и низкоуровневого I/O.
5. Кооперативное переключение #
Главная особенность asyncio — задачи переключаются только в точках ожидания.
Например:
async def good():
await asyncio.sleep(1)
Здесь задача отдаёт управление event loop.
А вот так плохо:
async def bad():
while True:
pass
Эта корутина не содержит await, поэтому она заблокирует event loop.
Правильная модель мышления:
asyncio не прерывает задачу насильно
задача должна сама дойти до await
только тогда event loop сможет выполнить что-то другое
6. Неблокирующий I/O #
asyncio особенно полезен для I/O-bound задач:
- HTTP-запросы
- WebSocket
- работа с БД через async-драйвер
- сетевые серверы
- ожидание таймеров
- большое количество соединений
Смысл:
пока одна задача ждёт ответ от сети
event loop выполняет другую задачу
Это отличается от синхронного кода:
time.sleep(1)
time.sleep() блокирует поток.
А это не блокирует event loop:
await asyncio.sleep(1)
7. asyncio — не про CPU-bound параллелизм
#
asyncio хорошо работает, когда задачи часто ждут I/O.
Но если задача активно грузит CPU:
async def cpu_bound():
total = 0
for i in range(100_000_000):
total += i
return total
Она будет занимать event loop и мешать другим задачам.
Для CPU-bound задач обычно используют:
- multiprocessing
- ProcessPoolExecutor
- run_in_executor
- asyncio.to_thread для выноса блокирующей функции в поток
Но сам asyncio не превращает CPU-bound код в параллельный код.
8. Один event loop — обычно один поток #
Обычно event loop работает в одном потоке.
один event loop
→ один поток
→ много задач
→ переключение через await
Поэтому в рамках одного loop не нужно думать о race condition так же, как в многопоточном коде, но проблемы общего состояния всё равно возможны, если несколько корутин изменяют одни и те же данные между await.
Пример потенциальной проблемы:
counter = 0
async def increment():
global counter
value = counter
await asyncio.sleep(0)
counter = value + 1
Здесь между чтением и записью есть await, значит другая корутина может вмешаться.
Для таких случаев есть async-примитивы:
lock = asyncio.Lock()
async with lock:
# критическая секция
...
9. Task cancellation #
Задачи в asyncio можно отменять.
task = asyncio.create_task(fetch_data())
task.cancel()
Отмена обычно приводит к выбрасыванию asyncio.CancelledError внутри корутины. Это важно для таймаутов, graceful shutdown и отмены долгих операций.
Типичный шаблон:
async def worker():
try:
while True:
await asyncio.sleep(1)
except asyncio.CancelledError:
# очистка ресурсов
raise
10. asyncio.gather() и конкурентный запуск
#
Для запуска нескольких awaitable-объектов часто используют asyncio.gather():
import asyncio
async def load_user():
await asyncio.sleep(1)
return "user"
async def load_orders():
await asyncio.sleep(1)
return "orders"
async def main():
user, orders = await asyncio.gather(
load_user(),
load_orders(),
)
print(user, orders)
asyncio.run(main())
Без конкурентности это заняло бы примерно 2 секунды:
load_user → 1 секунда
load_orders → 1 секунда
итого → 2 секунды
С gather() обе операции ожидания идут конкурентно:
load_user и load_orders ждут одновременно
итого примерно 1 секунда
Главная схема #
async def
→ создаёт корутинную функцию
coroutine object
→ результат вызова async-функции
Task
→ корутина, запланированная в event loop
await
→ точка ожидания и передачи управления event loop
event loop
→ планировщик, который переключает задачи
Future
→ объект будущего результата
Главное отличие от потоков #
threading
→ переключение контролирует ОС
→ вытесняющая многозадачность
→ возможны параллельные потоки, но в CPython мешает GIL
asyncio
→ переключение происходит через await
→ кооперативная многозадачность
→ один event loop обычно выполняет одну задачу Python-кода за раз
Итог:
asyncio — это модель конкурентности для I/O-bound задач.
event loop — планировщик.
coroutine — приостанавливаемая функция.
Task — корутина, поставленная на выполнение.
await — точка, где задача отдаёт управление.
Future — низкоуровневый контейнер будущего результата.
Использовать asyncio целесообразно, когда программа много ждёт внешние ресурсы и нужно эффективно обслуживать много таких ожиданий в одном потоке.
17. В какие моменты event loop переключается между задачами в асинхронном выполнении? #
Когда event loop переключается между задачами #
В asyncio переключение задач происходит не принудительно, а кооперативно: задача сама отдаёт управление event loop, обычно на await. Event loop в один момент выполняет только одну задачу; пока одна задача ждёт завершения Future/I/O/таймера, loop может выполнять другую задачу.
Главные моменты переключения:
1. На await, если ожидаемый объект ещё не готов
#
Например:
async def task_1():
print("A")
await some_io()
print("B")
Когда выполнение доходит до:
await some_io()
если some_io() ещё не завершился, текущая задача приостанавливается, а event loop берёт другую готовую задачу.
Схематично:
task_1 выполняется
↓
доходит до await some_io()
↓
some_io ещё не готов
↓
task_1 приостанавливается
↓
event loop запускает другую готовую task
2. На await asyncio.sleep(...)
#
asyncio.sleep() всегда приостанавливает текущую задачу и даёт возможность выполниться другим задачам. Даже await asyncio.sleep(0) используется как явная добровольная передача управления event loop.
import asyncio
async def worker(name):
for i in range(3):
print(name, i)
await asyncio.sleep(0)
asyncio.run(asyncio.gather(
worker("A"),
worker("B"),
))
Логика:
A 0
await sleep(0) → отдал управление
B 0
await sleep(0) → отдал управление
A 1
...
3. Когда задача ждёт сетевой ввод/вывод #
Типичный пример:
data = await reader.read(1024)
Пока данные из сокета не пришли, задача не занимает поток выполнения. Event loop следит за готовностью I/O и в это время выполняет другие задачи. asyncio как раз хорошо подходит для I/O-bound задач и сетевого кода.
4. Когда задача ждёт другую задачу #
result = await other_task
Если other_task ещё не завершилась, текущая задача приостанавливается.
task_1 ждёт task_2
↓
task_1 suspended
↓
event loop выполняет task_2 или другие готовые задачи
↓
task_2 завершилась
↓
task_1 продолжает выполнение
5. Когда задача ждёт Future #
Внутри asyncio многие операции сводятся к ожиданию Future.
result = await future
Если future ещё не готов, задача отдаёт управление. Когда future получает результат или исключение, event loop снова ставит задачу в очередь готовых к продолжению.
Важный момент: обычный код не переключается #
Вот здесь переключения не будет:
async def bad_task():
while True:
pass
Несмотря на async def, внутри нет await. Такая задача заблокирует event loop.
То же самое с тяжёлым CPU-кодом:
async def cpu_task():
for i in range(10**9):
...
Пока код не дойдёт до await, event loop не сможет переключиться на другие задачи.
Правильная короткая формулировка #
Event loop переключается между задачами только в точках ожидания, когда текущая coroutine делает await и ожидаемый объект не готов: I/O, sleep, Future, другая Task, блокировка asyncio.Lock, очередь asyncio.Queue и т.д.
asyncio = кооперативная многозадачность
не event loop насильно прерывает задачу,
а задача сама отдаёт управление через await
Мини-пример #
import asyncio
async def task_a():
print("A1")
await asyncio.sleep(1)
print("A2")
async def task_b():
print("B1")
await asyncio.sleep(1)
print("B2")
async def main():
await asyncio.gather(task_a(), task_b())
asyncio.run(main())
Примерный порядок:
A1
task_a дошла до await sleep(1) → отдала управление
B1
task_b дошла до await sleep(1) → отдала управление
через ~1 секунду обе задачи снова готовы
A2
B2
Итог: asyncio не делает параллельное выполнение Python-кода в одном потоке. Оно быстро переключает задачи в моменты, когда одна из них ждёт внешнее событие.
18. Как организовать одновременное выполнение нескольких coroutine-задач? #
Одновременное выполнение нескольких coroutine-задач в Python обычно организуют через asyncio:
import asyncio
async def fetch_user(user_id: int):
await asyncio.sleep(1)
return f"user:{user_id}"
async def main():
results = await asyncio.gather(
fetch_user(1),
fetch_user(2),
fetch_user(3),
)
print(results)
asyncio.run(main())
asyncio.run() запускает корутину как основную точку входа программы, создаёт event loop, выполняет awaitable-объект и закрывает цикл после завершения.
Важно: просто await подряд — не параллельное выполнение
#
Так задачи выполняются последовательно:
async def main():
result_1 = await fetch_user(1)
result_2 = await fetch_user(2)
result_3 = await fetch_user(3)
Сначала полностью ожидается fetch_user(1), потом fetch_user(2), потом fetch_user(3).
Для конкурентного выполнения нужно передать корутины в механизм планирования задач: asyncio.gather(), asyncio.TaskGroup или asyncio.create_task().
Способ 1: asyncio.gather()
#
Используется, когда нужно запустить несколько независимых задач и дождаться всех результатов.
import asyncio
async def task(name: str, delay: int):
print(f"{name}: start")
await asyncio.sleep(delay)
print(f"{name}: end")
return name
async def main():
results = await asyncio.gather(
task("A", 3),
task("B", 1),
task("C", 2),
)
print(results)
asyncio.run(main())
Примерный вывод:
A: start
B: start
C: start
B: end
C: end
A: end
['A', 'B', 'C']
asyncio.gather() запускает awaitable-объекты конкурентно; если переданы корутины, они автоматически планируются как задачи.
Способ 2: asyncio.create_task()
#
Используется, когда задачу нужно запустить сейчас, а дождаться результата позже.
import asyncio
async def load_data():
await asyncio.sleep(2)
return "data"
async def send_log():
await asyncio.sleep(1)
return "log sent"
async def main():
data_task = asyncio.create_task(load_data())
log_task = asyncio.create_task(send_log())
print("Задачи уже запущены")
data = await data_task
log = await log_task
print(data)
print(log)
asyncio.run(main())
Смысл:
create_task()
↓
корутина оборачивается в Task
↓
Task регистрируется в event loop
↓
event loop выполняет задачи конкурентно
Задача Task — это объект, через который event loop может управлять выполнением корутины. Документация Python отдельно выделяет coroutines, awaitables, tasks, cancellation и task groups как основные элементы работы с конкурентными корутинами.
Способ 3: asyncio.TaskGroup
#
В современном коде часто предпочтительнее TaskGroup, особенно когда задачи логически связаны.
import asyncio
async def worker(name: str, delay: int):
await asyncio.sleep(delay)
print(f"{name} done")
return name
async def main():
async with asyncio.TaskGroup() as tg:
task_1 = tg.create_task(worker("A", 2))
task_2 = tg.create_task(worker("B", 1))
task_3 = tg.create_task(worker("C", 3))
print(task_1.result())
print(task_2.result())
print(task_3.result())
asyncio.run(main())
TaskGroup даёт structured concurrency: задачи создаются внутри блока, а выход из блока означает, что все задачи завершены или корректно обработаны при ошибке.
Ограничение количества одновременных задач #
Например, нельзя отправлять 1000 HTTP-запросов сразу. Тогда используют asyncio.Semaphore.
import asyncio
sem = asyncio.Semaphore(3)
async def limited_task(number: int):
async with sem:
print(f"start {number}")
await asyncio.sleep(1)
print(f"end {number}")
async def main():
tasks = [
limited_task(i)
for i in range(10)
]
await asyncio.gather(*tasks)
asyncio.run(main())
Здесь одновременно выполняется максимум 3 задачи.
Главное отличие #
await coro()
означает: выполнить и дождаться эту корутину в текущем месте.
asyncio.create_task(coro())
означает: запланировать корутину как отдельную задачу в event loop.
await asyncio.gather(coro1(), coro2(), coro3())
означает: запустить несколько корутин конкурентно и дождаться всех.
Когда это реально полезно #
asyncio хорошо подходит для I/O-bound задач: сетевые запросы, работа с базой данных через async-драйвер, Redis, файловые операции через async-совместимые библиотеки, ожидание внешних сервисов. Официальная документация описывает asyncio как библиотеку для конкурентного кода через async/await, особенно подходящую для I/O-bound и сетевого кода.
Для CPU-bound задач вроде тяжёлых вычислений asyncio сам по себе не даст настоящего параллельного ускорения, потому что event loop выполняет Python-код в одном потоке и переключается только в точках ожидания.
Итог #
Основные варианты:
# 1. Просто и часто достаточно
await asyncio.gather(coro1(), coro2(), coro3())
# 2. Запустить сейчас, дождаться позже
task = asyncio.create_task(coro())
result = await task
# 3. Современный структурированный вариант
async with asyncio.TaskGroup() as tg:
task = tg.create_task(coro())
Для большинства случаев:
asyncio.gather()
— самый простой вариант.
Для более аккуратного production-кода с группой связанных задач:
asyncio.TaskGroup
— более безопасный и структурированный вариант.
19. Как event loop понимает, что можно запустить другую задачу? #
event loop не “угадывает”, что пора переключиться.
Он запускает задачу, а задача сама отдаёт управление, когда доходит до await, который ждёт ещё неготовый результат: I/O, таймер, Future, другую Task и т.д.
То есть переключение в asyncio — кооперативное.
Что происходит пошагово #
Допустим:
async def task1():
print("A")
await asyncio.sleep(1)
print("B")
async def task2():
print("C")
Когда event loop запускает task1, она выполняется до первого реального ожидания:
print("A")
await asyncio.sleep(1)
asyncio.sleep(1) создаёт ожидание таймера. Результата ещё нет, поэтому task1 приостанавливается и отдаёт управление обратно в event loop.
После этого event loop смотрит:
Есть ли готовые задачи/колбэки?
Есть ли завершившиеся Future?
Есть ли сработавшие таймеры?
Есть ли готовые I/O-события?
Если есть готовая task2, он запускает её.
Официальная документация Python описывает event loop как ядро asyncio, которое выполняет асинхронные задачи и callbacks, а также обрабатывает сетевой I/O и subprocesses.
Главная идея #
Задача выполняется до тех пор, пока не произойдёт одно из событий:
1. задача завершилась
2. задача выбросила исключение
3. задача дошла до await и ждёт неготовый результат
В третьем случае задача говорит event loop примерно следующее:
Я пока не могу продолжить.
Разбуди меня, когда этот Future/Task/I/O/таймер будет готов.
Future в asyncio как раз представляет будущий результат операции. Корутинa может ждать Future, пока тот не получит результат, исключение или отмену.
Как задача возвращается обратно #
Когда ожидаемый объект становится готовым:
таймер истёк
сокет получил данные
запрос к БД вернул результат
другая task завершилась
Future получил set_result()
Future помечается как done, и его callbacks добавляются обратно в очередь event loop. Документация указывает, что callback у Future вызывается после завершения Future, а если Future уже завершён, callback планируется через loop.call_soon().
Упрощённо:
task1 дошла до await
↓
task1 приостановлена
↓
event loop запускает другие готовые задачи
↓
ожидание task1 завершилось
↓
task1 снова попадает в очередь готовых задач
↓
event loop продолжает task1 после await
Важно: event loop не прерывает обычный код #
Вот так переключения не будет:
async def bad():
for i in range(10_000_000_000):
pass
Несмотря на async def, внутри нет await. Значит, задача не отдаёт управление. Event loop будет заблокирован, пока цикл не закончится.
Правильнее:
async def better():
for i in range(10_000_000):
if i % 1000 == 0:
await asyncio.sleep(0)
await asyncio.sleep(0) — явная точка уступки управления. Она даёт event loop возможность запустить другие готовые задачи.
Ментальная модель #
event loop
├─ очередь готовых задач
├─ таймеры
├─ I/O события
└─ callbacks
Task выполняется
└─ до ближайшего await, который реально ждёт
await на неготовом объекте
└─ задача засыпает
Future/таймер/I/O завершился
└─ задача снова становится готовой
Главное правило #
event loop понимает, что можно запустить другую задачу, потому что текущая задача дошла до await и приостановилась на неготовом awaitable-объекте.
Не async def делает переключение, а именно реальная точка ожидания:
await something_not_ready
20. Как в asyncio запускать блокирующие операции, чтобы не останавливать event loop? (run_in_excecutor) #
Правильное имя метода: run_in_executor, не run_in_excecutor.
В asyncio блокирующую синхронную функцию нельзя вызывать напрямую внутри async def, иначе она остановит event loop. Её нужно вынести в отдельный поток или процесс:
result = await loop.run_in_executor(None, blocking_func)
или, в современном коде для I/O-bound задач:
result = await asyncio.to_thread(blocking_func)
asyncio.to_thread() запускает обычную синхронную функцию в отдельном потоке и возвращает корутину, которую можно await-ить. Документация прямо указывает, что это нужно для I/O-bound функций, которые иначе заблокировали бы event loop.
Проблема #
Допустим, есть блокирующая функция:
import time
def blocking_io():
time.sleep(3)
return "done"
Если вызвать её напрямую:
async def main():
result = blocking_io()
print(result)
то time.sleep(3) заблокирует весь поток, в котором работает event loop. В это время другие корутины не смогут выполняться.
Вариант 1: asyncio.to_thread #
Для обычных блокирующих I/O-операций чаще всего достаточно asyncio.to_thread:
import asyncio
import time
def blocking_io():
time.sleep(3)
return "done"
async def background_task():
while True:
print("event loop жив")
await asyncio.sleep(0.5)
async def main():
task = asyncio.create_task(background_task())
result = await asyncio.to_thread(blocking_io)
print(result)
task.cancel()
asyncio.run(main())
Смысл:
event loop
├─ выполняет async-код
├─ отдаёт blocking_io в отдельный thread
└─ пока blocking_io работает, продолжает выполнять другие coroutine
То есть блокирующая функция всё ещё блокирующая, но блокирует не event loop, а отдельный поток.
Вариант 2: loop.run_in_executor #
run_in_executor — более низкоуровневый способ. Он позволяет явно выбрать executor: потоковый или процессный. По документации, loop.run_in_executor(executor, func, *args) запускает func в указанном executor; если передать None, используется executor по умолчанию, обычно ThreadPoolExecutor.
import asyncio
import time
def blocking_io():
time.sleep(3)
return "done"
async def main():
loop = asyncio.get_running_loop()
result = await loop.run_in_executor(
None,
blocking_io,
)
print(result)
asyncio.run(main())
Здесь:
None
означает:
используй default ThreadPoolExecutor
С аргументами #
Если функция принимает аргументы:
def read_file(path: str):
with open(path, "r", encoding="utf-8") as file:
return file.read()
то с to_thread:
content = await asyncio.to_thread(read_file, "data.txt")
с run_in_executor:
loop = asyncio.get_running_loop()
content = await loop.run_in_executor(
None,
read_file,
"data.txt",
)
Для keyword-аргументов с run_in_executor обычно используют functools.partial, потому что run_in_executor напрямую передаёт позиционные аргументы. Документация также указывает functools.partial() как способ передавать keyword-аргументы.
from functools import partial
content = await loop.run_in_executor(
None,
partial(open_file, path="data.txt", encoding="utf-8"),
)
ThreadPoolExecutor вручную #
Можно создать свой пул потоков:
import asyncio
from concurrent.futures import ThreadPoolExecutor
import time
def blocking_io(n: int):
time.sleep(2)
return n * 2
async def main():
loop = asyncio.get_running_loop()
with ThreadPoolExecutor(max_workers=5) as pool:
tasks = [
loop.run_in_executor(pool, blocking_io, i)
for i in range(10)
]
results = await asyncio.gather(*tasks)
print(results)
asyncio.run(main())
Так удобно контролировать количество потоков:
max_workers=5
значит одновременно будет работать не больше 5 блокирующих функций.
Для CPU-bound лучше ProcessPoolExecutor #
Если операция не I/O-bound, а CPU-bound, например тяжёлые вычисления, ThreadPoolExecutor обычно не даст нормального параллелизма в CPython из-за GIL. Документация asyncio.to_thread отдельно отмечает, что из-за GIL to_thread обычно подходит именно для I/O-bound задач, а CPU-bound — только в особых случаях, например если расширение освобождает GIL.
Для CPU-bound:
import asyncio
from concurrent.futures import ProcessPoolExecutor
def cpu_bound():
return sum(i * i for i in range(10 ** 7))
async def main():
loop = asyncio.get_running_loop()
with ProcessPoolExecutor() as pool:
result = await loop.run_in_executor(pool, cpu_bound)
print(result)
if __name__ == "__main__":
asyncio.run(main())
Важно: для ProcessPoolExecutor нужен guard:
if __name__ == "__main__":
Особенно на Windows. Это также указано в официальной документации Python для примера с process pool.
Что использовать на практике #
I/O-bound blocking code
├─ requests.get()
├─ time.sleep()
├─ обычное чтение файлов
├─ синхронные SDK
└─ синхронные DB/HTTP клиенты
↓
asyncio.to_thread() или ThreadPoolExecutor
CPU-bound blocking code
├─ тяжёлые вычисления
├─ обработка изображений
├─ парсинг больших данных
└─ криптография / сжатие / расчёты
↓
ProcessPoolExecutor
Главное правило #
Нельзя так:
async def handler():
data = requests.get("https://example.com").json()
return data
Потому что requests.get() блокирует event loop.
Лучше так:
import asyncio
import requests
def fetch_sync(url: str):
response = requests.get(url)
response.raise_for_status()
return response.json()
async def handler():
data = await asyncio.to_thread(fetch_sync, "https://example.com")
return data
Но ещё лучше — использовать настоящую async-библиотеку, например httpx.AsyncClient или aiohttp, потому что тогда не нужны дополнительные потоки.
Итог #
run_in_executor нужен, чтобы встроить обычный синхронный блокирующий код в asyncio-приложение:
loop = asyncio.get_running_loop()
result = await loop.run_in_executor(None, blocking_func)
А в современном коде для простого случая чаще пишут:
result = await asyncio.to_thread(blocking_func)
Разница по сути такая:
asyncio.to_thread()
удобная обёртка для запуска функции в отдельном потоке
loop.run_in_executor()
более гибкий низкоуровневый способ:
можно выбрать ThreadPoolExecutor или ProcessPoolExecutor
21. Как проверить, что объект является корутиной? #
Основной способ #
Чтобы проверить, что объект является именно объектом корутины, используют:
import inspect
inspect.iscoroutine(obj)
Пример:
import inspect
async def func():
return 123
coro = func()
print(inspect.iscoroutine(coro)) # True
print(inspect.iscoroutine(func)) # False
coro.close()
inspect.iscoroutine(obj) возвращает True, если объект является native coroutine object, то есть результатом вызова async def функции.
Важно: функция ≠ корутина #
async def func():
return 123
Сама func — это не корутина. Это coroutine function.
import inspect
async def func():
return 123
print(inspect.iscoroutinefunction(func)) # True
print(inspect.iscoroutine(func)) # False
А вот результат вызова func() — уже объект корутины:
coro = func()
print(inspect.iscoroutine(coro)) # True
coro.close()
Схема:
async def func()
↓
coroutine function
func()
↓
coroutine object
Для проверки функции используют inspect.iscoroutinefunction(), а для проверки уже созданной корутины — inspect.iscoroutine().
Проверка на awaitable #
Иногда нужно проверить не “корутина ли это”, а “можно ли это использовать с await”.
Тогда используют:
inspect.isawaitable(obj)
Пример:
import inspect
import asyncio
async def func():
return 123
coro = func()
task = asyncio.create_task(coro)
print(inspect.isawaitable(task)) # True
print(inspect.iscoroutine(task)) # False
await task
Почему так:
coroutine object
└─ awaitable
asyncio.Task
└─ тоже awaitable
└─ но не coroutine object
inspect.isawaitable(obj) проверяет, можно ли объект использовать в выражении await obj. В PEP 492 также отдельно различаются inspect.iscoroutine(obj), inspect.iscoroutinefunction(obj) и inspect.isawaitable(obj).
Практическое правило #
import inspect
if inspect.iscoroutine(obj):
print("Это объект корутины")
if inspect.iscoroutinefunction(obj):
print("Это async-функция")
if inspect.isawaitable(obj):
print("Это можно await-ить")
Разница:
inspect.iscoroutinefunction(func)
проверяет async def функцию
inspect.iscoroutine(func())
проверяет объект корутины
inspect.isawaitable(obj)
проверяет любой объект, который можно await-ить:
coroutine
Task
Future
объект с __await__()
Мини-пример #
import asyncio
import inspect
async def async_func():
return 1
def sync_func():
return 2
async def main():
coro = async_func()
task = asyncio.create_task(coro)
print(inspect.iscoroutinefunction(async_func)) # True
print(inspect.iscoroutinefunction(sync_func)) # False
print(inspect.iscoroutine(coro)) # False, уже обёрнут в Task
print(inspect.iscoroutine(task)) # False
print(inspect.isawaitable(task)) # True
await task
asyncio.run(main())
Но здесь есть важный момент: после передачи coro в asyncio.create_task(coro) корутина уже обёрнута в Task. Поэтому на практике чаще проверяют объект до обёртки:
coro = async_func()
print(inspect.iscoroutine(coro)) # True
coro.close()
Итог #
Для точной проверки объекта корутины:
inspect.iscoroutine(obj)
Для проверки async-функции:
inspect.iscoroutinefunction(obj)
Для проверки, можно ли объект передать в await:
inspect.isawaitable(obj)
22. Сколько процессов может быть у одного приложения? #
У одного приложения может быть:
1 процесс
или
много процессов
Фиксированного универсального числа нет. Количество процессов зависит от:
операционной системы
лимитов пользователя
лимитов контейнера / cgroup
доступной RAM
настроек приложения
модели запуска
Что значит “у приложения несколько процессов” #
Например, приложение может запустить дочерние процессы:
main process
├─ worker process 1
├─ worker process 2
├─ worker process 3
└─ worker process 4
В Python это обычно делается через multiprocessing, ProcessPoolExecutor, Gunicorn workers, Celery workers и т.д.
В документации Python multiprocessing.Process описан как объект, представляющий активность, выполняемую в отдельном процессе. То есть каждый Process — это отдельный OS-level процесс.
Есть ли жёсткий лимит #
На уровне языка Python — нет.
Python не говорит:
одно приложение может иметь максимум N процессов
Но лимиты есть на уровне ОС.
Например, в Linux есть RLIMIT_NPROC — лимит на количество существующих процессов, точнее на Linux это процессы/потоки, для конкретного real user ID. Если лимит достигнут, создание нового процесса через fork() может завершиться ошибкой EAGAIN.
То есть ограничение чаще выглядит не так:
у приложения максимум 100 процессов
а так:
у пользователя / контейнера / системы максимум N процессов или потоков
На Windows #
В Windows процесс создаётся через API вроде CreateProcess, который создаёт новый процесс и его основной поток.
Отдельный лимит может накладываться через Job Object. Например, у Job Object есть ActiveProcessLimit — лимит активных процессов внутри job. Если процесс пытаются добавить в job сверх этого лимита, операция может завершиться ошибкой, а процесс будет завершён.
То есть в Windows тоже нет нормального правила вида:
одно приложение = максимум N процессов
Лимит зависит от ресурсов и ограничений среды.
Практически для Python #
Можно создать несколько процессов так:
from multiprocessing import Process
import os
def worker():
print("PID:", os.getpid())
if __name__ == "__main__":
processes = []
for _ in range(4):
p = Process(target=worker)
p.start()
processes.append(p)
for p in processes:
p.join()
Получится:
1 главный процесс
4 дочерних процесса
Итого приложение фактически использует 5 процессов.
Для backend-приложений #
Например, если запустить FastAPI через Gunicorn:
gunicorn app.main:app -w 4 -k uvicorn.workers.UvicornWorker
то обычно будет:
master process
├─ worker 1
├─ worker 2
├─ worker 3
└─ worker 4
То есть минимум 5 процессов: один master и четыре worker-процесса.
Как выбирать количество процессов #
Обычно количество worker-процессов подбирают по CPU и типу нагрузки:
CPU-bound задачи
→ примерно по количеству CPU cores
I/O-bound задачи
→ можно больше, но часто лучше async + меньше процессов
Web backend
→ несколько worker-процессов + async внутри каждого worker
Пример:
8 CPU cores
→ 4-8 worker-процессов как стартовая точка
Но слишком много процессов — плохо:
много RAM
много context switching
больше overhead на IPC
сложнее синхронизация
больше нагрузка на БД/Redis
Главное отличие от потоков #
процесс
отдельная память
отдельный PID
дороже создать
лучше изоляция
обходит GIL для CPU-bound задач
поток
общая память внутри процесса
дешевле создать
в CPython ограничен GIL для Python-кода
Итог #
У одного приложения может быть сколько угодно процессов в рамках ограничений ОС и ресурсов.
Правильнее отвечать так:
Минимум: 1 процесс.
Максимум: не задан Python или самим понятием "приложение".
Он ограничен ОС, лимитами пользователя/контейнера, памятью, CPU и настройками запуска.
В Python каждый multiprocessing.Process создаёт отдельный процесс, а в production backend несколько процессов часто используют как worker-процессы для параллельной обработки запросов.
23. Какие инструменты синхронизации потоков есть в Python? | Для чего нужны Lock, RLock, Semaphore, Event, Condition #
Зачем нужны инструменты синхронизации #
В Python потоки внутри одного процесса могут обращаться к общим данным:
counter += 1
shared_list.append(item)
cache[key] = value
Проблема в том, что несколько потоков могут одновременно менять одно и то же состояние. Это приводит к race condition: результат зависит от порядка выполнения потоков. GIL не отменяет необходимость синхронизации, потому что он не делает всю пользовательскую логику атомарной и не защищает твои структуры данных на уровне бизнес-операций. В документации threading отдельно описаны примитивы синхронизации: Lock, RLock, Semaphore, BoundedSemaphore, Event, Condition, Barrier.
Основные инструменты #
Lock
обычная взаимная блокировка
RLock
рекурсивная блокировка
Semaphore
ограничитель количества одновременных потоков
BoundedSemaphore
Semaphore с защитой от лишних release()
Event
флаг-сигнал между потоками
Condition
ожидание некоторого состояния под lock
Barrier
ожидание группы потоков в одной точке
Lock, RLock, Condition, Semaphore и BoundedSemaphore можно использовать через with, то есть как context manager. Это безопаснее, потому что release() будет вызван даже при исключении.
Lock #
Lock нужен, когда только один поток должен выполнять критическую секцию.
import threading
lock = threading.Lock()
counter = 0
def increment():
global counter
for _ in range(100_000):
with lock:
counter += 1
Без Lock несколько потоков могут одновременно читать и менять counter, из-за чего часть инкрементов потеряется.
Схема:
Thread 1 ── acquire lock ── работает ── release lock
Thread 2 ── ждёт ─── работает после освобождения
Использовать, когда нужно защитить:
общий счётчик
общий dict/list/set
запись в общий файл
общий cache
общий объект с изменяемым состоянием
RLock #
RLock — reentrant lock, то есть рекурсивная блокировка.
Она нужна, когда один и тот же поток может захватить lock несколько раз.
import threading
lock = threading.RLock()
def outer():
with lock:
inner()
def inner():
with lock:
print("Работает")
С обычным Lock такой код может зависнуть:
lock = threading.Lock()
def outer():
with lock:
inner()
def inner():
with lock:
print("deadlock")
Почему:
outer() захватил Lock
inner() пытается захватить тот же Lock
тот же поток ждёт сам себя
RLock хранит владельца и счётчик захватов. Один и тот же поток может вызвать acquire() несколько раз, но должен столько же раз вызвать release(). Это поведение описано в документации Python для RLock.
Использовать, когда:
методы одного класса вызывают друг друга
и каждый метод должен быть thread-safe
Пример:
import threading
class BankAccount:
def __init__(self):
self._balance = 0
self._lock = threading.RLock()
def deposit(self, amount):
with self._lock:
self._balance += amount
def deposit_bonus(self, amount):
with self._lock:
self.deposit(amount)
self.deposit(10)
Semaphore #
Semaphore — ограничитель количества одновременных доступов.
Он нужен, когда ресурсом может пользоваться не один поток, а ограниченное количество потоков.
Например, максимум 3 потока одновременно могут делать сетевой запрос:
import threading
import time
semaphore = threading.Semaphore(3)
def request():
with semaphore:
print("start")
time.sleep(1)
print("end")
threads = [
threading.Thread(target=request)
for _ in range(10)
]
for thread in threads:
thread.start()
for thread in threads:
thread.join()
Идея:
Lock → пускает только 1 поток
Semaphore → пускает N потоков
Пример использования:
ограничить количество подключений к API
ограничить количество одновременных загрузок файлов
ограничить доступ к пулу ресурсов
BoundedSemaphore #
BoundedSemaphore похож на Semaphore, но защищает от ошибки лишнего release().
import threading
sem = threading.BoundedSemaphore(2)
sem.acquire()
sem.release()
sem.release() # ошибка, если release() вызван слишком много раз
Обычный Semaphore может увеличить счётчик выше начального значения. BoundedSemaphore проверяет, чтобы счётчик не превысил начальное значение. Это полезно для поиска багов в логике синхронизации.
Event #
Event — это флаг-сигнал между потоками.
Один поток ждёт:
event.wait()
Другой поток сообщает:
event.set()
Пример:
import threading
import time
event = threading.Event()
def worker():
print("Жду сигнал")
event.wait()
print("Сигнал получен")
thread = threading.Thread(target=worker)
thread.start()
time.sleep(2)
event.set()
thread.join()
Схема:
Thread 1: event.wait() ── ждёт
Thread 2: event.set() ── подаёт сигнал
Thread 1: продолжает работу
У Event есть внутренний флаг. set() устанавливает его в True, clear() сбрасывает в False, wait() блокирует поток, пока флаг не станет True.
Использовать, когда:
один поток должен дождаться старта другого
нужно сообщить потокам "можно начинать"
нужно сообщить "пора завершаться"
нужно сделать простой сигнал без передачи данных
Пример остановки worker-потока:
import threading
import time
stop_event = threading.Event()
def worker():
while not stop_event.is_set():
print("working")
time.sleep(1)
print("stopped")
thread = threading.Thread(target=worker)
thread.start()
time.sleep(3)
stop_event.set()
thread.join()
Condition #
Condition нужен, когда поток должен ждать не просто сигнал, а выполнение некоторого условия.
Например:
ждать, пока очередь не станет непустой
ждать, пока появятся данные
ждать, пока состояние объекта изменится
Пример producer/consumer:
import threading
import time
condition = threading.Condition()
items = []
def consumer():
with condition:
while not items:
condition.wait()
item = items.pop()
print("got:", item)
def producer():
time.sleep(2)
with condition:
items.append("data")
condition.notify()
t1 = threading.Thread(target=consumer)
t2 = threading.Thread(target=producer)
t1.start()
t2.start()
t1.join()
t2.join()
Как это работает:
consumer:
захватывает condition lock
проверяет items
если items пустой:
condition.wait()
отпускает lock
засыпает
producer:
захватывает condition lock
добавляет item
вызывает notify()
отпускает lock
consumer:
просыпается
снова захватывает lock
продолжает работу
В документации указано важное поведение: wait() освобождает lock и блокирует поток, а после пробуждения снова захватывает lock перед возвратом из wait().
Обычно Condition используют с циклом while, а не с if:
with condition:
while not items:
condition.wait()
item = items.pop()
Так безопаснее, потому что поток может проснуться, но нужное условие всё ещё может быть ложным.
Lock vs RLock vs Semaphore vs Event vs Condition #
Lock
"Только один поток может войти в эту секцию"
RLock
"Только один поток, но этот же поток может войти повторно"
Semaphore
"Не больше N потоков одновременно"
Event
"Ждите, пока кто-то подаст сигнал"
Condition
"Ждите, пока под lock выполнится конкретное условие"
Практическая таблица #
| Инструмент | Что делает | Когда использовать |
|---|---|---|
Lock | Пускает только один поток | Защита общего состояния |
RLock | То же, но один поток может захватить повторно | Вложенные thread-safe методы |
Semaphore | Пускает максимум N потоков | Лимит подключений/запросов/ресурсов |
BoundedSemaphore | Как Semaphore, но ловит лишний release() | Более безопасный semaphore |
Event | Сигнал-флаг между потоками | Старт/стоп/готовность |
Condition | Ожидание условия + lock | Producer/consumer, очереди, состояние |
Barrier | Ждёт, пока соберётся группа потоков | Синхронный старт/этапы обработки |
Главное правило #
Для защиты данных:
with lock:
# работа с общим состоянием
Для ожидания сигнала:
event.wait()
Для ограничения параллельного доступа:
with semaphore:
# работа с ограниченным ресурсом
Для ожидания изменения состояния:
with condition:
while not condition_is_true:
condition.wait()
В реальном коде часто используют не голые Condition, а готовую потокобезопасную очередь queue.Queue, потому что она уже внутри использует нужную синхронизацию для producer/consumer-сценариев.
24. Что такое корутины и задачи? #
В asyncio есть два важных понятия:
Coroutine
объект, который описывает асинхронную операцию
Task
обёртка над coroutine, которая ставит её на выполнение в event loop
Официальная документация Python относит coroutines, Tasks и Futures к awaitable-объектам, то есть объектам, которые можно использовать с await.
Coroutine function и coroutine object #
Когда ты пишешь:
async def get_data():
return 123
это coroutine function, то есть асинхронная функция.
Но сама по себе она ещё не выполняется.
Когда ты вызываешь её:
coro = get_data()
создаётся coroutine object.
Схема:
async def get_data()
↓
coroutine function
get_data()
↓
coroutine object
Документация Python прямо разделяет эти два смысла слова “coroutine”: coroutine function — это async def функция, а coroutine object — объект, возвращаемый вызовом этой функции.
Корутина не запускается сама #
Важный момент:
async def get_data():
print("started")
return 123
coro = get_data()
После строки:
coro = get_data()
код внутри get_data() ещё не выполнился.
Чтобы корутина начала выполняться, её нужно:
await-ить
или
обернуть в Task
Пример через await:
import asyncio
async def get_data():
print("started")
return 123
async def main():
result = await get_data()
print(result)
asyncio.run(main())
Здесь await get_data() запускает корутину и ждёт её результат.
Что такое Task #
Task — это объект, который берёт корутину и планирует её выполнение в event loop.
task = asyncio.create_task(get_data())
asyncio.create_task(coro) оборачивает корутину в Task, ставит её на выполнение и возвращает объект задачи.
Пример:
import asyncio
async def get_data():
await asyncio.sleep(1)
return 123
async def main():
task = asyncio.create_task(get_data())
print("task created")
result = await task
print(result)
asyncio.run(main())
Схема:
get_data()
↓
coroutine object
↓
asyncio.create_task(...)
↓
Task
↓
event loop выполняет эту задачу
Главное отличие coroutine от Task #
Coroutine
просто объект асинхронной операции
сам по себе не запланирован в event loop
не выполняется, пока его не await-нут или не обернут в Task
Task
coroutine, привязанная к event loop
автоматически запланирована на выполнение
может выполняться конкурентно с другими Task
Task — это Future-like объект, который запускает Python-корутину в event loop. Если корутина внутри Task ждёт Future, выполнение корутины приостанавливается, а потом продолжается, когда Future завершится.
await coroutine — это не параллельность #
Пример последовательного выполнения:
import asyncio
async def job(name):
print(f"{name} started")
await asyncio.sleep(1)
print(f"{name} finished")
async def main():
await job("A")
await job("B")
asyncio.run(main())
Здесь сначала полностью выполняется A, потом B.
Примерно так:
A started
A finished
B started
B finished
То есть простой await означает:
запусти корутину
дождись её завершения
потом иди дальше
Task даёт конкурентное выполнение #
Чтобы задачи выполнялись одновременно с точки зрения event loop, их нужно создать заранее:
import asyncio
async def job(name):
print(f"{name} started")
await asyncio.sleep(1)
print(f"{name} finished")
async def main():
task_a = asyncio.create_task(job("A"))
task_b = asyncio.create_task(job("B"))
await task_a
await task_b
asyncio.run(main())
Теперь обе задачи запланированы в event loop.
Примерная схема:
task A starts
task A awaits sleep
event loop switches to task B
task B starts
task B awaits sleep
event loop continues other work
То есть Task нужна, когда ты хочешь не просто выполнить одну корутину, а запустить несколько корутин конкурентно.
Через gather #
Часто несколько задач запускают через asyncio.gather():
import asyncio
async def job(name):
await asyncio.sleep(1)
return name
async def main():
results = await asyncio.gather(
job("A"),
job("B"),
job("C"),
)
print(results)
asyncio.run(main())
Если в asyncio.gather() передать корутины, они автоматически будут запланированы как задачи.
Task — это не поток #
Важно:
Task ≠ Thread
Task ≠ Process
asyncio.Task не создаёт новый системный поток. Все задачи обычно выполняются в одном потоке event loop и переключаются кооперативно: задача отдаёт управление, когда доходит до await.
один event loop
├─ Task A
├─ Task B
├─ Task C
└─ Task D
Event loop по очереди даёт задачам выполняться, а когда задача приостанавливается на await, запускает другую работу. Концептуальный обзор asyncio описывает event loop как механизм, который берёт работу из очереди, запускает её, а когда она приостанавливается или завершается, переходит к другой работе.
Что можно делать с Task #
У Task есть управление состоянием:
task = asyncio.create_task(job())
task.cancel() # отменить задачу
task.done() # завершена ли задача
task.result() # получить результат, если задача завершена
task.exception() # получить исключение, если оно было
Пример:
import asyncio
async def job():
await asyncio.sleep(1)
return 100
async def main():
task = asyncio.create_task(job())
result = await task
print(task.done()) # True
print(result) # 100
asyncio.run(main())
Практическое сравнение #
| Объект | Что это | Запускается сам? | Можно await? |
|---|---|---|---|
async def func | coroutine function | Нет | Нет |
func() | coroutine object | Нет | Да |
asyncio.create_task(func()) | Task | Да, планируется в event loop | Да |
Итог #
Корутина — это асинхронная операция, созданная вызовом async def функции:
coro = some_async_func()
Task — это корутина, поставленная на выполнение в event loop:
task = asyncio.create_task(some_async_func())
Главная разница:
coroutine
описывает работу
task
запускает эту работу через event loop и позволяет управлять её выполнением
25. Потокобезопасен ли python? #
Нет, Python-код сам по себе не становится потокобезопасным.
Но CPython — стандартная реализация Python — имеет GIL, который защищает внутренности интерпретатора от одновременного выполнения Python bytecode несколькими потоками. Это не защищает твою бизнес-логику от race condition.
Что реально защищает GIL #
В обычном CPython с включённым GIL:
несколько потоков есть
↓
но Python bytecode в один момент времени выполняет только один поток
↓
интерпретатор и базовые объекты не ломаются от параллельного доступа
То есть GIL помогает CPython безопасно работать с объектами на внутреннем уровне, но не делает операции над общими данными логически безопасными. Документация прямо говорит, что блокировки всё равно нужны для thread-safety в pure-Python коде.
Что не потокобезопасно #
Например:
counter = 0
def inc():
global counter
counter += 1
counter += 1 выглядит как одна операция, но логически это несколько шагов:
прочитать counter
прибавить 1
записать обратно
Между этими шагами поток может быть переключён, поэтому результат может быть неправильным.
Правильно:
import threading
counter = 0
lock = threading.Lock()
def inc():
global counter
with lock:
counter += 1
А встроенные list/dict/set? #
Некоторые одиночные операции в CPython атомарны, например:
lst.append(x)
d[key] = value
x = lst.pop()
В FAQ Python указано, что операции над встроенными типами, которые «выглядят атомарными», в CPython обычно действительно атомарны, потому что переключение потоков происходит между bytecode-инструкциями. Но составные операции уже не безопасны: i = i + 1, D[x] = D[x] + 1, L.append(L[-1]). В документации прямо сказано: когда сомневаешься — используй mutex/lock.
Важный момент про Python 3.13+ #
Начиная с Python 3.13, у CPython есть free-threaded build, где GIL можно отключить. Это позволяет потокам реально выполняться параллельно на разных ядрах, но это отдельная сборка, а не обычное поведение по умолчанию.
В free-threaded Python встроенные типы вроде dict, list, set используют внутренние блокировки для защиты от конкурентных модификаций, но это всё равно не означает, что любая логика приложения автоматически потокобезопасна.
Итог #
Python / CPython безопасен на уровне интерпретатора
≠
твой многопоточный код потокобезопасен
Правило:
Есть общее изменяемое состояние между потоками
→ нужен Lock / RLock / Semaphore / Queue / Condition / Event
Для I/O-bound задач потоки в Python подходят нормально.
Для CPU-bound задач обычные потоки в CPython ограничены GIL, поэтому чаще используют multiprocessing или ProcessPoolExecutor.
26. Как запустить синхронный код в асинхронном проекте #
Проблема
Прямой синхронный вызов в асинхронном коде блокирует весь event loop, останавливая выполнение всех других задач.
Решение: 3 подхода
- Использование run_in_executor() с ThreadPoolExecutor
Для сетевых запросов, работы с файлами, обращений к БД:
import asyncio import requests
async def fetch_data(url): loop = asyncio.get_event_loop() response = await loop.run_in_executor( None, # Стандартный пул потоков requests.get, url ) return response.json() Суть: Блокирующая операция выполняется в фоновом потоке, не блокируя основной event loop.
- Явное управление пулом потоков
Когда нужно ограничить количество одновременных операций:
async def process_requests(urls): with ThreadPoolExecutor(max_workers=3) as executor: loop = asyncio.get_event_loop() tasks = [] for url in urls: task = loop.run_in_executor(executor, requests.get, url) tasks.append(task) responses = await asyncio.gather(*tasks) return [r.json() for r in responses] Суть: Создаем пул с ограничением потоков для контроля нагрузки.
- ProcessPoolExecutor для CPU-bound задач
Для вычислений, обработки данных, ML:
import asyncio from concurrent.futures import ProcessPoolExecutor
def heavy_calculation(data): return sum(i * i for i in range(data))
async def main(): with ProcessPoolExecutor() as executor: loop = asyncio.get_event_loop() result = await loop.run_in_executor( executor, heavy_calculation, 1000000 ) return result Суть: Вычисления в отдельном процессе, обход GIL для истинной параллельности.
Ключевые правила
I/O операции → ThreadPoolExecutor Вычисления → ProcessPoolExecutor Всегда предпочитайте асинхронные библиотеки когда возможно Никогда не делайте прямые синхронные вызовы в асинхронном коде
27. Что такое зелёные потоки? #
Зеленые потоки (green threads) являются примитивным уровнем асинхронного программирования.
Зеленый поток — это обычный поток, за исключением того, что переключения между потоками производятся в коде приложения, а не в процессоре.Gevent — известная Python-библиотека для использования зеленых потоков неблокирующего ввода-вывода Eventlet. Gevent.monkey изменяет поведение стандартных библиотек Python таким образом, что они позволяют выполнять неблокирующие операции ввода-вывода.
Вот пример использования Gevent для одновременного обращения к нескольким URL-адресам:
import gevent.monkey
from urllib.request import urlopen
gevent.monkey.patch_all()
urls = ['http://www.google.com', 'http://www.yandex.ru', 'http://www.python.org']
def print_head(url):
print('Starting {}'.format(url))
data = urlopen(url).read()
print('{}: {} bytes: {}'.format(url, len(data), data))
jobs = [gevent.spawn(print_head, _url) for _url in urls]
gevent.wait(jobs)
Как видите, API-интерфейс Gevent выглядит так же, как и потоки.
Однако за кадром он использует сопрограммы (coroutines), а не потоки, и запускает их в цикле событий (event loop) для постановки в очередь. Это значит, что вы получаете преимущества потоков, без понимания сопрограмм, но вы не избавляетесь от проблем, связанных с потоками.Gevent — хорошая библиотека, но только для тех, кто понимает, как работают потоки.
28. Сколько ядер CPU использует asyncio в Python? #
Обычно asyncio использует 1 ядро CPU для выполнения Python-кода event loop.
Точнее:
1 event loop
→ работает в 1 потоке
→ этот поток выполняется на 1 ядре CPU в конкретный момент времени
В документации Python указано, что event loop выполняется в потоке, обычно в главном, и все callbacks/tasks этого event loop выполняются в этом же потоке. Пока одна Task выполняется в event loop, другая задача в этом же loop не выполняется.
Важное уточнение #
asyncio даёт конкурентность, а не CPU-параллелизм.
То есть asyncio умеет эффективно переключаться между задачами, когда они ждут I/O:
задача A ждёт сеть
→ event loop переключился на задачу B
задача B ждёт БД
→ event loop переключился на задачу C
Но это не значит, что несколько Python-корутин одновременно считают на разных ядрах.
async def task_1():
...
async def task_2():
...
Эти задачи в одном event loop не выполняются физически одновременно на разных ядрах. Они выполняются по очереди, просто переключение происходит быстро и в точках await.
Пример #
import asyncio
async def task(name):
print(f"{name}: start")
await asyncio.sleep(1)
print(f"{name}: end")
async def main():
await asyncio.gather(
task("A"),
task("B"),
task("C"),
)
asyncio.run(main())
Здесь три задачи выполняются конкурентно:
A start
B start
C start
ожидание sleep
A end
B end
C end
Но это не значит, что Python занял 3 ядра. Один event loop всё равно работает в одном потоке.
Когда asyncio может использовать больше ядер #
Сам по себе asyncio — нет. Но async-приложение может задействовать больше ядер через дополнительные механизмы.
1. Через несколько процессов #
Например, в FastAPI/Gunicorn/Uvicorn можно запустить несколько worker-процессов:
worker 1 → event loop → ядро CPU
worker 2 → event loop → ядро CPU
worker 3 → event loop → ядро CPU
worker 4 → event loop → ядро CPU
Каждый процесс имеет свой Python-интерпретатор, свой event loop и может выполняться на отдельном ядре.
2. Через ProcessPoolExecutor
#
Для тяжёлых CPU-bound задач:
import asyncio
from concurrent.futures import ProcessPoolExecutor
def cpu_heavy(n: int) -> int:
total = 0
for i in range(n):
total += i * i
return total
async def main():
loop = asyncio.get_running_loop()
with ProcessPoolExecutor() as pool:
result = await loop.run_in_executor(pool, cpu_heavy, 10_000_000)
print(result)
asyncio.run(main())
run_in_executor() позволяет запускать блокирующий код в пуле потоков или процессов. Для CPU-bound задач в документации asyncio приводится вариант с ProcessPoolExecutor.
3. Через нативные библиотеки #
Некоторые библиотеки могут сами использовать несколько ядер внутри C/C++/Rust-кода:
NumPy
OpenCV
PyTorch
TensorFlow
Но это уже не сам asyncio. Это внутренняя параллельность конкретной библиотеки.
Что с потоками #
Можно запустить синхронный код через:
await asyncio.to_thread(sync_func)
или:
await loop.run_in_executor(None, sync_func)
Но для обычного Python CPU-bound кода потоки не дадут нормального параллелизма из-за GIL. Они полезнее для блокирующего I/O: файлов, сети, синхронных SDK, старых библиотек.
Итог #
Один asyncio event loop
→ обычно 1 поток
→ в конкретный момент выполняется на 1 ядре CPU
→ много корутин не означают много ядер
Для использования нескольких ядер нужны:
CPU-bound задачи
→ multiprocessing / ProcessPoolExecutor / несколько worker-процессов
I/O-bound задачи
→ asyncio достаточно, даже на одном ядре
Главная мысль:
asyncio масштабирует ожидание I/O,
а не вычисления по ядрам CPU.
29. Что такое event loop в асинхронности? #
Что такое event loop #
event loop — это центральный механизм асинхронного выполнения в asyncio.
Он отвечает за то, чтобы:
запускать coroutine-задачи
переключаться между ними
следить за I/O-операциями
выполнять callbacks
возвращать управление задачам, когда они готовы продолжить работу
В документации Python event loop описан как ядро каждого asyncio-приложения: он запускает асинхронные задачи и callbacks, выполняет сетевые I/O-операции и работает с subprocess.
Простая аналогия #
Event loop — это диспетчер задач.
У него есть несколько задач:
Task A: ждёт ответ от базы данных
Task B: ждёт HTTP-запрос
Task C: ждёт таймер
Task D: уже готова продолжить выполнение
Event loop не сидит без дела, пока одна задача ждёт. Он переключается на другую задачу, которая уже может выполняться.
Task A дошла до await
↓
A временно остановлена
event loop смотрит:
есть ли другая готовая задача?
↓
Task B продолжает выполнение
Как это выглядит в коде #
import asyncio
async def task(name: str, delay: int):
print(f"{name}: start")
await asyncio.sleep(delay)
print(f"{name}: end")
async def main():
await asyncio.gather(
task("A", 2),
task("B", 1),
task("C", 3),
)
asyncio.run(main())
Здесь asyncio.run(main()) создаёт и запускает event loop для верхнего уровня программы. Документация Python указывает, что asyncio.run() используется для запуска top-level async entry point.
Что происходит пошагово #
1. Запускается main()
2. asyncio.gather() регистрирует три coroutine как задачи
3. Task A доходит до await asyncio.sleep(2)
4. Task A отдаёт управление event loop
5. Event loop запускает Task B
6. Task B доходит до await asyncio.sleep(1)
7. Task B отдаёт управление event loop
8. Event loop запускает Task C
9. Когда sleep у B закончился — event loop возвращает выполнение B
10. Потом возвращает A и C, когда они тоже готовы
Главная мысль:
await — это точка, где coroutine может уступить управление event loop.
Event loop не выполняет всё одновременно #
Важно: event loop не делает Python-код параллельным сам по себе.
Обычно схема такая:
один event loop
↓
один поток
↓
одна задача реально выполняется в конкретный момент
Но когда задача ждёт I/O, например сеть или БД, event loop может выполнять другие задачи.
Параллельность CPU:
несколько задач реально считаются одновременно на разных ядрах
Конкурентность asyncio:
задачи быстро чередуются, пока часть из них ждёт I/O
asyncio в официальной документации описывается как библиотека для конкурентного кода через async/await, особенно подходящая для I/O-bound и сетевого кода.
Что event loop отслеживает #
Event loop обычно работает с такими вещами:
готовые задачи
ожидающие задачи
таймеры
сетевые события
callbacks
Future-объекты
subprocess
executor для потоков/процессов
Например, когда coroutine делает:
await asyncio.sleep(1)
она не блокирует поток на 1 секунду. Она говорит event loop:
"Поставь меня на паузу и вернись ко мне примерно через 1 секунду"
Почему нельзя блокировать event loop #
Плохой пример:
import asyncio
import time
async def bad():
time.sleep(5) # блокирует весь event loop
print("done")
asyncio.run(bad())
Пока работает time.sleep(5), event loop не может переключиться на другие задачи.
Правильно:
import asyncio
async def good():
await asyncio.sleep(5)
print("done")
asyncio.run(good())
Или, если нужно запустить обычный синхронный блокирующий код:
import asyncio
import time
def sync_work():
time.sleep(5)
return "done"
async def main():
result = await asyncio.to_thread(sync_work)
print(result)
asyncio.run(main())
Event loop vs coroutine vs task #
coroutine
→ асинхронная функция, созданная через async def
task
→ coroutine, запланированная на выполнение в event loop
event loop
→ механизм, который запускает task и переключается между ними
Пример:
async def hello():
return "hello"
Это coroutine-функция.
coro = hello()
Это coroutine-объект.
task = asyncio.create_task(hello())
Это уже задача, которую event loop может выполнять. Документация Python отдельно указывает, что простой вызов coroutine не запускает её выполнение; для запуска используются asyncio.run(), await или создание Task.
Короткая схема #
asyncio.run(main())
↓
создаётся event loop
↓
main() запускается как coroutine
↓
создаются task
↓
task выполняются до await
↓
на await задача отдаёт управление
↓
event loop запускает другую готовую task
↓
когда I/O готово — задача продолжается
Итог #
event loop — это диспетчер асинхронных задач.
Он не делает CPU-вычисления параллельными, но позволяет эффективно использовать время ожидания:
пока одна coroutine ждёт сеть / БД / файл / таймер,
event loop запускает другую coroutine.
Поэтому asyncio особенно полезен для I/O-bound задач: сетевые запросы, базы данных, WebSocket, HTTP-серверы, очереди, большое количество одновременных соединений.
30. На чём основана асинхронность в Python под капотом? #
Асинхронность в Python основана на 4 главных вещах:
1. coroutine-объекты
2. event loop
3. cooperative scheduling
4. неблокирующий I/O через системные механизмы ОС
То есть asyncio — это не магия и не автоматические потоки. Это модель, где задачи сами отдают управление в точках await, а event loop решает, какую задачу продолжить дальше.
1. Coroutine: функция, которую можно приостанавливать #
Когда ты пишешь:
async def fetch_data():
await some_io()
return "done"
fetch_data() не выполняется сразу как обычная функция.
coro = fetch_data()
Здесь создаётся coroutine-объект.
async def
↓
создаёт coroutine-функцию
вызов async-функции
↓
создаёт coroutine-объект
await / Task / event loop
↓
запускает или продолжает coroutine
PEP 492 ввёл нативные корутины и синтаксис async / await; async def функции являются корутинами.
2. await: точка добровольной остановки
#
await — это место, где coroutine может приостановиться и отдать управление event loop.
async def main():
print("before")
await asyncio.sleep(1)
print("after")
На await asyncio.sleep(1) coroutine говорит:
"Я сейчас жду.
Пока меня не продолжай.
Можешь запустить другую задачу."
Важно: переключение не происходит в любой случайный момент. Оно происходит кооперативно — когда задача доходит до await и ожидаемый объект ещё не готов.
3. Event loop: диспетчер асинхронных задач #
event loop хранит и обслуживает:
готовые задачи
спящие задачи
таймеры
сетевые события
Future-объекты
callbacks
Упрощённо:
event loop:
взять готовую задачу
выполнить её до ближайшего await
если задача ждёт I/O — поставить её на ожидание
если другая задача готова — запустить её
когда I/O готово — вернуть задачу в очередь готовых
Документация Python описывает asyncio как библиотеку для конкурентного кода через async/await; она включает высокоуровневые API для запуска корутин, управления event loop, сетевого I/O, subprocess и синхронизации.
4. Cooperative scheduling #
В asyncio используется кооперативная многозадачность.
ОС-потоки:
операционная система может вытеснить поток почти в любой момент
asyncio-задачи:
задача сама отдаёт управление через await
Пример:
import asyncio
async def task_a():
print("A start")
await asyncio.sleep(1)
print("A end")
async def task_b():
print("B start")
await asyncio.sleep(1)
print("B end")
async def main():
await asyncio.gather(task_a(), task_b())
asyncio.run(main())
Схема:
task_a start
↓
await sleep
↓
event loop переключается на task_b
↓
task_b start
↓
await sleep
↓
event loop ждёт ближайшее событие
↓
продолжает готовые задачи
В концептуальном обзоре asyncio показано, что coroutine можно обернуть в Task, а Task планируется на выполнение через event loop.
5. Future: объект будущего результата #
Под капотом многие операции в asyncio завязаны на Future.
Future — это объект, который представляет результат, который появится позже.
Future:
пока результата нет → pending
результат появился → done
ошибка появилась → exception
отменили → cancelled
Coroutine может ждать Future:
result = await future
Пока Future не готов, coroutine приостанавливается. Когда Future получает результат, event loop возвращает coroutine в очередь готовых задач.
В документации Python asyncio.Future описан как awaitable-объект, представляющий eventual result асинхронной операции.
6. Неблокирующий I/O и мультиплексирование #
Самое важное под капотом: asyncio эффективно работает с I/O не потому, что создаёт поток на каждый запрос, а потому что использует неблокирующий I/O.
Упрощённо:
обычный socket.recv()
↓
поток зависает, пока данные не придут
non-blocking socket
↓
если данных нет, поток не зависает
event loop ждёт уведомления от ОС
Для ожидания множества I/O-событий Python использует механизмы мультиплексирования I/O. Модуль selectors предоставляет высокоуровневое эффективное I/O-мультиплексирование поверх примитивов модуля select.
На уровне ОС это может быть:
Linux → epoll
macOS → kqueue
Windows → IOCP / ProactorEventLoop
Идея такая:
event loop спрашивает ОС:
"сообщи мне, когда один из этих сокетов будет готов"
ОС сообщает:
"сокет #4 готов к чтению"
event loop:
"значит, задачу, которая ждала сокет #4, можно продолжить"
7. Почему это работает без множества потоков #
Потому что при сетевом I/O программа большую часть времени не считает, а ждёт:
ждёт ответ от БД
ждёт HTTP-ответ
ждёт данные из socket
ждёт Redis
ждёт Kafka/RabbitMQ
asyncio использует это время ожидания:
задача A ждёт сеть
↓
event loop запускает задачу B
задача B ждёт БД
↓
event loop запускает задачу C
Поэтому asyncio хорошо подходит для I/O-bound задач, но не делает CPU-bound вычисления автоматически параллельными.
8. Что происходит при asyncio.run(main())
#
import asyncio
async def main():
print("hello")
asyncio.run(main())
Упрощённо:
asyncio.run(main())
↓
создаёт event loop
↓
создаёт coroutine main()
↓
запускает её в event loop
↓
выполняет задачи до завершения
↓
закрывает event loop
asyncio.run() в документации описан как высокоуровневый способ запуска asyncio-кода, построенный поверх event loop.
9. Общая схема под капотом #
async def функция
↓
coroutine object
↓
Task
↓
event loop
↓
выполнение до await
↓
ожидание Future / I/O / таймера
↓
selector / ОС сообщает о готовности
↓
Future становится done
↓
Task возвращается в очередь готовых
↓
coroutine продолжается
Итог #
Асинхронность в Python под капотом основана не на параллельном выполнении корутин, а на управляемом приостановлении и возобновлении выполнения.
async/await
→ синтаксис для корутин
coroutine
→ объект, который можно приостановить
Task
→ coroutine, запланированная в event loop
Future
→ будущий результат операции
event loop
→ планировщик задач
selectors / epoll / kqueue / IOCP
→ механизм ОС для ожидания I/O без блокировки потока
Главная формула:
asyncio = корутины + event loop + cooperative scheduling + non-blocking I/O
31. Разница между корутинами и функциями #
Обычная функция выполняется сразу и до конца, пока не вернёт результат или не выбросит исключение.
Корутина — это функция, выполнение которой можно приостановить и потом продолжить. В Python корутина обычно создаётся через async def, а запускается через await, asyncio.run() или asyncio.create_task(). В документации Python термин “coroutine” используется для двух вещей: coroutine function — функция async def, и coroutine object — объект, полученный при вызове этой функции.
Обычная функция #
def get_data():
return "data"
result = get_data()
print(result)
Что происходит:
get_data() вызвана
↓
код внутри функции сразу выполняется
↓
return возвращает результат
↓
result получает "data"
Обычная функция при вызове сразу начинает выполняться.
Корутина #
import asyncio
async def get_data():
await asyncio.sleep(1)
return "data"
coro = get_data()
print(coro)
Здесь get_data() не выполняет тело функции сразу. Она создаёт coroutine object.
Чтобы реально выполнить корутину, нужно:
result = await get_data()
или:
asyncio.run(get_data())
или создать задачу:
task = asyncio.create_task(get_data())
Документация Python отдельно показывает, что простой вызов coroutine function создаёт coroutine object, а выполнение происходит через await или запуск в event loop.
Главное отличие #
Обычная функция:
вызвал → выполнилась → вернула результат
Корутина:
вызвал → получил coroutine object
await → корутина начала/продолжила выполнение
Пример:
def regular_func():
print("regular")
return 10
async def coroutine_func():
print("coroutine")
return 20
a = regular_func()
b = coroutine_func()
print(a)
print(b)
Результат будет примерно такой:
regular
10
<coroutine object coroutine_func at ...>
Почему?
regular_func()
→ сразу выполнилась
coroutine_func()
→ не выполнилась, а вернула coroutine object
await — точка приостановки
#
Корутина может приостанавливаться на await:
async def handler():
print("start")
await asyncio.sleep(1)
print("end")
Схема:
handler начал выполнение
↓
print("start")
↓
await asyncio.sleep(1)
↓
корутина приостановилась
↓
event loop может выполнять другие задачи
↓
через 1 секунду handler продолжился
↓
print("end")
Обычная функция так не умеет. Если внутри обычной функции стоит блокирующая операция:
import time
def handler():
print("start")
time.sleep(1)
print("end")
то поток просто блокируется на time.sleep(1).
Разница в возвращаемом результате #
Обычная функция:
def add(a, b):
return a + b
result = add(2, 3)
result = 5
Корутина:
async def add(a, b):
return a + b
result = add(2, 3)
result = <coroutine object ...>
Чтобы получить 5, нужно:
result = await add(2, 3)
или:
result = asyncio.run(add(2, 3))
Таблица различий #
| Критерий | Обычная функция | Корутина |
|---|---|---|
| Объявление | def | async def |
| Что возвращает при вызове | результат выполнения | coroutine object |
| Выполняется ли сразу при вызове | да | нет |
Может использовать await | нет | да |
| Может приостанавливаться | нет | да |
| Кто управляет выполнением | обычный поток выполнения | event loop |
| Основная область применения | обычная логика, CPU-код, синхронный код | I/O-bound задачи, сеть, БД, таймеры |
| Можно ли запустить конкурентно | только через потоки/процессы | через asyncio.create_task, gather |
Проверка через inspect
#
import inspect
def regular_func():
pass
async def coroutine_func():
pass
print(inspect.iscoroutinefunction(regular_func))
print(inspect.iscoroutinefunction(coroutine_func))
Результат:
False
True
inspect.iscoroutinefunction() проверяет, является ли объект coroutine function, то есть функцией, определённой через async def. inspect.iscoroutine() проверяет уже coroutine object, полученный после вызова coroutine function.
Важный пример #
import asyncio
async def load_user():
print("load user start")
await asyncio.sleep(1)
print("load user end")
return {"id": 1}
async def load_orders():
print("load orders start")
await asyncio.sleep(1)
print("load orders end")
return [1, 2, 3]
async def main():
user_task = asyncio.create_task(load_user())
orders_task = asyncio.create_task(load_orders())
user = await user_task
orders = await orders_task
print(user, orders)
asyncio.run(main())
Здесь load_user() и load_orders() могут выполняться конкурентно, потому что обе корутины отдают управление event loop на await asyncio.sleep(1).
С обычными функциями без потоков/процессов такого переключения не будет.
Итог #
Функция:
обычный вызываемый блок кода,
который выполняется сразу и возвращает результат.
Корутина:
асинхронная функция,
которая при вызове создаёт coroutine object,
может приостанавливаться на await
и выполняется под управлением event loop.
Самая короткая разница:
def func()
→ вызов сразу выполняет код
async def func()
→ вызов создаёт coroutine object
→ выполнение начинается через await / Task / event loop
32. Что такое Race condition (состояние гонки)? #
Что такое Race condition #
Race condition — это ошибка конкурентного выполнения, когда результат программы зависит от того, в каком порядке несколько потоков, процессов или async-задач получили доступ к общему ресурсу.
Проще:
Есть общий ресурс
↓
Несколько исполнителей меняют его одновременно
↓
Порядок выполнения непредсказуем
↓
Результат может стать неправильным
Классическое определение: состояние гонки возникает, когда поведение системы зависит от порядка или времени неконтролируемых событий, и из-за этого возможен неожиданный или неправильный результат.
Простой пример #
Допустим, есть общий счётчик:
counter = 0
Два потока одновременно делают:
counter += 1
Кажется, что после двух операций должно быть:
counter == 2
Но фактически может получиться:
counter == 1
Почему? Потому что counter += 1 — это не одна атомарная операция.
Упрощённо она состоит из нескольких шагов:
1. Прочитать counter
2. Прибавить 1
3. Записать новое значение обратно
Как возникает ошибка #
counter = 0
Поток A читает counter → 0
Поток B читает counter → 0
Поток A считает 0 + 1 → 1
Поток B считает 0 + 1 → 1
Поток A записывает counter = 1
Поток B записывает counter = 1
Итог: counter = 1, хотя ожидали 2
Это и есть состояние гонки.
Почему это называется “гонка” #
Потому что несколько исполнителей как будто “соревнуются”, кто первым успеет:
первым прочитать данные
первым изменить данные
первым записать результат
Если результат зависит от этого порядка — есть риск race condition.
Где бывает Race condition #
Race condition может быть не только в потоках.
Она бывает в:
многопоточности
многопроцессности
асинхронном коде
работе с БД
файловой системе
кэше
очередях
распределённых системах
Например:
два запроса одновременно списывают деньги
два worker'а одновременно обрабатывают одну задачу
два пользователя одновременно меняют один документ
две async-задачи одновременно обновляют общий список
Race condition в asyncio #
В asyncio нет настоящего параллельного выполнения Python-кода внутри одного event loop, но race condition всё равно возможна.
Причина — переключение между задачами на await.
import asyncio
counter = 0
async def increment():
global counter
value = counter
await asyncio.sleep(0) # здесь задача отдаёт управление
counter = value + 1
async def main():
await asyncio.gather(
increment(),
increment(),
)
print(counter)
asyncio.run(main())
Ожидаем:
2
Но можно получить:
1
Почему:
Task A читает counter = 0
Task A доходит до await и отдаёт управление
Task B читает counter = 0
Task B доходит до await и отдаёт управление
Task A записывает counter = 1
Task B записывает counter = 1
В документации asyncio.Lock прямо указано, что lock используется для гарантии эксклюзивного доступа к общему ресурсу.
Как исправить в asyncio #
Использовать asyncio.Lock.
import asyncio
counter = 0
lock = asyncio.Lock()
async def increment():
global counter
async with lock:
value = counter
await asyncio.sleep(0)
counter = value + 1
async def main():
await asyncio.gather(
increment(),
increment(),
)
print(counter)
asyncio.run(main())
Теперь критическая секция защищена:
async with lock:
только одна task может менять counter
Race condition в threading #
Для потоков используется threading.Lock.
import threading
counter = 0
lock = threading.Lock()
def increment():
global counter
with lock:
value = counter
counter = value + 1
threading.Lock имеет состояния locked и unlocked, а методы acquire() и release() используются для захвата и освобождения блокировки.
Критическая секция #
Критическая секция — это участок кода, где происходит доступ к общему изменяемому состоянию.
with lock:
counter += 1
Здесь критическая секция:
прочитать counter
изменить counter
записать counter обратно
Главное правило:
Общий изменяемый ресурс
+ конкурентный доступ
+ нет синхронизации
= риск race condition
Как предотвращают Race condition #
Основные способы:
1. Lock / mutex
2. Semaphore
3. Queue
4. Транзакции в БД
5. Атомарные операции
6. Избегание общего изменяемого состояния
7. Immutable-структуры
8. Actor model / message passing
В Python для потоков есть Lock, RLock, Semaphore, Event, Condition и другие примитивы синхронизации в модуле threading.
Race condition vs Deadlock #
Это разные проблемы.
| Проблема | Суть |
|---|---|
| Race condition | Несколько исполнителей одновременно меняют общий ресурс, и результат зависит от порядка выполнения |
| Deadlock | Исполнители заблокировали друг друга и никто не может продолжить работу |
Пример race condition:
два потока одновременно увеличили counter
и потеряли одно обновление
Пример deadlock:
Поток A держит lock_1 и ждёт lock_2
Поток B держит lock_2 и ждёт lock_1
Итог #
Race condition — это ситуация, когда несколько потоков, процессов или async-задач работают с одним общим ресурсом, а итог зависит от случайного порядка их выполнения.
Короткая формула:
Race condition =
shared mutable state
+ concurrent access
+ no proper synchronization
В Python она может быть и в threading, и в multiprocessing, и в asyncio. GIL не отменяет race condition, потому что логическая операция вроде counter += 1 всё равно может состоять из нескольких шагов и требовать синхронизации.
33. Может ли Python-приложение использовать больше одного ядра CPU? #
Да, Python-приложение может использовать больше одного ядра CPU.
Но важно различать:
1. Один Python-процесс + обычные потоки
→ чаще всего не дают параллельного выполнения Python-кода на нескольких ядрах из-за GIL
2. Несколько процессов
→ могут реально использовать несколько ядер
3. Нативные библиотеки
→ могут использовать несколько ядер внутри C/C++/Rust-кода
4. Free-threaded CPython
→ новая сборка CPython без GIL, где потоки могут выполняться параллельно
1. Обычный CPython и потоки #
В стандартном CPython с включённым GIL несколько потоков могут существовать, но только один поток в один момент времени выполняет Python bytecode.
import threading
То есть такая программа может создать несколько потоков:
Thread 1
Thread 2
Thread 3
Thread 4
Но для CPU-bound Python-кода они обычно не будут считать одновременно на 4 ядрах.
Документация Python прямо указывает: GIL ограничивает выгоду от threading для CPU-bound задач, потому что только один поток может выполнять Python bytecode за раз.
2. Для нескольких ядер обычно используют процессы #
Классический способ задействовать несколько ядер в Python — multiprocessing.
from multiprocessing import Pool
def cpu_task(x: int) -> int:
return x * x
if __name__ == "__main__":
with Pool() as pool:
result = pool.map(cpu_task, range(10))
print(result)
Схема:
Процесс 1 → ядро CPU 1
Процесс 2 → ядро CPU 2
Процесс 3 → ядро CPU 3
Процесс 4 → ядро CPU 4
multiprocessing обходит GIL за счёт subprocess вместо потоков и позволяет полноценно использовать несколько процессоров/ядер на машине.
3. Через ProcessPoolExecutor
#
В async-проекте или обычном коде часто используют ProcessPoolExecutor.
import asyncio
from concurrent.futures import ProcessPoolExecutor
def cpu_heavy(n: int) -> int:
total = 0
for i in range(n):
total += i * i
return total
async def main():
loop = asyncio.get_running_loop()
with ProcessPoolExecutor() as pool:
result = await loop.run_in_executor(pool, cpu_heavy, 10_000_000)
print(result)
asyncio.run(main())
Это хороший вариант для CPU-bound задач:
asyncio event loop
↓
отдаёт тяжёлую функцию в ProcessPoolExecutor
↓
функция выполняется в отдельном процессе
↓
отдельный процесс может занять другое ядро CPU
ProcessPoolExecutor использует multiprocessing, что позволяет обходить GIL; при этом выполняемые и возвращаемые объекты должны быть сериализуемыми через pickle.
4. Asyncio сам по себе несколько ядер не использует #
Один asyncio event loop обычно работает в одном потоке.
1 event loop
↓
1 thread
↓
1 CPU core в конкретный момент
asyncio полезен не для CPU-параллелизма, а для I/O-конкурентности:
запрос к БД ждёт ответ
↓
event loop переключается на другую задачу
HTTP-запрос ждёт сеть
↓
event loop переключается на другую задачу
То есть:
asyncio хорошо:
сеть, БД, Redis, WebSocket, HTTP, очереди
asyncio плохо:
тяжёлые вычисления без выноса в процессы
5. Веб-приложение может использовать несколько ядер через workers #
Например, FastAPI/Uvicorn/Gunicorn можно запускать несколькими worker-процессами:
worker 1 → свой процесс → свой event loop → ядро 1
worker 2 → свой процесс → свой event loop → ядро 2
worker 3 → свой процесс → свой event loop → ядро 3
worker 4 → свой процесс → свой event loop → ядро 4
То есть один worker сам по себе обычно не распараллелит CPU-bound Python-код, но несколько worker-процессов могут загрузить несколько ядер.
6. Нативные библиотеки могут использовать несколько ядер #
Даже в одном Python-процессе несколько ядер могут использоваться, если тяжёлая работа выполняется не самим Python bytecode, а нативной библиотекой.
Примеры:
NumPy
OpenCV
PyTorch
TensorFlow
scikit-learn
Такие библиотеки часто выполняют вычисления внутри C/C++/Fortran/Rust-кода и могут отпускать GIL или использовать свои внутренние потоки.
7. Новый вариант: free-threaded CPython #
Начиная с Python 3.13, у CPython появилась поддержка отдельной free-threaded сборки, где GIL отключён. Такая сборка позволяет потокам реально выполняться параллельно на доступных CPU-ядрах.
Но важное уточнение:
обычная сборка CPython
→ GIL включён
free-threaded сборка
→ GIL отключён
→ это отдельный режим/сборка, не обычный дефолтный Python
Глоссарий Python описывает free-threaded build как сборку CPython с поддержкой free threading, которая конфигурируется через --disable-gil.
Итог #
Python-приложение может использовать больше одного ядра CPU, но способ зависит от модели:
| Способ | Использует несколько ядер для Python-кода? | Когда применять |
|---|---|---|
asyncio | обычно нет | I/O-bound задачи |
threading в обычном CPython | обычно нет для CPU-bound | I/O-bound задачи |
multiprocessing | да | CPU-bound задачи |
ProcessPoolExecutor | да | CPU-bound задачи из async/обычного кода |
| несколько web workers | да | веб-приложения |
| NumPy/PyTorch/OpenCV | да, если библиотека так реализована | тяжёлые вычисления |
| free-threaded CPython | да | потоки без GIL, новый отдельный режим |
Короткая формула:
Для I/O-bound:
asyncio или threading
Для CPU-bound:
multiprocessing / ProcessPoolExecutor / несколько worker-процессов
Для обычного CPython:
потоки ≠ полноценное использование нескольких ядер для Python bytecode
34. Как организовать обмен данными между процессами, если процессы не разделяют память напрямую? #
Если процессы не разделяют память напрямую, обмен данными делают через IPC — Inter-Process Communication.
В Python основные варианты:
1. Queue
2. Pipe
3. Manager
4. Shared memory
5. Сокеты
6. Файлы / БД / Redis / RabbitMQ / Kafka
Для обычного Python-кода чаще всего используют multiprocessing.Queue или multiprocessing.Pipe.
1. multiprocessing.Queue
#
Queue — самый удобный способ передавать данные между процессами.
from multiprocessing import Process, Queue
def worker(q: Queue):
q.put("hello from child process")
if __name__ == "__main__":
q = Queue()
p = Process(target=worker, args=(q,))
p.start()
message = q.get()
print(message)
p.join()
Схема:
parent process
↓ q.get()
Queue
child process
↑ q.put(...)
Важно: объект, который кладётся в multiprocessing.Queue, сериализуется. То есть он превращается в байты, передаётся другому процессу, а там восстанавливается обратно. В документации Python указано, что очереди multiprocessing являются thread-safe и process-safe, а объекты, помещённые в очередь, сериализуются.
Подходит для:
передачи задач worker-процессам
получения результатов
producer/consumer схемы
фоновой обработки
2. multiprocessing.Pipe
#
Pipe — это канал связи между двумя процессами.
from multiprocessing import Process, Pipe
def worker(conn):
conn.send("hello from child")
conn.close()
if __name__ == "__main__":
parent_conn, child_conn = Pipe()
p = Process(target=worker, args=(child_conn,))
p.start()
message = parent_conn.recv()
print(message)
p.join()
Схема:
process A ←──────── Pipe ────────→ process B
Pipe() возвращает пару соединённых объектов connection. По умолчанию канал двусторонний, то есть оба конца могут отправлять и принимать данные.
Pipe удобен, когда:
процессов мало
нужна связь 1 к 1
нужен простой запрос-ответ между parent и child
3. multiprocessing.Manager
#
Manager создаёт отдельный серверный процесс, который хранит Python-объекты, а другие процессы работают с ними через proxy-объекты.
from multiprocessing import Process, Manager
def worker(shared_list):
shared_list.append("item from child")
if __name__ == "__main__":
with Manager() as manager:
shared_list = manager.list()
p = Process(target=worker, args=(shared_list,))
p.start()
p.join()
print(list(shared_list))
Схема:
process A ── proxy ──┐
↓
manager process
↑
process B ── proxy ──┘
Плюс: удобно работать с list, dict, Queue, Lock и другими объектами.
Минус: медленнее, чем Queue или Pipe, потому что доступ идёт через proxy и отдельный manager-процесс.
4. Shared memory #
Формально процессы не разделяют обычную память Python-объектов, но можно специально создать участок общей памяти.
Для этого есть multiprocessing.shared_memory.
from multiprocessing import Process
from multiprocessing import shared_memory
def worker(name):
shm = shared_memory.SharedMemory(name=name)
shm.buf[0] = 100
shm.close()
if __name__ == "__main__":
shm = shared_memory.SharedMemory(create=True, size=10)
p = Process(target=worker, args=(shm.name,))
p.start()
p.join()
print(shm.buf[0])
shm.close()
shm.unlink()
Схема:
process A ─────┐
↓
shared memory block
↑
process B ─────┘
multiprocessing.shared_memory предоставляет SharedMemory для выделения и управления общей памятью, доступной нескольким процессам.
Подходит для:
больших массивов данных
числовых данных
избежания лишнего копирования
NumPy-массивов между процессами
Но это не обычный Python-объект вроде списка или словаря. Там нужно самостоятельно контролировать формат данных и синхронизацию.
5. Value и Array
#
Для простых общих значений можно использовать multiprocessing.Value и multiprocessing.Array.
from multiprocessing import Process, Value
def increment(counter):
with counter.get_lock():
counter.value += 1
if __name__ == "__main__":
counter = Value("i", 0)
processes = [
Process(target=increment, args=(counter,))
for _ in range(10)
]
for p in processes:
p.start()
for p in processes:
p.join()
print(counter.value)
Здесь Value("i", 0) создаёт общее целочисленное значение.
Важно: при изменении общего значения нужна синхронизация через lock, иначе можно получить race condition.
6. Внешний брокер или хранилище #
В реальных backend-системах процессы часто обмениваются данными не напрямую, а через внешний сервис:
Redis
RabbitMQ
Kafka
PostgreSQL
файлы
HTTP API
Unix socket / TCP socket
Например:
process A
↓
кладёт задачу в Redis / RabbitMQ
process B
↓
забирает задачу и обрабатывает
Это особенно полезно, когда процессы находятся:
в разных Docker-контейнерах
на разных серверах
в разных worker-приложениях
7. Что выбрать #
| Задача | Что использовать |
|---|---|
| Передать задачи worker-процессам | multiprocessing.Queue |
| Получить результаты от worker’ов | multiprocessing.Queue |
| Связь 1 к 1 между двумя процессами | multiprocessing.Pipe |
Общий dict / list между процессами | multiprocessing.Manager |
| Большие массивы без копирования | multiprocessing.shared_memory |
| Простое общее число | multiprocessing.Value |
| Простая общая таблица чисел | multiprocessing.Array |
| Обмен между сервисами / контейнерами | Redis / RabbitMQ / Kafka / БД / HTTP |
Главная разница #
Queue / Pipe
→ данные передаются сообщениями
→ обычно сериализация и копирование
Manager
→ общий объект через proxy
→ удобно, но медленнее
Shared memory
→ общий участок памяти
→ быстрее для больших данных, но сложнее
Redis / RabbitMQ / Kafka
→ обмен между независимыми процессами/сервисами
→ хорошо для backend-архитектуры
Итог #
Процессы не видят обычную память друг друга, поэтому обмен делают через специальные механизмы:
message passing:
Queue, Pipe, sockets, broker
shared state through proxy:
Manager
explicit shared memory:
shared_memory, Value, Array
external storage:
DB, Redis, файлы
Для большинства задач в Python начинай с multiprocessing.Queue.
Для CPU-bound worker’ов:
parent process
↓ кладёт задачи в Queue
worker processes
↓ считают
result Queue
↓ возвращают результат
Для больших числовых данных лучше смотреть в сторону shared_memory, чтобы не копировать огромные объекты между процессами.
35. Можно ли использовать общую память в multiprocessing без внешних хранилищ? #
Да, можно. В multiprocessing есть встроенные механизмы общей памяти, без Redis, БД, файлов, очередей сообщений и других внешних хранилищ.
Основные варианты #
1. multiprocessing.Value
#
Для одного общего значения:
from multiprocessing import Process, Value
def increment(counter):
for _ in range(100_000):
with counter.get_lock():
counter.value += 1
if __name__ == "__main__":
counter = Value("i", 0)
p1 = Process(target=increment, args=(counter,))
p2 = Process(target=increment, args=(counter,))
p1.start()
p2.start()
p1.join()
p2.join()
print(counter.value)
Value создаёт ctypes-объект в общей памяти. По умолчанию он оборачивается синхронизирующей обёрткой с lock. Но операция counter.value += 1 не атомарна, поэтому для корректного инкремента нужен get_lock(). Это прямо указано в документации Python.
2. multiprocessing.Array
#
Для массива фиксированной длины:
from multiprocessing import Process, Array
def worker(arr):
arr[0] = 100
arr[1] = 200
if __name__ == "__main__":
arr = Array("i", [1, 2, 3])
p = Process(target=worker, args=(arr,))
p.start()
p.join()
print(list(arr)) # [100, 200, 3]
Array создаёт массив в общей памяти. По умолчанию доступ синхронизируется lock-ом, но если передать lock=False, доступ уже не будет автоматически process-safe.
3. multiprocessing.shared_memory.SharedMemory
#
Это более низкоуровневый и гибкий вариант. Он появился в Python 3.8 и позволяет создавать блок общей памяти, к которому могут подключаться разные процессы по имени. ( Python documentation)
from multiprocessing import Process
from multiprocessing import shared_memory
def worker(shm_name):
shm = shared_memory.SharedMemory(name=shm_name)
shm.buf[0] = 99
shm.close()
if __name__ == "__main__":
shm = shared_memory.SharedMemory(create=True, size=10)
shm.buf[0] = 1
p = Process(target=worker, args=(shm.name,))
p.start()
p.join()
print(shm.buf[0]) # 99
shm.close()
shm.unlink()
Здесь shm.buf — это memoryview на общий блок памяти. Один процесс меняет байты, другой видит изменения. После работы важно вызвать close(), а когда память больше не нужна — unlink(), иначе можно получить утечку ресурса.
Важное ограничение #
Общая память не означает, что обычные Python-объекты автоматически становятся общими.
Например, нельзя просто положить в shared memory обычный dict, list или объект класса и безопасно менять его как общий объект. Общая память хранит байты / ctypes-структуры / массивы. Для сложных Python-объектов обычно используют:
multiprocessing.Manager()
Queue / Pipe
сериализацию
Redis / БД / другое внешнее хранилище
Manager тоже встроен в multiprocessing, но это не настоящая общая память. Он запускает отдельный серверный процесс, который хранит объекты, а другие процессы работают с ними через proxy.
Когда что использовать #
Value
→ один общий счётчик, флаг, число
Array
→ небольшой массив фиксированной длины
shared_memory.SharedMemory
→ большие данные, бинарные буферы, NumPy-массивы, высокая производительность
Manager
→ удобно шарить dict/list, но медленнее и это proxy, а не прямой доступ к памяти
Короткий вывод #
Да, в multiprocessing можно использовать общую память без внешних хранилищ.
Основные инструменты:
multiprocessing.Value
multiprocessing.Array
multiprocessing.shared_memory.SharedMemory
multiprocessing.sharedctypes
Но нужно помнить: процессы всё равно имеют разные адресные пространства, поэтому shared memory надо создавать явно, а при конкурентной записи защищать доступ через Lock, иначе возможны race condition.
36. Как решать проблемы race condition в Python без GIL? #
Race condition решается не GIL-ом, а правильной синхронизацией доступа к общему состоянию.
В Python без GIL, например в free-threaded CPython, потоки могут реально выполняться параллельно на разных ядрах. Начиная с Python 3.13 есть сборка CPython с отключаемым GIL, а в Python 3.14 она продолжает развиваться. Встроенные типы вроде dict, list, set имеют внутренние защиты от некоторых конкурентных модификаций, но документация прямо рекомендует использовать threading.Lock и другие примитивы синхронизации, а не полагаться на внутренние детали реализации.
1. Защищать критическую секцию через Lock #
Проблемный код:
counter += 1
Это не одна операция на уровне логики. Это примерно:
прочитать counter
прибавить 1
записать обратно
Два потока могут одновременно прочитать старое значение и перезаписать результат друг друга.
Правильно:
import threading
counter = 0
lock = threading.Lock()
def increment():
global counter
for _ in range(100_000):
with lock:
counter += 1
threading.Lock — базовый примитив синхронизации. Когда один поток захватил lock, остальные блокируются до его освобождения. Объекты Lock, RLock, Condition, Semaphore можно использовать через with, что безопаснее, чем вручную вызывать acquire() / release().
2. Делать shared state минимальным #
Плохой подход:
users = {}
def update_user(user_id, data):
users[user_id]["balance"] += data["amount"]
Здесь общий dict, вложенные структуры и изменение баланса. Без синхронизации это легко приводит к race condition.
Лучше:
users = {}
users_lock = threading.Lock()
def update_user(user_id, amount):
with users_lock:
user = users[user_id]
user["balance"] += amount
Ещё лучше — не давать разным потокам менять одну и ту же структуру напрямую.
3. Использовать очереди вместо общей изменяемой памяти #
Часто лучший способ убрать race condition — не делить память, а передавать сообщения.
from queue import Queue
from threading import Thread
queue = Queue()
def worker():
while True:
item = queue.get()
try:
process(item)
finally:
queue.task_done()
def process(item):
print(item)
Thread(target=worker, daemon=True).start()
queue.put("task-1")
queue.put("task-2")
queue.join()
Идея:
не несколько потоков меняют один объект
а один поток владеет состоянием
остальные отправляют ему задачи
Это обычно проще и безопаснее, чем ставить lock вокруг каждой операции.
4. Для multiprocessing использовать multiprocessing.Lock #
Для процессов GIL вообще не решает race condition, потому что процессы имеют разные интерпретаторы и разные адресные пространства. Если используется shared memory, Value, Array или sharedctypes, синхронизация всё равно нужна.
Пример:
from multiprocessing import Process, Value
def worker(counter):
for _ in range(100_000):
with counter.get_lock():
counter.value += 1
if __name__ == "__main__":
counter = Value("i", 0)
p1 = Process(target=worker, args=(counter,))
p2 = Process(target=worker, args=(counter,))
p1.start()
p2.start()
p1.join()
p2.join()
print(counter.value)
Документация multiprocessing прямо указывает, что операции вида += над shared value не атомарны, потому что включают чтение и запись. Для атомарного инкремента нужно использовать lock, например with counter.get_lock():.
5. Использовать правильный примитив под задачу #
Lock
→ защитить одну критическую секцию
RLock
→ нужен повторный захват lock тем же потоком
Semaphore
→ ограничить количество одновременных доступов к ресурсу
Condition
→ ждать изменения состояния
Event
→ один поток сообщает другим, что событие произошло
Queue
→ передавать задачи/данные без прямого общего состояния
Condition всегда связан с lock-ом: поток ждёт изменения состояния, временно отпускает lock, затем снова захватывает его после пробуждения. Это типичный механизм для producer-consumer сценариев.
6. Не полагаться на “атомарность” встроенных структур #
Даже если отдельная операция с list или dict сейчас выглядит безопасной, нельзя строить на этом логику.
Плохо:
if key not in cache:
cache[key] = load_value()
Между проверкой и записью другой поток может уже добавить этот ключ.
Правильно:
cache = {}
cache_lock = threading.Lock()
def get_or_load(key):
with cache_lock:
if key not in cache:
cache[key] = load_value(key)
return cache[key]
Главное правило #
Любая операция вида:
прочитать → проверить → изменить → записать
должна быть защищена синхронизацией,
если объект доступен из нескольких потоков или процессов.
GIL мог случайно скрывать часть проблем, но он не был нормальным механизмом защиты бизнес-логики. В коде, который должен корректно работать без GIL, нужно проектировать доступ к данным явно: Lock, Queue, Condition, immutable-данные, ownership-модель или изоляция состояния.
37. Что такое deadlock и какие есть методы решения этой проблемы? #
Что такое deadlock #
deadlock — это взаимная блокировка: несколько потоков, процессов или задач ждут друг друга, и никто не может продолжить выполнение.
Классический пример:
Thread A:
захватил lock_1
ждёт lock_2
Thread B:
захватил lock_2
ждёт lock_1
В итоге:
A не отпускает lock_1, потому что ждёт lock_2
B не отпускает lock_2, потому что ждёт lock_1
Оба зависают.
Условия возникновения deadlock #
В теории ОС обычно выделяют 4 условия Коффмана:
1. Mutual exclusion
Ресурс может быть занят только одним потоком/процессом.
2. Hold and wait
Поток уже держит один ресурс и ждёт другой.
3. No preemption
Ресурс нельзя насильно отобрать, его должен освободить владелец.
4. Circular wait
Есть цикл ожидания: A ждёт B, B ждёт C, C ждёт A.
Deadlock становится возможным, когда эти условия выполняются одновременно. Это стандартная модель объяснения взаимных блокировок в ОС.
Пример deadlock в Python #
import threading
import time
lock_a = threading.Lock()
lock_b = threading.Lock()
def worker_1():
with lock_a:
time.sleep(0.1)
with lock_b:
print("worker_1 done")
def worker_2():
with lock_b:
time.sleep(0.1)
with lock_a:
print("worker_2 done")
t1 = threading.Thread(target=worker_1)
t2 = threading.Thread(target=worker_2)
t1.start()
t2.start()
t1.join()
t2.join()
Проблема:
worker_1 держит lock_a и ждёт lock_b
worker_2 держит lock_b и ждёт lock_a
threading.Lock блокирует поток, если lock уже занят другим потоком. Поэтому такой код может зависнуть навсегда.
Как решать deadlock #
1. Всегда брать lock-и в одном порядке #
Это главный практический способ.
Плохо:
# Поток 1
with lock_a:
with lock_b:
...
# Поток 2
with lock_b:
with lock_a:
...
Правильно:
# Все потоки берут lock-и только так:
with lock_a:
with lock_b:
...
То есть в проекте должен быть единый порядок захвата:
lock_user
↓
lock_account
↓
lock_transaction
А не где-то наоборот.
2. Не держать lock дольше, чем нужно #
Плохо:
with lock:
data = load_from_network()
result = heavy_calculation(data)
shared_state["result"] = result
Здесь lock удерживается во время сетевого запроса и тяжёлого вычисления.
Лучше:
data = load_from_network()
result = heavy_calculation(data)
with lock:
shared_state["result"] = result
Правило:
Внутри lock — только работа с общим состоянием.
Всё медленное — вне lock.
3. Использовать with, а не ручной acquire() / release()
#
Плохо:
lock.acquire()
do_something()
lock.release()
Если внутри do_something() будет исключение, release() может не выполниться.
Лучше:
with lock:
do_something()
В документации Python прямо рекомендуется использовать lock-и как context manager, то есть через with, когда это возможно.
Альтернатива вручную:
lock.acquire()
try:
do_something()
finally:
lock.release()
4. Использовать timeout при захвате lock #
Можно не ждать бесконечно:
acquired = lock.acquire(timeout=2)
if acquired:
try:
do_something()
finally:
lock.release()
else:
print("Не удалось получить lock")
Такой подход не всегда полностью решает проблему, но помогает не зависать навсегда и даёт шанс обработать ситуацию.
5. Использовать RLock, если один и тот же поток должен брать lock повторно
#
Обычный Lock не является реентерабельным. Если тот же поток попытается захватить его второй раз, он сам себя заблокирует.
Проблемный пример:
lock = threading.Lock()
def outer():
with lock:
inner()
def inner():
with lock:
print("inner")
outer() уже захватил lock, потом вызывает inner(), а inner() снова пытается взять тот же lock.
Для таких случаев используют RLock:
lock = threading.RLock()
def outer():
with lock:
inner()
def inner():
with lock:
print("inner")
RLock можно захватывать несколько раз из одного и того же потока, но освобождать его нужно столько же раз. Документация Python указывает, что несоответствие количества acquire() и release() может привести к deadlock.
6. Не вызывать join() на самом себе
#
Поток не может ждать завершения самого себя.
threading.current_thread().join()
Это логически невозможно:
поток ждёт, пока сам завершится
но он не может завершиться, потому что ждёт
Python прямо запрещает join() текущего потока, потому что это вызвало бы deadlock.
То же самое относится и к процессам: процесс не может сделать join() самого себя.
7. Осторожно с multiprocessing.Queue
#
В multiprocessing deadlock может возникнуть, если процесс положил данные в Queue, а родительский процесс сделал join() до того, как эти данные были прочитаны.
Проблемная схема:
child process кладёт много данных в Queue
parent делает child.join()
Queue не очищена
child не может завершиться
parent ждёт child
Документация Python предупреждает: если используется multiprocessing.Queue, нужно убедиться, что все элементы из очереди будут извлечены до join(). Иначе можно получить deadlock.
8. Не ждать результат задачи внутри той же маленькой executor-группы #
Пример проблемы:
from concurrent.futures import ThreadPoolExecutor
def task():
future = executor.submit(lambda: 123)
return future.result()
executor = ThreadPoolExecutor(max_workers=1)
future = executor.submit(task)
print(future.result())
Здесь один worker занят task(), а внутри task() ожидается другая задача. Но свободного worker-а нет.
worker ждёт новую задачу
новая задача ждёт свободный worker
свободного worker нет
В документации concurrent.futures есть предупреждение, что ожидание Future внутри задач executor-а может привести к deadlock; для ProcessPoolExecutor вызов методов Executor или Future из выполняемой задачи прямо указан как причина deadlock.
9. Использовать очереди и модель владельца состояния #
Часто лучше не делить один объект между потоками, а сделать один поток владельцем состояния.
Плохо:
много потоков напрямую меняют общий dict
Лучше:
один поток владеет dict
остальные отправляют ему команды через Queue
Пример:
from queue import Queue
from threading import Thread
commands = Queue()
state = {}
def state_owner():
while True:
key, value = commands.get()
try:
state[key] = value
finally:
commands.task_done()
Thread(target=state_owner, daemon=True).start()
commands.put(("user_id", 123))
commands.join()
Так меньше lock-ов — меньше риска deadlock.
10. Логировать зависания и диагностировать #
Практические признаки deadlock:
программа не падает, но зависла
CPU почти не используется
потоки висят в ожидании lock/join/result
логи остановились на одном месте
Для диагностики в Python обычно помогают:
import faulthandler
faulthandler.dump_traceback_later(10, repeat=True)
Он будет периодически печатать stack trace потоков. По stack trace часто видно, где поток завис: на lock.acquire(), join(), future.result() или ожидании очереди.
Главное правило #
Deadlock решают не тем, что “добавляют больше lock-ов”, а наоборот — уменьшают количество мест, где код может заблокироваться.
Практический чеклист:
1. Все lock-и брать в одном порядке.
2. Держать lock минимальное время.
3. Использовать with lock.
4. Не делать I/O внутри lock.
5. Использовать timeout там, где зависание критично.
6. Использовать RLock только когда реально нужна реентерабельность.
7. Не ждать Future/result внутри того же маленького pool-а.
8. В multiprocessing сначала читать Queue, потом делать join.
9. По возможности заменять shared state на Queue/message passing.
Самая частая причина deadlock в реальном коде — нарушение порядка захвата lock-ов. Самый надёжный практический метод — заранее определить и строго соблюдать единый порядок блокировок.
38. Что такое Семафор #
Семафоры служат для управления доступом к общим ресурсам или для ограничения количества одновременно выполняемых задач. Они позволяют избежать перегрузки системы и управлять параллельным выполнением задач, обеспечивая более эффективное использование ресурсов.
Как работают семафоры в asyncio:
- Создание: Семафор создается с указанием максимального количества задач, которые могут одновременно выполняться. Например, semaphore = asyncio.Semaphore(2) создаст семафор, который позволяет одновременно работать только двум задачам.
- Блокировка и освобождение: Перед выполнением асинхронной задачи, требующей доступа к ресурсу или имеющей ограничение по количеству одновременных запусков, необходимо “захватить” семафор, используя конструкцию async with semaphore:. Это заблокирует семафор, уменьшив его счетчик, и позволит задаче выполняться.
- Ожидание освобождения: Если семафор уже занят (счетчик равен нулю), задача будет ожидать, пока другая задача не освободит его (выполнит release(), увеличив счетчик). После освобождения семафора, задача сможет его захватить и продолжить выполнение.
- Освобождение после выполнения: После завершения выполнения асинхронной задачи, семафор автоматически освобождается (если использовалась конструкция async with), и счетчик семафора увеличивается.
Пример использования:
import asyncio
async def worker(name, semaphore):
print(f'Задача {name}: ожидает семафор')
async with semaphore:
print(f'Задача {name}: захватила семафор')
await asyncio.sleep(1)
print(f'Задача {name}: освобождает семафор')
async def main():
semaphore = asyncio.Semaphore(2)
tasks = [worker(i, semaphore) for i in range(5)]
await asyncio.gather(*tasks)
if __name__ == "__main__":
asyncio.run(main())
В этом примере, несмотря на то, что создано 5 задач, только 2 из них могут выполняться одновременно из-за установленного в семафоре лимита в 2. Остальные задачи будут ожидать, пока одна из двух работающих задач не завершится и не освободит семафор.
Преимущества использования семафоров:
Управление ресурсами: Семафоры позволяют контролировать доступ к ресурсам, которые могут быть ограничены по количеству одновременно использующих их задач (например, соединения с базой данных, файлы).
Ограничение параллелизма: Семафоры позволяют ограничить количество одновременно выполняемых задач, предотвращая перегрузку системы и обеспечивая более стабильную работу.
Синхронизация задач: Семафоры могут использоваться для координации выполнения задач, например, для обеспечения того, чтобы определенная задача не начиналась, пока не завершится другая.
Простота использования: Semaphore предоставляет удобный и интуитивно понятный способ управления параллелизмом
39. Как задать timeout для asyncio задачи? #
Основной способ: asyncio.wait_for()
#
Для одной coroutine/задачи чаще всего используют asyncio.wait_for():
import asyncio
async def fetch_data():
await asyncio.sleep(10)
return "done"
async def main():
try:
result = await asyncio.wait_for(fetch_data(), timeout=3)
print(result)
except TimeoutError:
print("Задача не успела выполниться за 3 секунды")
asyncio.run(main())
asyncio.wait_for(aw, timeout) ждёт awaitable не дольше указанного количества секунд. Если время истекло, задача отменяется и выбрасывается TimeoutError. Это описано в официальной документации Python.
Если задача уже создана через create_task()
#
import asyncio
async def worker():
await asyncio.sleep(10)
return "ok"
async def main():
task = asyncio.create_task(worker())
try:
result = await asyncio.wait_for(task, timeout=3)
print(result)
except TimeoutError:
print("Timeout")
print(task.cancelled()) # True, если задача реально отменилась
asyncio.run(main())
Важно: при timeout wait_for() вызывает отмену задачи. Но отмена в asyncio не мгновенная: coroutine должна дойти до ближайшего await, где будет выброшен CancelledError.
Современный способ: asyncio.timeout() в Python 3.11+
#
Начиная с Python 3.11 можно задавать timeout на блок кода:
import asyncio
async def main():
try:
async with asyncio.timeout(3):
await asyncio.sleep(10)
print("completed")
except TimeoutError:
print("Блок не успел выполниться за 3 секунды")
asyncio.run(main())
Это удобно, когда нужно ограничить по времени не одну функцию, а несколько операций внутри блока. asyncio.timeout() добавлен в Python 3.11.
Timeout для нескольких задач #
import asyncio
async def task_1():
await asyncio.sleep(2)
return "task_1"
async def task_2():
await asyncio.sleep(5)
return "task_2"
async def main():
try:
result = await asyncio.wait_for(
asyncio.gather(task_1(), task_2()),
timeout=3,
)
print(result)
except TimeoutError:
print("Не все задачи успели выполниться")
asyncio.run(main())
Здесь gather() запускает несколько awaitable, а wait_for() ограничивает общее время ожидания.
Если нужно НЕ отменять задачу при timeout #
Иногда нужно ограничить только ожидание, но не убивать саму задачу. Тогда используют asyncio.shield():
import asyncio
async def long_task():
await asyncio.sleep(5)
return "done"
async def main():
task = asyncio.create_task(long_task())
try:
await asyncio.wait_for(asyncio.shield(task), timeout=2)
except TimeoutError:
print("Ожидание истекло, но задача продолжает работать")
result = await task
print(result)
asyncio.run(main())
Кратко #
await asyncio.wait_for(coro(), timeout=5)
или в Python 3.11+:
async with asyncio.timeout(5):
await coro()
Выбирай так:
wait_for() → timeout на конкретную coroutine/task
asyncio.timeout() → timeout на блок async-кода, Python 3.11+
shield() → timeout ожидания без отмены самой задачи
40. Что такое hyper threading (гипер потоки)? #
Что такое Hyper-Threading #
Hyper-Threading — это технология Intel, при которой одно физическое ядро процессора представляется операционной системе как два логических процессора.
То есть ОС видит не только реальные ядра, но и дополнительные “виртуальные” потоки выполнения.
Пример:
4 физических ядра + Hyper-Threading
↓
8 логических потоков / logical processors
Intel описывает Hyper-Threading как аппаратную технологию, позволяющую запускать больше одного потока на одном ядре. Больше потоков — больше работы может выполняться параллельно.
Простая аналогия #
Физическое ядро — это как один работник.
Без Hyper-Threading:
1 ядро → 1 поток задач
С Hyper-Threading:
1 ядро → 2 потока задач
Но это не значит, что одно ядро превращается в два полноценных ядра. Второй поток использует свободные ресурсы того же самого физического ядра.
Как это работает внутри #
У ядра процессора есть разные блоки:
ядро CPU
├─ блоки выполнения инструкций
├─ кэш
├─ регистры
├─ планировщик инструкций
└─ другие внутренние ресурсы
Когда один поток простаивает, например ждёт данные из памяти, второй поток может использовать часть свободных ресурсов ядра.
Смысл такой:
Поток A ждёт данные
↓
CPU не простаивает полностью
↓
Поток B использует свободные ресурсы
В серверной документации Intel также указывается, что с Hyper-Threading каждое физическое ядро может выполнять два потока, которые ОС видит как логические процессоры.
Hyper-Threading vs физические ядра #
Важно:
8 потоков ≠ 8 полноценных ядер
Например:
4 ядра / 8 потоков
Это значит:
4 реальных физических ядра
8 логических процессоров для ОС
Производительность обычно будет выше, чем у 4 ядра / 4 потока, но ниже, чем у настоящих 8 физических ядер.
Где помогает #
Hyper-Threading полезен, когда программа может распараллеливать работу:
рендеринг
архивация
компиляция кода
виртуальные машины
серверные приложения
браузер с большим числом вкладок
многопоточные вычисления
В таких задачах CPU может лучше загружать свои внутренние ресурсы.
Где почти не помогает #
Если программа в основном использует один поток, Hyper-Threading почти не даст прироста:
старые игры
однопоточные скрипты
часть Python-кода под GIL
простые офисные задачи
Например, если программа умеет нормально использовать только одно ядро, дополнительные логические потоки не сделают её в 2 раза быстрее.
Важно для Python #
В CPython обычный Python-код ограничен GIL, поэтому два Python-потока не выполняют Python-байткод параллельно на разных ядрах.
То есть Hyper-Threading не сделает обычный многопоточный CPU-bound Python-код в 2 раза быстрее.
Пример CPU-bound задачи:
def heavy_calc():
total = 0
for i in range(100_000_000):
total += i
return total
Для такого в Python чаще используют:
multiprocessing
ProcessPoolExecutor
расширения на C/Rust
NumPy, если вычисления уходят в нативный код
А вот для I/O-bound задач потоки всё ещё полезны:
запросы в сеть
работа с файлами
ожидание базы данных
Кратко #
Hyper-Threading — это не дополнительные физические ядра.
Это возможность одному физическому ядру обрабатывать 2 логических потока.
Главная идея:
меньше простоев ядра
лучше загрузка CPU
выше производительность в многопоточных задачах
Но:
2 потока на ядро ≠ x2 производительности
Обычно прирост зависит от задачи: где-то заметный, где-то почти нулевой.
41. Что такое главный поток в Python и как он влияет на дополнительные потоки? #
Что такое главный поток #
Главный поток в Python — это поток, в котором изначально запускается программа.
Когда ты запускаешь файл:
python main.py
код начинает выполняться в главном потоке:
print("start")
А дополнительные потоки создаются уже из него:
import threading
def worker():
print("Работа в отдельном потоке")
thread = threading.Thread(target=worker)
thread.start()
print("Главный поток продолжает работу")
В threading есть функция threading.main_thread(), которая возвращает объект главного потока программы. Документация также указывает, что threading.current_thread() возвращает текущий поток выполнения.
Главный поток не является “командиром” всех потоков #
Важно: главный поток не управляет каждым шагом дополнительных потоков напрямую.
После запуска:
thread.start()
дополнительный поток начинает выполняться отдельно, а планированием занимается ОС и интерпретатор Python.
Упрощённо:
main thread
├─ запускает thread_1
├─ запускает thread_2
└─ продолжает свой код
thread_1
└─ выполняет свою функцию
thread_2
└─ выполняет свою функцию
Главный поток может:
создать поток
запустить поток
ждать поток через join()
передать данные
использовать Lock/Event/Queue для синхронизации
Но он не “крутит” дополнительные потоки вручную.
Как главный поток влияет на завершение программы #
Если главный поток завершился, это не всегда значит, что процесс сразу завершится.
Зависит от типа дополнительных потоков:
обычные потоки daemon=False → программа ждёт их завершения
daemon-потоки daemon=True → программа может завершиться, не дожидаясь их
Пример с обычным потоком:
import threading
import time
def worker():
time.sleep(3)
print("worker завершился")
thread = threading.Thread(target=worker)
thread.start()
print("main завершил свой код")
Вывод будет примерно такой:
main завершил свой код
worker завершился
Хотя главный поток дошёл до конца, Python-процесс продолжил жить, потому что обычный поток ещё работал.
Пример с daemon=True:
import threading
import time
def worker():
time.sleep(3)
print("worker завершился")
thread = threading.Thread(target=worker, daemon=True)
thread.start()
print("main завершил свой код")
Здесь программа может завершиться до того, как worker успеет вывести сообщение.
join() — главный поток ждёт другой поток
#
Чтобы главный поток явно дождался завершения другого потока, используют join():
import threading
import time
def worker():
time.sleep(2)
print("worker done")
thread = threading.Thread(target=worker)
thread.start()
print("main ждёт worker")
thread.join()
print("main продолжает работу после worker")
Схема:
main thread
↓
создал worker
↓
thread.start()
↓
thread.join()
↓
main заблокирован, пока worker не завершится
↓
main продолжает выполнение
Главное ограничение: GIL #
В стандартном CPython есть GIL — Global Interpreter Lock.
Из-за него в обычной сборке CPython только один поток может выполнять Python-байткод в один момент времени. Поэтому потоки в CPython не дают полноценного параллельного ускорения для CPU-bound задач. Официальная документация прямо указывает, что для CPU-bound задач обычно лучше использовать multiprocessing или ProcessPoolExecutor, а threading хорошо подходит для I/O-bound задач.
Пример CPU-bound:
def calculate():
total = 0
for i in range(100_000_000):
total += i
Для такого два потока обычно не дадут честного ускорения x2.
Пример I/O-bound:
import requests
response = requests.get("https://example.com")
Здесь поток может ждать сеть, а другой поток в это время может выполнять другую работу.
Сигналы обычно обрабатываются в главном потоке #
Ещё одна важная особенность: Python-обработчики сигналов выполняются в главном Python-потоке, даже если сигнал был получен другим потоком. Также только главный поток может устанавливать новый обработчик сигнала.
Например, это важно для:
Ctrl+C / KeyboardInterrupt
signal.signal(...)
graceful shutdown
Поэтому часто главный поток отвечает за корректное завершение программы, а worker-потоки получают команду остановиться через Event, Queue или общий флаг с синхронизацией.
Правильная остановка worker-потока #
Плохой подход — пытаться насильно убить поток.
Обычный подход:
import threading
import time
stop_event = threading.Event()
def worker():
while not stop_event.is_set():
print("working...")
time.sleep(1)
print("worker stopped")
thread = threading.Thread(target=worker)
thread.start()
time.sleep(3)
stop_event.set()
thread.join()
print("main finished")
Смысл:
main thread
↓
создаёт Event
↓
запускает worker
↓
через некоторое время вызывает stop_event.set()
↓
worker сам замечает сигнал остановки
↓
main ждёт worker через join()
Кратко #
Главный поток — поток, с которого начинается выполнение Python-программы.
Он влияет на дополнительные потоки так:
создаёт и запускает их
может ждать их через join()
может координировать их через Lock/Event/Queue
обычно обрабатывает сигналы
его завершение влияет на жизнь daemon-потоков
Но:
он не управляет выполнением потоков вручную
не делает CPU-bound Python-код реально параллельным из-за GIL
не должен насильно убивать worker-потоки
Упрощённая картина:
main thread
├─ управляет запуском и завершением
├─ обрабатывает сигналы
├─ координирует worker-потоки
└─ сам тоже является обычным потоком выполнения
worker threads
├─ выполняют свои функции
├─ делят память процесса
├─ конкурируют за GIL при выполнении Python-кода
└─ должны синхронизироваться при работе с общими данными
42. Как дождаться завершения всех потоков в Python? | threading.Thread.join() как работает? #
Чтобы дождаться завершения всех потоков, нужно сохранить объекты Thread в список, запустить их через start(), а потом пройтись по ним и вызвать join():
import threading
import time
def worker(number: int):
print(f"Поток {number} начал работу")
time.sleep(2)
print(f"Поток {number} завершился")
threads = []
for i in range(5):
thread = threading.Thread(target=worker, args=(i,))
threads.append(thread)
thread.start()
for thread in threads:
thread.join()
print("Все потоки завершились")
Схема:
main thread
↓
создал 5 потоков
↓
запустил 5 потоков через start()
↓
вызвал join() для каждого
↓
дождался завершения всех
↓
продолжил выполнение
Метод join() блокирует вызывающий поток до тех пор, пока поток, у которого вызвали join(), не завершится. В официальной документации Python это описано как ожидание завершения потока: нормально, через необработанное исключение или до истечения timeout.
Как работает thread.join()
#
Пример:
thread.start()
thread.join()
Означает:
запусти поток
↓
останови текущий поток выполнения здесь
↓
жди, пока thread завершится
↓
после этого продолжай код дальше
Важно: join() не останавливает поток и не “убивает” его.
Он только ждёт.
thread.start() → запускает поток
thread.join() → ждёт завершения потока
Почему join() обычно вызывают отдельным циклом
#
Правильно:
threads = []
for i in range(5):
thread = threading.Thread(target=worker, args=(i,))
threads.append(thread)
thread.start()
for thread in threads:
thread.join()
Так все потоки сначала запускаются, а потом главный поток ждёт их завершения.
Неправильно, если нужна параллельность:
for i in range(5):
thread = threading.Thread(target=worker, args=(i,))
thread.start()
thread.join()
Здесь каждый поток будет запускаться и сразу ожидаться:
запустил поток 1 → дождался
запустил поток 2 → дождался
запустил поток 3 → дождался
То есть параллельность почти пропадает.
join() с timeout
#
Можно ждать поток не бесконечно, а ограниченное время:
thread.join(timeout=3)
Но важный момент: join() всегда возвращает None. Поэтому, чтобы понять, завершился поток или timeout истёк, нужно проверить is_alive():
import threading
import time
def worker():
time.sleep(10)
thread = threading.Thread(target=worker)
thread.start()
thread.join(timeout=3)
if thread.is_alive():
print("Поток всё ещё работает")
else:
print("Поток завершился")
Документация Python отдельно указывает: если используется timeout, после join() нужно вызвать is_alive(), потому что сам join() не возвращает признак timeout.
Можно ли вызывать join() несколько раз
#
Да.
thread.join()
thread.join()
thread.join()
Это разрешено. Если поток уже завершился, последующие join() вернут управление почти сразу. В документации прямо указано, что поток можно join() много раз.
Когда join() вызовет ошибку
#
Нельзя вызвать join() у самого текущего потока:
import threading
threading.current_thread().join()
Это привело бы к deadlock:
поток ждёт сам себя
↓
он никогда не завершится
↓
программа зависает
Поэтому Python выбрасывает RuntimeError. Также RuntimeError будет, если вызвать join() у потока до его запуска через start().
Что с daemon-потоками #
Обычные потоки по умолчанию daemon=False.
thread = threading.Thread(target=worker)
Такие потоки не дадут программе завершиться, пока сами не закончат работу.
А daemon-потоки:
thread = threading.Thread(target=worker, daemon=True)
могут быть резко остановлены при завершении программы, если остались только daemon-потоки. Python предупреждает, что ресурсы daemon-потоков — например файлы или транзакции БД — могут быть не освобождены корректно. Для нормального завершения лучше использовать обычные потоки и механизм сигнала, например threading.Event.
Правильная остановка долгих потоков #
Если поток работает в цикле, его не надо пытаться “убивать”. Нужно дать ему сигнал остановиться:
import threading
import time
stop_event = threading.Event()
def worker():
while not stop_event.is_set():
print("Работаю...")
time.sleep(1)
print("Поток корректно завершился")
thread = threading.Thread(target=worker)
thread.start()
time.sleep(3)
stop_event.set()
thread.join()
print("Главный поток завершился")
Схема:
main thread
↓
запускает worker
↓
через 3 секунды вызывает stop_event.set()
↓
worker сам выходит из цикла
↓
main thread ждёт его через join()
Главное:
join() ждёт завершения потока
join() не запускает поток
join() не убивает поток
join(timeout=...) не сообщает результат напрямую
после join(timeout=...) проверяют is_alive()
нельзя делать join() текущего потока
для всех потоков сначала start(), потом отдельным циклом join()
43. Можно ли вручную переключить выполнение с одного потока на другой? #
Можно ли вручную переключить поток #
В обычном Python-коде — нет.
Нельзя сделать что-то вроде:
switch_to(thread_2)
и принудительно передать выполнение конкретному потоку.
Потоки планируются:
операционной системой
+
интерпретатором CPython через GIL
Модуль threading даёт API для создания и управления потоками, но не даёт ручного управления планировщиком потоков. Потоки выполняются конкурентно внутри одного процесса и делят память процесса.
Как переключение происходит на самом деле #
Упрощённо:
thread_1 выполняется
↓
ОС / CPython решают, что пора дать время другому потоку
↓
thread_2 получает выполнение
↓
потом снова может выполниться thread_1
В CPython есть GIL. Только поток, владеющий GIL, может работать с Python-объектами и выполнять Python C API. Интерпретатор регулярно пытается переключать потоки между инструкциями байткода; для блокирующего I/O GIL также освобождается.
Можно ли повлиять на переключение #
Можно только косвенно.
Например, можно задать интервал переключения:
import sys
sys.setswitchinterval(0.001)
Это не означает:
переключись прямо сейчас на thread_2
Это означает примерно:
попробуй чаще давать шанс другим Python-потокам
Но конкретный поток всё равно выбирает не твой код. Это остаётся задачей планировщика ОС и рантайма CPython. Документация CPython прямо связывает регулярные переключения между потоками с sys.setswitchinterval().
Как “уступить” выполнение другому потоку #
Иногда используют короткую паузу:
import time
time.sleep(0)
Идея такая:
текущий поток добровольно отдаёт шанс планировщику
Но это тоже не гарантирует, что следующим выполнится нужный тебе поток.
Пример:
import threading
import time
def worker(name):
for i in range(3):
print(name, i)
time.sleep(0)
t1 = threading.Thread(target=worker, args=("A",))
t2 = threading.Thread(target=worker, args=("B",))
t1.start()
t2.start()
t1.join()
t2.join()
Вывод может быть разным:
A 0
B 0
A 1
B 1
...
или другим. Порядок не гарантируется.
Как правильно управлять порядком потоков #
Если нужен не “ручной switch”, а контролируемый порядок работы, используют синхронизацию:
Lock
RLock
Event
Condition
Semaphore
Queue
Barrier
Пример через Event:
import threading
event = threading.Event()
def thread_1():
print("thread_1: подготовка")
event.set()
def thread_2():
event.wait()
print("thread_2: начал после thread_1")
t1 = threading.Thread(target=thread_1)
t2 = threading.Thread(target=thread_2)
t2.start()
t1.start()
t1.join()
t2.join()
Здесь мы не переключаем поток вручную. Мы задаём правило:
thread_2 не продолжает работу, пока thread_1 не вызовет event.set()
Главное различие #
Ручное переключение потока → нельзя
Косвенно повлиять на планирование → можно
Управлять порядком через синхронизацию → правильно
Кратко #
В Python нельзя вручную сказать интерпретатору:
останови этот поток и запусти вот тот
Можно только:
запустить поток через start()
ждать через join()
синхронизировать через Lock/Event/Queue
косвенно влиять через sleep() или sys.setswitchinterval()
Но конкретное переключение между потоками остаётся за ОС и CPython.
44. Когда операционная система переключает поток? | Какие ситуации приводят к переключению потоков? #
Операционная система переключает поток, когда текущий поток больше не должен или не может продолжать выполняться на CPU.
Это называется context switch — переключение контекста:
thread_1 выполняется на CPU
↓
ОС сохраняет его состояние
↓
ОС выбирает другой готовый поток
↓
thread_2 продолжает выполнение
Планировщик ОС решает, какой поток получит следующий квант процессорного времени. В Windows это прямо описывается как выбор конкурирующего потока для следующего processor time slice; в Linux обычные задачи планируются через CFS/EEVDF-подобную логику справедливого распределения CPU-времени.
Основные ситуации переключения потоков #
1. Истёк квант времени #
Поток работал слишком долго, и ОС дала шанс другому потоку:
thread_1 работает
↓
его квант времени закончился
↓
ОС переключает CPU на thread_2
Это типичная ситуация для вытесняющей многозадачности:
поток сам не обязан уступать CPU
ОС может прервать его принудительно
В Windows планировщик использует приоритеты и процессорные временные интервалы для выбора следующего потока.
2. Поток ушёл в ожидание I/O #
Например, поток ждёт:
чтение файла
ответ от сети
ответ от базы данных
ввод с клавиатуры
Пример:
data = socket.recv(4096)
Если данных ещё нет, поток блокируется:
thread_1 ждёт сеть
↓
CPU не должен простаивать
↓
ОС запускает другой готовый поток
Для Python это особенно важно: при блокирующем I/O CPython может освобождать GIL, чтобы другие Python-потоки могли выполняться. Официальная документация Python указывает, что GIL освобождается вокруг потенциально блокирующих I/O-операций, например чтения или записи файла.
3. Поток ждёт lock / mutex / semaphore #
Если поток пытается захватить занятый lock:
lock.acquire()
и lock уже удерживается другим потоком, текущий поток может перейти в состояние ожидания:
thread_1 хочет lock
↓
lock занят thread_2
↓
thread_1 блокируется
↓
ОС запускает другой поток
Это происходит не потому, что поток “сам переключился”, а потому что он больше не готов к выполнению.
4. Появился более приоритетный поток #
Если стал готов поток с более высоким приоритетом, ОС может вытеснить текущий поток:
thread_1 работает, приоритет обычный
↓
thread_2 стал готов, приоритет выше
↓
ОС переключается на thread_2
В Windows потоки планируются на основе scheduling priority; уровни приоритета влияют на то, какой поток будет выбран для выполнения.
5. Поток вызвал sleep()
#
Когда поток вызывает:
time.sleep(1)
он добровольно говорит:
мне не нужен CPU примерно 1 секунду
Схема:
thread_1 → sleep()
↓
thread_1 временно не runnable
↓
ОС выполняет другой поток
time.sleep(0) иногда используют как способ “уступить” планировщику шанс переключиться, но это не гарантирует, что следующим выполнится конкретный поток.
6. Поток завершился #
Когда поток закончил функцию:
def worker():
return
его больше не надо планировать:
thread_1 завершился
↓
ОС выбирает другой runnable-поток
7. Поток ждёт событие / условие / очередь #
Например:
event.wait()
queue.get()
condition.wait()
Если данных или события ещё нет, поток блокируется:
thread_1 ждёт event
↓
event ещё не установлен
↓
thread_1 спит
↓
CPU получает другой поток
Это нормальный способ координировать потоки без постоянной проверки в цикле.
Состояния потока упрощённо #
RUNNING
поток сейчас выполняется на CPU
READY / RUNNABLE
поток готов выполняться, но ждёт CPU
BLOCKED / WAITING
поток ждёт I/O, lock, event, sleep и т.д.
TERMINATED
поток завершился
Переключение обычно происходит так:
RUNNING → READY
истёк квант или вытеснен более приоритетным потоком
RUNNING → BLOCKED
поток ждёт I/O, lock, sleep, event
RUNNING → TERMINATED
поток завершил работу
В Python есть дополнительный уровень: GIL #
В CPython поток ОС может получить CPU, но для выполнения Python-байткода ему ещё нужен GIL.
Упрощённо:
ОС дала CPU thread_2
↓
thread_2 хочет выполнять Python-код
↓
но GIL у thread_1
↓
thread_2 ждёт GIL
Официальная документация Python указывает, что из-за GIL только один поток может выполнять Python bytecode одновременно, поэтому threading ограничен для CPU-bound задач.
CPython регулярно пытается переключать Python-потоки; это связано с sys.setswitchinterval():
import sys
print(sys.getswitchinterval())
sys.getswitchinterval() возвращает интервал переключения потоков интерпретатора в секундах.
Важно:
переключение OS-thread
≠
передача GIL другому Python-thread
ОС управляет потоками на уровне ядра, а GIL управляет доступом Python-потоков к выполнению Python-байткода.
Пример на Python #
import threading
import time
def worker(name):
for i in range(3):
print(name, i)
time.sleep(1)
t1 = threading.Thread(target=worker, args=("A",))
t2 = threading.Thread(target=worker, args=("B",))
t1.start()
t2.start()
t1.join()
t2.join()
Здесь time.sleep(1) приводит к тому, что поток временно не использует CPU:
A работает
↓
A вызывает sleep()
↓
ОС может выполнить B
↓
B вызывает sleep()
↓
ОС может вернуться к A
Но порядок вывода не гарантирован:
A 0
B 0
A 1
B 1
или:
B 0
A 0
B 1
A 1
Главное #
ОС переключает поток обычно в таких случаях:
истёк квант времени
поток ушёл в I/O-ожидание
поток ждёт lock / event / condition / queue
поток вызвал sleep()
появился поток с более высоким приоритетом
поток завершился
произошло вытеснение планировщиком
Для Python нужно помнить отдельно:
ОС переключает системные потоки
CPython дополнительно ограничивает выполнение Python-кода через GIL
То есть поток может существовать и даже быть запланирован ОС, но выполнять Python-байткод он сможет только тогда, когда получит GIL.
45. Как выполнить много корутин (asyncio.gather)? #
asyncio.gather() нужен, чтобы запустить несколько awaitable-объектов конкурентно и дождаться их всех.
Базовый шаблон:
import asyncio
async def worker(number: int):
await asyncio.sleep(1)
return f"done {number}"
async def main():
results = await asyncio.gather(
worker(1),
worker(2),
worker(3),
)
print(results)
asyncio.run(main())
Вывод:
['done 1', 'done 2', 'done 3']
asyncio.gather(*aws) запускает awaitable-объекты конкурентно. Если среди них coroutine, она автоматически оборачивается в Task. Результаты возвращаются списком в том же порядке, в котором awaitable были переданы в gather().
Почему это конкурентно #
Без gather():
result_1 = await worker(1)
result_2 = await worker(2)
result_3 = await worker(3)
Схема:
worker(1) → ждём завершения
worker(2) → ждём завершения
worker(3) → ждём завершения
С gather():
results = await asyncio.gather(
worker(1),
worker(2),
worker(3),
)
Схема:
worker(1) ┐
worker(2) ├─ выполняются конкурентно
worker(3) ┘
↓
gather ждёт завершения всех
Важно: это не значит, что они выполняются параллельно на разных ядрах. В asyncio конкурентность достигается за счёт переключения между корутинами в моменты await.
Много корутин из списка #
Часто корутины лежат в списке. Тогда нужен оператор распаковки *:
import asyncio
async def fetch(user_id: int):
await asyncio.sleep(1)
return {"id": user_id}
async def main():
user_ids = [1, 2, 3, 4, 5]
coroutines = [fetch(user_id) for user_id in user_ids]
results = await asyncio.gather(*coroutines)
print(results)
asyncio.run(main())
Без * будет ошибка, потому что gather() ждёт отдельные awaitable-объекты:
await asyncio.gather(coro1, coro2, coro3)
а не один список:
await asyncio.gather([coro1, coro2, coro3]) # неправильно
Правильно:
await asyncio.gather(*[coro1, coro2, coro3])
Можно передавать уже созданные Task #
Так тоже можно:
import asyncio
async def worker(number: int):
await asyncio.sleep(1)
return number * 10
async def main():
tasks = [
asyncio.create_task(worker(1)),
asyncio.create_task(worker(2)),
asyncio.create_task(worker(3)),
]
results = await asyncio.gather(*tasks)
print(results)
asyncio.run(main())
Но для простого случая create_task() не обязателен:
await asyncio.gather(worker(1), worker(2), worker(3))
Потому что gather() сам запланирует coroutine как задачи.
Что будет при ошибке #
По умолчанию, если одна корутина выбросит исключение, это исключение будет проброшено наружу из gather():
import asyncio
async def good():
await asyncio.sleep(1)
return "ok"
async def bad():
await asyncio.sleep(0.5)
raise ValueError("error")
async def main():
try:
result = await asyncio.gather(good(), bad())
print(result)
except ValueError as error:
print(f"Поймали ошибку: {error}")
asyncio.run(main())
При return_exceptions=False, это значение по умолчанию, первое исключение сразу передаётся ожидающему коду. Остальные awaitable не отменяются автоматически и продолжают выполнение.
Собрать ошибки как результаты #
Можно сделать так, чтобы ошибки попали в список результатов:
import asyncio
async def good():
await asyncio.sleep(1)
return "ok"
async def bad():
await asyncio.sleep(0.5)
raise ValueError("error")
async def main():
results = await asyncio.gather(
good(),
bad(),
return_exceptions=True,
)
print(results)
asyncio.run(main())
Пример результата:
['ok', ValueError('error')]
При return_exceptions=True исключения обрабатываются как обычные результаты и собираются в итоговый список.
Ограничение количества одновременных задач #
Если задач много, например 10 000 HTTP-запросов, не стоит запускать все сразу. Обычно используют asyncio.Semaphore:
import asyncio
sem = asyncio.Semaphore(3)
async def fetch(number: int):
async with sem:
print(f"start {number}")
await asyncio.sleep(1)
print(f"end {number}")
return number
async def main():
tasks = [fetch(i) for i in range(10)]
results = await asyncio.gather(*tasks)
print(results)
asyncio.run(main())
Здесь одновременно будут выполняться максимум 3 корутины внутри блока:
async with sem:
...
Timeout на gather #
Можно ограничить общее время выполнения:
import asyncio
async def worker(number: int):
await asyncio.sleep(number)
return number
async def main():
try:
results = await asyncio.wait_for(
asyncio.gather(
worker(1),
worker(2),
worker(5),
),
timeout=3,
)
print(results)
except TimeoutError:
print("Не все задачи успели завершиться")
asyncio.run(main())
Схема:
gather(...) запускает несколько задач
wait_for(...) ограничивает общее ожидание
Кратко #
results = await asyncio.gather(coro1(), coro2(), coro3())
Для списка:
results = await asyncio.gather(*coroutines)
Поведение:
gather запускает awaitable конкурентно
ждёт завершения всех
возвращает список результатов
порядок результатов соответствует порядку аргументов
при ошибке по умолчанию пробрасывает исключение
return_exceptions=True собирает ошибки в список
Главная ошибка новичков:
await asyncio.gather(coroutines) # неправильно
await asyncio.gather(*coroutines) # правильно
46. Как отслеживать статус асинхронной задачи, если она может зависнуть? #
Основные способы #
В asyncio статус задачи обычно отслеживают через объект Task:
task = asyncio.create_task(some_coro())
У Task есть методы:
task.done() # задача завершилась: успешно, с ошибкой или отменой
task.cancelled() # задача была отменена
task.result() # получить результат, если задача завершилась успешно
task.exception() # получить исключение, если задача упала
Task является разновидностью Future, поэтому у него есть методы проверки завершения, результата и исключения. Официальная документация также указывает, что coroutine, переданная в create_task(), планируется как задача в event loop.
Вариант 1: ждать задачу с timeout #
Самый простой способ защититься от зависания:
import asyncio
async def long_operation():
await asyncio.sleep(10)
return "done"
async def main():
task = asyncio.create_task(long_operation())
try:
result = await asyncio.wait_for(task, timeout=3)
print("Результат:", result)
except TimeoutError:
print("Задача выполнялась слишком долго")
print("Отменена:", task.cancelled())
asyncio.run(main())
asyncio.wait_for() ждёт awaitable ограниченное время. Если timeout истёк, задача отменяется, а наружу выбрасывается TimeoutError. Отмена в asyncio доставляется в coroutine через CancelledError при ближайшей возможности выполнения.
Вариант 2: проверять статус вручную #
Можно периодически смотреть, завершилась ли задача:
import asyncio
async def long_operation():
await asyncio.sleep(5)
return "ok"
async def main():
task = asyncio.create_task(long_operation())
while not task.done():
print("Задача ещё работает...")
await asyncio.sleep(1)
if task.cancelled():
print("Задача была отменена")
elif task.exception() is not None:
print("Задача упала с ошибкой:", task.exception())
else:
print("Результат:", task.result())
asyncio.run(main())
Схема:
create_task()
↓
периодически проверяем task.done()
↓
если done=True:
cancelled() → была отменена
exception() → завершилась с ошибкой
result() → завершилась успешно
Вариант 3: timeout без автоматической отмены задачи #
Иногда нужно ограничить ожидание, но не отменять саму задачу. Тогда используют asyncio.shield():
import asyncio
async def long_operation():
await asyncio.sleep(5)
return "done"
async def main():
task = asyncio.create_task(long_operation())
try:
await asyncio.wait_for(asyncio.shield(task), timeout=2)
except TimeoutError:
print("Ожидание истекло, но задача продолжает работать")
print("Статус после timeout:", task.done())
result = await task
print("Итоговый результат:", result)
asyncio.run(main())
Без shield() timeout через wait_for() отменяет задачу. С shield() отменяется ожидание, но исходная задача продолжает выполняться. Документация описывает shield() как способ защитить awaitable от внешней отмены.
Вариант 4: современный timeout-блок Python 3.11+ #
В Python 3.11+ можно использовать asyncio.timeout():
import asyncio
async def main():
try:
async with asyncio.timeout(3):
await asyncio.sleep(10)
print("completed")
except TimeoutError:
print("Блок выполнялся слишком долго")
asyncio.run(main())
Это удобно, когда нужно ограничить не одну задачу, а целый блок async-кода. asyncio.timeout() — асинхронный context manager для ограничения времени выполнения блока.
Важный нюанс: “зависла” или заблокировала event loop #
Timeout сработает нормально, если coroutine зависла на await:
await asyncio.sleep(999)
await session.get(...)
await queue.get()
Но если внутри coroutine запущен долгий синхронный код без await, event loop будет заблокирован:
async def bad_task():
while True:
pass
В таком случае event loop не сможет нормально переключаться между задачами, и timeout тоже может не сработать вовремя.
Плохо:
async def bad_task():
# CPU-bound цикл блокирует event loop
while True:
pass
Лучше:
async def better_task():
while True:
await asyncio.sleep(0)
Но для настоящей CPU-bound работы лучше использовать:
asyncio.to_thread()
ProcessPoolExecutor
ThreadPoolExecutor
отдельный процесс
Практический шаблон watchdog #
Если нужно наблюдать за задачей и отменять её при превышении лимита:
import asyncio
async def worker():
try:
await asyncio.sleep(10)
return "done"
except asyncio.CancelledError:
print("worker: получил отмену")
raise
async def watch_task(task: asyncio.Task, timeout: float):
try:
return await asyncio.wait_for(task, timeout=timeout)
except TimeoutError:
print("watchdog: задача зависла или выполняется слишком долго")
if not task.done():
task.cancel()
try:
await task
except asyncio.CancelledError:
print("watchdog: задача отменена")
return None
async def main():
task = asyncio.create_task(worker())
result = await watch_task(task, timeout=3)
print("result:", result)
print("done:", task.done())
print("cancelled:", task.cancelled())
asyncio.run(main())
Основные проверки:
task.done() # завершилась ли задача
task.cancelled() # была ли отменена
task.exception() # какое исключение было
task.result() # результат успешной задачи
Главное правило:
Если задача может зависнуть на await → используй wait_for() / asyncio.timeout().
Если задача блокирует event loop синхронным кодом → timeout может не помочь вовремя.
47. Как управлять асинхронными задачами внутри HTTP request? #
Внутри HTTP request асинхронными задачами нужно управлять по-разному в зависимости от того, нужен ли результат задачи прямо для ответа.
Нужен результат в response
↓
await / gather / TaskGroup + timeout
Не нужен результат прямо сейчас
↓
BackgroundTasks / очередь задач / отдельный worker
Задача долгая или важная
↓
не держать HTTP request, а отдавать 202 + job_id
1. Когда результат нужен для ответа #
Например, endpoint должен сходить в несколько сервисов/БД и вернуть объединённый результат.
Для Python 3.11+ лучше использовать asyncio.TaskGroup():
import asyncio
from fastapi import FastAPI
app = FastAPI()
async def get_user(user_id: int):
await asyncio.sleep(1)
return {"id": user_id, "name": "Alex"}
async def get_notifications(user_id: int):
await asyncio.sleep(1)
return [{"id": 1, "text": "Hello"}]
@app.get("/users/{user_id}/dashboard")
async def dashboard(user_id: int):
async with asyncio.timeout(3):
async with asyncio.TaskGroup() as tg:
user_task = tg.create_task(get_user(user_id))
notifications_task = tg.create_task(get_notifications(user_id))
return {
"user": user_task.result(),
"notifications": notifications_task.result(),
}
Что здесь происходит:
HTTP request пришёл
↓
создали несколько async-задач
↓
ждём их завершения внутри request
↓
собираем результат
↓
возвращаем response
TaskGroup даёт структурированную конкурентность: задачи создаются внутри блока, а выход из блока означает, что связанные задачи завершены или корректно обработаны. asyncio.timeout() ограничивает время выполнения блока. Оба механизма описаны в официальной документации asyncio.
2. Вариант через asyncio.gather()
#
Для простых случаев можно использовать gather():
import asyncio
from fastapi import FastAPI
app = FastAPI()
async def fetch_profile(user_id: int):
await asyncio.sleep(1)
return {"user_id": user_id}
async def fetch_stats(user_id: int):
await asyncio.sleep(1)
return {"posts": 10}
@app.get("/users/{user_id}")
async def get_user_data(user_id: int):
try:
profile, stats = await asyncio.wait_for(
asyncio.gather(
fetch_profile(user_id),
fetch_stats(user_id),
),
timeout=3,
)
except TimeoutError:
return {"error": "request timeout"}
return {
"profile": profile,
"stats": stats,
}
asyncio.gather() конкурентно запускает awaitable-объекты и возвращает результаты в порядке переданных аргументов. asyncio.wait_for() ограничивает ожидание и при timeout отменяет awaitable.
3. Не делать create_task() без контроля
#
Плохой вариант:
@app.post("/bad")
async def bad_endpoint():
asyncio.create_task(do_something())
return {"ok": True}
Проблемы:
задача может упасть, а ошибку никто не обработает
задача может жить дольше request
при перезапуске процесса задача потеряется
сложно контролировать timeout и отмену
сложно корректно закрывать ресурсы
create_task() внутри request допустим только когда ты сохраняешь ссылку на задачу, контролируешь её жизненный цикл, обрабатываешь ошибки и понимаешь, что задача живёт в памяти текущего процесса. Документация asyncio описывает create_task() как способ запланировать coroutine как Task, но сама задача не становится надёжной фоновой job-очередью.
4. Когда задача должна выполниться после ответа #
В FastAPI для коротких фоновых действий есть BackgroundTasks:
from fastapi import BackgroundTasks, FastAPI
app = FastAPI()
def write_log(user_id: int):
with open("log.txt", "a", encoding="utf-8") as file:
file.write(f"user_id={user_id}\n")
@app.post("/users/{user_id}/action")
async def user_action(user_id: int, background_tasks: BackgroundTasks):
background_tasks.add_task(write_log, user_id)
return {"status": "accepted"}
Схема:
request пришёл
↓
endpoint добавил background task
↓
response отправлен клиенту
↓
background task выполняется после response
FastAPI использует механизм фоновых задач Starlette. Starlette указывает, что background task прикрепляется к response и запускается только после отправки response.
5. Для тяжёлых задач нужен worker, а не HTTP request #
Для долгих или важных задач лучше не держать HTTP request открытым.
Правильная схема:
POST /reports
↓
создать job_id
↓
положить задачу в Redis/RabbitMQ/Celery/RQ/arq
↓
вернуть 202 Accepted + job_id
GET /reports/{job_id}
↓
проверить статус задачи
↓
вернуть pending / success / failed
Пример ответа:
{
"job_id": "abc123",
"status": "queued"
}
FastAPI прямо предупреждает: для тяжёлых фоновых вычислений и задач, которые должны выполняться в нескольких процессах или серверах, лучше использовать более серьёзные инструменты вроде Celery с брокером сообщений, например Redis или RabbitMQ.
6. Как контролировать зависание внутри request #
Используй timeout:
import asyncio
from fastapi import HTTPException
@app.get("/external")
async def external_data():
try:
async with asyncio.timeout(5):
result = await call_external_service()
return result
except TimeoutError:
raise HTTPException(status_code=504, detail="External service timeout")
Для HTTP API обычно логичнее возвращать:
504 Gateway Timeout → внешний сервис не ответил вовремя
408 Request Timeout → клиент слишком долго отправлял request
500 Internal Server Error → внутренняя ошибка сервера
7. Как ограничить количество одновременных задач #
Например, нельзя запускать 1000 запросов к внешнему сервису одновременно. Используй asyncio.Semaphore:
import asyncio
sem = asyncio.Semaphore(10)
async def limited_call(item_id: int):
async with sem:
return await call_external_service(item_id)
@app.post("/batch")
async def batch_endpoint(ids: list[int]):
async with asyncio.timeout(10):
results = await asyncio.gather(
*(limited_call(item_id) for item_id in ids)
)
return {"results": results}
Схема:
100 задач создано
↓
но одновременно внутрь async with sem проходят только 10
↓
остальные ждут
Это защищает приложение и внешние сервисы от перегрузки.
8. Что делать при отключении клиента #
Для long-polling или streaming иногда нужно проверять, не отключился ли клиент:
from fastapi import Request
@app.get("/long")
async def long_request(request: Request):
for i in range(100):
if await request.is_disconnected():
return {"status": "client disconnected"}
await asyncio.sleep(1)
return {"status": "done"}
Starlette/FastAPI позволяют проверить состояние соединения через await request.is_disconnected(), что полезно для long-polling и streaming-сценариев.
9. Не использовать ресурсы request в долгой фоновой задаче #
Плохая идея:
@app.post("/bad")
async def bad(background_tasks: BackgroundTasks, db: AsyncSession = Depends(get_db)):
background_tasks.add_task(do_work, db)
return {"ok": True}
Проблема: db, request context, transaction, dependency scope могут быть закрыты после завершения request.
Лучше передавать простые данные:
@app.post("/good")
async def good(background_tasks: BackgroundTasks, user_id: int):
background_tasks.add_task(do_work, user_id)
return {"ok": True}
А внутри фоновой задачи открыть новое соединение/сессию:
async def do_work(user_id: int):
async with async_session_maker() as session:
...
Практическое правило #
await внутри request
↓
когда результат нужен для response
TaskGroup / gather
↓
когда нужно выполнить несколько async-операций конкурентно
asyncio.timeout / wait_for
↓
когда задача может зависнуть
BackgroundTasks
↓
короткая best-effort задача после response
Celery / RQ / arq / отдельный worker
↓
долгая, тяжёлая или важная задача
Semaphore
↓
ограничение параллельности
job_id + status endpoint
↓
правильная модель для долгих операций
Главное: HTTP request не должен превращаться в бесконтрольный контейнер для долгоживущих задач. Для задач, от которых зависит ответ, жди их явно. Для задач после ответа используй background-механизм. Для надёжных долгих задач — очередь и отдельный worker.