Многопоточность и GIL
Следующим шагом после профилирования, векторизации и JIT-компиляции является использование всех ядер процессора, а не одного. В Python это непросто, и рассмотрение необходимо начать с основ: что представляет собой процесс, что представляет собой поток и почему в Python второе не даёт ожидаемого результата.
Необходимый минимум: процессы и потоки
Процесс
Процесс — запущенная программа, и операционная система выделяет каждому процессу собственное, изолированное от остальных состояние:
- виртуальное адресное пространство, то есть собственную память, скрытую от остальных процессов;
- указатель на исполняемую инструкцию;
- стек вызовов;
- системные ресурсы, например открытые файловые дескрипторы.
Несколько задач, которые необходимо выполнять одновременно, можно распределить по разным процессам. Процессы, разведённые по своим адресным пространствам, не мешают друг другу, однако и обмен данными между ними затруднён: общей памяти у них нет.
Поток
Поток исполняется независимо от других потоков, как и процесс, но существует внутри процесса и разделяет с ним адресное пространство и системные ресурсы. Два потока одного процесса могут свободно работать с общими данными, поскольку они видят одни и те же объекты, размещённые в общей памяти.
Это удобно до тех пор, пока два потока не обратятся к одним данным одновременно: в этом случае возникают гонки, о которых сказано ниже.
Порядком исполнения как процессов, так и потоков распоряжается операционная система, поочерёдно выделяющая каждому по нескольку тактов процессора.
Модуль threading
Поток в Python является обычным системным потоком, и его исполнением управляет операционная система, а не интерпретатор. Создать поток можно классом Thread, объявленным в модуле стандартной библиотеки threading.
import time
from threading import Thread
def countdown(n):
for i in range(n):
print(n - i - 1, "left")
time.sleep(1)
t = Thread(target=countdown, args=(3,))
t.start()
2 left
1 left
0 left
Другой способ создать поток — наследование.
class CountdownThread(Thread):
def __init__(self, n):
super().__init__()
self.n = n
def run(self): # вызывается методом start
for i in range(self.n):
print(self.n - i - 1, "left")
time.sleep(1)
t = CountdownThread(3)
t.start()
2 left
1 left
0 left
Этот подход ограничивает повторное использование кода: функциональность, заключённая в класс CountdownThread, доступна только в отдельном потоке.
У потока есть имя, по умолчанию Thread-N; имя, заданное разработчиком, лучше читается в журналах.
Thread().name
'Thread-6'
Thread(name="NumberCruncher").name
'NumberCruncher'
У каждого активного потока есть идентификатор, неотрицательное число, выданное операционной системой и уникальное среди всех активных потоков.
t = Thread()
t.start()
t.ident
123145519448064
Дождаться завершения потока позволяет метод join: вызывающий поток блокируется, пока t не завершит работу; повторный join на уже завершившемся потоке возвращается мгновенно.
t = Thread(target=time.sleep, args=(5, ))
t.start()
t.join() # блокирует на 5 секунд
t.join() # выполняется мгновенно
Проверить, жив ли поток, можно методом is_alive.
t = Thread(target=time.sleep, args=(5, ))
t.start()
t.is_alive() # False через 5 секунд
True
Поток, созданный с аргументом daemon=True, называется демоном: при выходе из интерпретатора он уничтожается автоматически, а не удерживает программу от завершения до окончания своей работы.
t = Thread(target=time.sleep, args=(5,), daemon=True)
t.start()
t.daemon
True
Встроенного способа принудительно завершить поток в Python нет, и это осознанное решение. Корректное завершение почти всегда означает освобождение ресурсов: поток мог открыть файл, дескриптор которого необходимо закрыть, или захватить примитив синхронизации, который необходимо освободить. Принудительное завершение оставило бы всё это в неопределённом состоянии, поэтому запрос на завершение передаётся потоку флагом, выставленным снаружи.
class Task:
def __init__(self):
self._running = True
def terminate(self):
self._running = False
def run(self, n):
while self._running:
...
Набор примитивов синхронизации в модуле threading стандартный:
Lock, обычный мьютекс, обеспечивающий эксклюзивный доступ к разделяемому состоянию.RLock, рекурсивный мьютекс, разрешающий потоку, владеющему мьютексом, захватывать его более одного раза.Semaphore, вариация мьютекса, захватываемая не более фиксированного числа раз.BoundedSemaphore, семафор, следящий за тем, чтобы его захватывали и освобождали одинаковое число раз.
Все примитивы синхронизации реализуют единый интерфейс:
- метод
acquireзахватывает примитив синхронизации, - а метод
releaseосвобождает его.
from threading import Lock
class SharedCounter:
def __init__(self, value):
self._value = value
self._lock = Lock()
def increment(self, delta=1):
self._lock.acquire()
self._value += delta
self._lock.release()
def get(self):
return self._value
Модуль queue
Обмен данными между потоками с ручной расстановкой мьютексов неудобен. Обычно применяется готовая очередь; модуль queue реализует несколько потокобезопасных вариантов.
Queue, очередь FIFO,LifoQueue, очередь LIFO, то есть стек,PriorityQueue, очередь с приоритетом: элементы, обычно пары (priority, item), выдаются не в порядке поступления, а по возрастанию приоритета.
Все методы, изменяющие состояние, работают под мьютексом. Queue хранит элементы в deque, а LifoQueue и PriorityQueue — в обычном списке, защищённом тем же мьютексом.
Классическая схема «производитель-потребитель» на очереди приведена ниже.
def worker(q):
while True:
item = q.get() # блокирующе ждёт следующий
do_something(item) # элемент
q.task_done() # уведомляет очередь о выполнении
def master(q):
for item in source():
q.put(item)
# блокирующе ждёт, пока все элементы
# очереди не будут обработаны
q.join()
Модуль futures
Модуль concurrent.futures предоставляет абстрактный класс Executor и его реализацию над пулом потоков, ThreadPoolExecutor. Интерфейс исполнителя состоит из трёх методов.
from concurrent.futures import *
executor = ThreadPoolExecutor(max_workers=4)
executor.submit(print, "Hello, world!")
Hello, world!
<Future at 0x1043b1d90 state=running>
list(executor.map(print, ["Knock?", "Knock!"]))
Knock?
Knock!
[None, None]
executor.shutdown()
Исполнители поддерживают протокол менеджеров контекста: пул, открытый в with, закрывать вручную не требуется.
with ThreadPoolExecutor(max_workers=4) as executor:
...
Метод Executor.submit возвращает объект Future, обёртку над ещё не завершённым вычислением.
С Future можно выполнить следующие действия.
with ThreadPoolExecutor(max_workers=4) as executor:
f = executor.submit(sorted, [4, 3, 1, 2])
- Запросить статус вычисления:
f.running(), f.done(), f.cancelled()
(False, True, False)
- Блокирующе дождаться результата вычисления:
print(f.result())
[1, 2, 3, 4]
print(f.exception())
None
- Добавить функцию, вызываемую после завершения вычисления:
f.add_done_callback(print)
<Future at 0x1043b7f10 state=finished returned list>
Пример с модулем futures: integrate
Проверим это на вычислительной задаче: численное интегрирование методом прямоугольников, циклом по заданному разбиению отрезка.
import math
def integrate(f, a, b, *, n_iter=1000):
acc = 0
step = (b - a) / n_iter
for i in range(n_iter):
acc += f(a + i * step) * step
return acc
integrate(math.cos, 0, math.pi / 2)
1.0007851925466296
from functools import partial
def integrate_async(f, a, b, *, n_jobs, n_iter=1000):
executor = ThreadPoolExecutor(max_workers=n_jobs)
spawn = partial(executor.submit, integrate, f,
n_iter=n_iter // n_jobs)
step = (b-a)/n_jobs
fs=[spawn(a+i*step,a+(i+1)*step)
for i in range(n_jobs)]
return sum(f.result() for f in as_completed(fs))
integrate_async(math.cos, 0, math.pi / 2, n_jobs=2)
1.0007851925466305
Параллелизм и конкурентность
Сравним производительность последовательной и параллельной версий integrate магической командой timeit.
%%timeit -n100
integrate(math.cos, 0, math.pi / 2, n_iter=10**6)
154 ms ± 2.28 ms per loop (mean ± std. dev. of 7 runs, 100 loops each)
%%timeit -n100
integrate_async(math.cos, 0, math.pi / 2, n_iter=10**6, n_jobs=2)
142 ms ± 2.11 ms per loop (mean ± std. dev. of 7 runs, 100 loops each)
Глобальная блокировка интерпретатора
Два потока вместо одного дали выигрыш в проценты вместо ожидаемого двукратного ускорения. Причиной является GIL (global interpreter lock, глобальная блокировка интерпретатора).
GIL представляет собой мьютекс, гарантирующий, что в каждый момент времени байт-код исполняет только один поток. Внутреннее состояние интерпретатора (в первую очередь счётчики ссылок у каждого объекта) не защищено от одновременного доступа, и глобальная блокировка является самым простым способом не позволить двум потокам одновременно нарушить его целостность. Платой за простоту является результат приведённого выше замера.
Следствие: в обычной сборке CPython потоки не ускоряют вычисления. Сколько бы их ни было, байт-код исполняется по очереди, а к вычислениям добавляются накладные расходы на переключение, из-за которых многопоточная версия иногда оказывается даже медленнее однопоточной.
Влияние GIL на разные задачи
Мешает ли GIL, зависит от того, чем занята программа.
Если программа вычисляет, потоки бесполезны. Обойти GIL можно двумя путями: перейти к отдельным процессам, у каждого из которых свой интерпретатор со своим GIL (об этом сказано ниже), или перейти к скомпилированному коду, способному освобождать GIL.
Если программа ожидает, то есть читает файл, получает данные по сети, опрашивает прибор, на время ожидания поток освобождает GIL, и остальные потоки работают. Для такой нагрузки потоки подходят, и проигрыша от блокировки нет.
NumPy на длинных операциях также освобождает GIL: пока перемножаются большие матрицы, интерпретатор свободен для остальных потоков. Поэтому многопоточность в сочетании с NumPy работает лучше, чем можно было бы ожидать.
Оговорка «в обычной сборке CPython» существенна. Начиная с 3.13 интерпретатор собирается и вовсе без глобальной блокировки (PEP 703): такая сборка устанавливается рядом с обычной и называется python3.13t, а с 3.14 режим свободных потоков официально поддержан (PEP 779), и потоки в нём вычисляют параллельно, без процессов и без Cython. Платой являются 5–10 %, потерянные на однопоточном коде, и неготовность части C-расширений, написанных в расчёте на то, что блокировка защитит их сама. Определить, в какой сборке выполняется программа, позволяет sys._is_gil_enabled(), возвращающий False в свободнопоточной сборке. Пока она не стала используемой по умолчанию, всё изложенное в настоящей главе относится к обычной сборке.
C и Cython: освобождение GIL
Скомпилированный код может освободить GIL явно, на то время, пока он не обращается к объектам Python. В Cython для этого предусмотрен блок with nogil, внутри которого допустимы только типы C, зато остальные потоки работают параллельно.
%load_ext Cython
%%cython
from libc.math cimport cos
def integrate(f, double a, double b, long n_iter):
cdef double acc = 0
cdef double step=(b-a)/n_iter
cdef long i
with nogil:
for i in range(n_iter):
acc += cos(a + i * step) * step
return acc
Первый аргумент здесь оставлен только ради совместимости с integrate_async, а сама функция жёстко привязана к cos: вызвать функцию Python f внутри with nogil невозможно.
%%timeit -n100
integrate_async(math.cos, 0, math.pi / 2, n_iter=10**6, n_jobs=2)
5.88 ms ± 126 µs per loop (mean ± std. dev. of 7 runs, 100 loops each)
Те же два потока, тот же integrate_async, и 5.88 мс вместо 142. Этот выигрыш является составным. Потоков два, поэтому на снятие GIL приходится не более двукратного ускорения; остальное, примерно двенадцатикратное, обеспечила компиляция в машинный код. Таким образом, сначала применяется Cython, и только затем распараллеливается уже скомпилированный код.
О замерах ниже. Ячейка выше переопределила имя
integrate: теперь под ним находится скомпилированная Cython-версия. Замеры двух следующих разделов — про процессы и проjoblib— сняты не с неё, а с исходной, чистой Python-версии в свежем ядре; для их повторения необходимо перезапустить блокнот, пропустив ячейку с%%cython. В противном случае передать Cython-функцию в пул процессов и вовсе не удастся: она находится в модуле_cython_magic_<hash>из~/.ipython/cython, который отсутствует вsys.path.
Модуль multiprocessing
Процессы как способ обойти GIL
Поскольку глобальная блокировка одна на интерпретатор, самым прямым способом её обхода является запуск нескольких интерпретаторов. Именно это и обеспечивают процессы: у каждого свой GIL, друг другу они не мешают, и на восьми ядрах восемь процессов вычисляют одновременно.
Платой является изоляция. Процессы не разделяют память, поэтому всё, что необходимо передать между ними, приходится сериализовать и копировать. Для расчёта, в котором каждому процессу передаётся часть работы и забирается число, это несущественно; для задачи с большим общим состоянием накладные расходы на пересылку могут поглотить весь выигрыш.
За работу с процессами отвечает модуль multiprocessing, устроенный аналогично threading.
import multiprocessing as mp
p = mp.Process(target=countdown, args=(5, ))
p.start()
4 left
3 left
2 left
1 left
0 left
Примитивы синхронизации те же: мьютексы, семафоры, условные переменные. Для обмена данными между процессами предусмотрен Pipe, соединение двух процессов поверх сокетов.
def ponger(conn):
conn.send("pong")
parent_conn, child_conn = mp.Pipe()
p = mp.Process(target=ponger, args=(child_conn, ))
p.start()
parent_conn.recv()
'pong'
p.join()
Процессы и производительность
Реализация integrate_async на пуле потоков работала долго; рассмотрим пул процессов.
from concurrent.futures import ProcessPoolExecutor
def integrate_async(executor, f, a, b, *, n_jobs, n_iter=1000):
spawn = partial(executor.submit, integrate, f,
n_iter=n_iter // n_jobs)
step = (b - a) / n_jobs
fs=[spawn(a + i * step, a + (i + 1) * step)
for i in range(n_jobs)]
return sum(f.result() for f in as_completed(fs))
Пул принимается аргументом, а не создаётся внутри. Запуск
процесса стоит десятки миллисекунд: на macOS и Windows интерпретатор запускается
заново, целиком, а начиная с Python 3.14 и на Linux место fork занял
forkserver, порождающий процессы от отдельного чистого сервера. Наследования
всего состояния родителя больше нет нигде, поэтому и функция, передаваемая
в пул, должна находиться в импортируемом модуле, а не в ячейке блокнота. Пул,
созданный внутри измеряемой функции, превратит замер вычислений в замер запуска интерпретатора.
Кроме того, ProcessPoolExecutor, оставленный без закрытия, оставляет процессы работающими,
и за сотню итераций %%timeit их число на машине станет значительным.
with ProcessPoolExecutor(max_workers=2) as executor:
integrate_async(executor, math.cos, 0, math.pi / 2,
n_iter=10**6, n_jobs=2) # прогрев
%timeit -n5 integrate_async(executor, math.cos, 0, math.pi / 2, \
n_iter=10**6, n_jobs=2)
21 ms ± 0.2 ms per loop (mean ± std. dev. of 5 runs, 5 loops each)
В этом же прогоне последовательный расчёт занимает 44 мс, а два процесса дают 21 мс, то есть двукратное ускорение на паре ядер: у каждого процесса свой GIL. С числами из разделов выше эту пару сопоставлять нельзя: она снята отдельно.
Пакет joblib
Пакет joblib предоставляет параллельный аналог цикла for для независимых итераций, которые необходимо распределить по ядрам. Выбор между потоками и процессами задаётся аргументом backend конструктора.
Замеры ниже выполнены на исходной, чистой Python-версии integrate. Если в блокноте
уже выполнена ячейка с %%cython, имя integrate указывает на скомпилированную
функцию, и числа получатся на порядок меньше. Сравнивать их с результатами
до Cython нельзя; для повторения замеров отсюда необходимо определить integrate заново.
from joblib import Parallel, delayed
def integrate_async(f, a, b, *, n_jobs, n_iter=1000, backend=None):
step = (b - a) / n_jobs
with Parallel(n_jobs=n_jobs, backend=backend) as parallel:
fs = (delayed(integrate)(f, a + i * step,
a + (i + 1) * step,
n_iter=n_iter // n_jobs)
for i in range(n_jobs))
return sum(parallel(fs)) # внутри with: пул переиспользуется
%%timeit -n100
integrate_async(math.cos, 0, math.pi / 2, n_iter=10**6, n_jobs=2, backend="threading")
104 ms ± 280 µs per loop (mean ± std. dev. of 7 runs, 100 loops each)
%%timeit -n100
integrate_async(math.cos, 0, math.pi / 2, n_iter=10**6, n_jobs=2, backend="multiprocessing")
290 ms ± 1.13 ms per loop (mean ± std. dev. of 7 runs, 100 loops each)
Сам по себе joblib ничего не ускоряет: он распределяет работу, которая упирается в тот же GIL на потоках или в ту же сериализацию на процессах.
JIT-компилятор numba
Numba из главы «Скорость выполнения программ» обходит GIL одновременно с ускорением: декоратор @jit(nopython=True, parallel=True) компилирует функцию в машинный код, освобождающий блокировку, после чего цикл prange распределяется по ядрам.
import math
from numba import jit, prange
@jit(nopython=True, parallel=True, fastmath=True)
def integrate(a, b, n_iter=1000):
acc = 0
step = (b - a) / n_iter
for i in prange(n_iter):
acc += math.cos(a + i * step) * step
return acc
%%timeit -n100
integrate(0, math.pi / 2, n_iter=10**6)
5.11 ms ± 983 µs per loop (mean ± std. dev. of 7 runs, 100 loops each)
Резюме
Перед выбором инструмента необходимо ответить на один вопрос: программа вычисляет или ожидает?
Если программа ожидает сети, диска, прибора, ответа пользователя, используются потоки. GIL на время ожидания освобождается, потоки дёшевы, и десятки одновременно висящих ожиданий обходятся почти бесплатно. Если ожиданий предстоят тысячи, требуется перейти к асинхронности, рассматриваемой в следующей главе.
Если программа вычисляет, потоки в обычной сборке CPython не помогут, сколько бы их ни было. Существует три пути: вынести расчёт в отдельные процессы через multiprocessing, где у каждого свой интерпретатор и свой GIL; перейти к скомпилированному коду на Cython или C, способному освобождать GIL; или, что обычно проще всего, поручить расчёт NumPy, освобождающему блокировку самостоятельно.
Прежде чем распараллеливать, необходимо измерить. Часто узкое место находится не там, где предполагалось, и одна правка алгоритма даёт больше, чем любое количество потоков.