Параллелизм и многопоточность низкоуровнего кода asyncio в Python
Цикл событий выполняется в потоке (обычно в основном потоке) и выполняет все обратные вызовы и задачи в своем потоке. Пока задача Task выполняется в цикле событий, никакие другие задачи не могут выполняться в том же потоке. Когда задача Task выполняется с оператором await , то выполняющаяся задача приостанавливается, а цикл обработки событий выполняет следующую задачу.
Чтобы запланировать обратный вызов из другого потока ОС, следует использовать метод loop.call_soon_threadsafe() . Пример:
loop.call_soon_threadsafe(callback, *args)
Почти все объекты модуля asyncio не являются потокобезопасными, что обычно не является проблемой, если нет кода, который работает с ними извне. Если такой код необходим, то для вызова низкоуровневого API, следует использовать метод loop.call_soon_threadsafe() , например:
loop.call_soon_threadsafe(future.cancel)
Чтобы запланировать объект сопрограммы из другого потока ОС, следует использовать функцию asyncio.run_coroutine_threadsafe() . Она возвращает concurrent.futures.Future для доступа к результату:
async def coro_func(): return await asyncio.sleep(1, 42) # Позже в другом потоке ОС: future = asyncio.run_coroutine_threadsafe(coro_func(), loop) # Ждите результата: result = future.result()
Для обработки сигналов и выполнения подпроцессов, цикл событий должен выполняться в основном потоке.
Метод loop.run_in_executor() можно использовать с concurrent.futures.ThreadPoolExecutor для выполнения блокирующего кода в другом потоке ОС без блокировки основного потока ОС, в котором выполняется цикл событий.
В настоящее время нет возможности запланировать сопрограммы или обратные вызовы непосредственно из другого процесса (например, запущенного с многопроцессорной обработкой). Однако, API-интерфейс модуля asyncio для обслуживания subprocess позволяют запускать процесс и взаимодействовать с ним из цикла событий.
Наконец, вышеупомянутый метод loop.run_in_executor() также можно использовать с concurrent.futures.ProcessPoolExecutor для выполнения кода в другом процессе.
Пример запуска сопрограммы в другом потоке:
В примере функция worker() запускается явно в задаче в текущем цикле событий и отправляется при помощи функции asyncio.run_coroutine_threadsafe() в новый цикл событий, созданный в отдельном потоке.
import asyncio, threading async def worker(name, delay): th_name = threading.current_thread().name print(f'Start name>; ожидание delay>; поток: th_name>') res = await asyncio.sleep(delay, result=delay) print(f'Done name>; ожидание delay>') return name, res async def main(new_loop): results = [] # Функция `run_coroutine_threadsafe()` проталкивает # `worker()` в поток с циклом событий `new_loop` future1 = asyncio.run_coroutine_threadsafe(worker('Thread', 1), new_loop) results.append(future1) task = asyncio.create_task(worker('Task', 1.5)) future2 = await task results.append(future2) print('\nРезультаты:') for future in results: if type(future) == tuple: print(f'Задача future[0]>; результат future[1]>') else: res = future.result() print(f'Задача res[0]>; результат res[1]>') # останавливаем цикл событий `new_loop` new_loop.call_soon_threadsafe(new_loop.stop) if __name__ == '__main__': # получаем новый цикл событий new_loop = asyncio.new_event_loop() # создаем поток с запущенным новым циклом событий thread = threading.Thread(target=new_loop.run_forever) thread.start() asyncio.run(main(new_loop)) # Start Task; ожидание 1.5; поток: MainThread # Start Thread; ожидание 1; поток: Thread-1 # Done Thread; ожидание 1 # Done Task; ожидание 1.5 # Результаты: # Задача Thread; результат 1 # Задача Task; результат 1.5
Пример запуска нескольких команд в терминале из асинхронного кода.
Так как все функции для запуска подпроцесса модуля asyncio являются асинхронными, то легко выполнять и контролировать несколько подпроцессов, выполняемых параллельно.
import asyncio async def run(cmd): proc = await asyncio.create_subprocess_shell( cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) stdout, stderr = await proc.communicate() print(f'[cmd!r> exited with proc.returncode>]') if stdout: print(f'[stdout]\nstdout.decode()>') if stderr: print(f'[stderr]\nstderr.decode()>') async def main(): await asyncio.gather( run('ls /zzz'), run('sleep 1; echo "hello"')) asyncio.run(main()) # ['ls /zzz' завершилась с кодом 2] # [stderr] # ls: невозможно получить доступ к '/zzz': Нет такого файла или каталога # ['sleep 1; echo "Привет"' завершилась с кодом 0] # [stdout] # Привет
- КРАТКИЙ ОБЗОР МАТЕРИАЛА.
- Сопрограммы и механизмы их запуска модулем asyncio
- Что такое аwaitable объект модуля asyncio
- Функция run() модуля asyncio
- Менеджер контекста Runner() модуля asyncio
- Функция create_task() модуля asyncio
- Группы задач TaskGroup() модуля asyncio
- Класс Task() модуля asyncio
- Функция sleep() модуля asyncio
- Функция gather() модуля asyncio
- Функция shield() модуля asyncio
- Асинхронный менеджер timeout() модуля asyncio
- Асинхронный менеджер timeout_at() модуля asyncio
- Функция wait_for() модуля asyncio
- Функция as_completed() модуля asyncio
- Функция wait() модуля asyncio
- Функция to_thread() модуля asyncio
- Функция run_coroutine_threadsafe() модуля asyncio
- Функции current_task() и all_tasks() модуля asyncio
- Использование очереди asyncio.Queue
- Примитивы синхронизации задач в asyncio
- Запуск внешних программ из кода asyncio
- Работа с сетевыми соединениями модуля asyncio
- Объект Future модуля asyncio Python и связанные функции
- Создание и получение текущего цикла событий, модуль asyncio
- Создание, запуск и остановка цикла событий модуля asyncio
- Планирование обратных вызовов из цикла событий asyncio
- Создание Future и Task из цикла событий asyncio
- Немедленное выполнение задач модулем asyncio
- Создание пулов потоков и процессов из цикла событий asyncio
- Создание TCP, UDP и Unix соединений из цикла событий asyncio
- Создание сетевых серверов из цикла событий asyncio
- Создание субпроцесса из цикла событий asyncio
- Работа с сокетами напрямую из цикла событий asyncio
- Передача файлов из цикла событий asyncio
- Наблюдение за дескрипторами файлов из цикла событий asyncio
- DNS запросы из цикла событий asyncio
- Сигналы Unix в циклах событий asyncio
- Параллелизм и многопоточность в цикле событий asyncio
- Объекты Transport и Protocol в цикле событий asyncio
- Включение режима отладки в asyncio
- Обработка исключений в цикле событий модуля asyncio
- Исключения модуля asyncio Python
Как в одном потоке получить данные из других потоков?
Я запускаю 15 потоков. Из них 1 поток я отправляю пополнять количество задач функцией update_tasks() А остальные 14 должны решать эти задачи функцией do_work() Новые задачи поступают каждые 10 секунд. При этом, в свежеполученных задачах постоянно приходят те, что уже выполнены. Поэтому функция update_tasks() сравнивает свежеполученные задачи с теми, что уже сделаны, и теми, что находятся сейчас в очереди. А теперь сама проблема: do_work() решает задачу примерно за 1 минуту и 14 задач, которые сейчас в работе снова поступают в список новых задач. Потому что они еще не появились в списке решенных задач. Но и в q.queue их уже нет. Эта конструкция их пропускает:
asks = set(get_new_tasks()) - set(get_completed_tasks) - set(list(q.queue))
Подскажите, пожалуйста, как можно одним потоком получить задачи, которые в это время выполняют другие потоки, для того чтобы избежать повторного выполнения одних и тех же задач?
Многопоточность в Python
Иногда случается так, что необходимо выполнять какие-то процедуры параллельно друг-другу и независимо друг от друга – например, получать курсы с разных бирж. Или проверять состояние заказов на разных парах. Или грабить текст с чужих сайтов. Или заливать текст на чужие сайты. Или распознавать капчи.. Да мало ли чего!
В этой статье попрактикуемся в написании таких скриптов.
Без многопоточности
Давайте начнем с бесполезной задачи – будем получать книгу ордеров разных валютных пар на одной бирже. Для демонстрации возьмем Exmo. И для начала давайте без многопоточности – сделаем всё последовательно.
Нам понадобится функция, которая будет делать всю работу, и которой мы будем передавать нужные пары. Вот так будет выглядеть нулевая версия скрипта:
import time import requests def get_rates(pair): local_start_time = time.time() try: requests.get("https://api.exmo.com/v1/order_book/?pair= &limit=1000".format(pair=pair)) except Exception as e: print(e) print("Пара , время работы функции: ".format(pair=pair, t=time.time()-local_start_time)) get_rates('BTC_LTC')
Запустим его и узнаем, что получение одной пары занимает примерно половину секунды (на самом деле, когда как)

Давайте добавим несколько пар, будем прогонять их в цикле, и заодно замерим общее время работы скрипта:
import time import requests pairs = ['BTC_LTC', 'BTC_ETH', 'BTC_USD', 'BTC_EUR', 'BTC_PLN', 'BCH_BTC', 'EOS_BTC', 'EOS_USD', 'BCH_RUB', 'BCH_ETH'] def get_rates(pair): local_start_time = time.time() try: requests.get("https://api.exmo.com/v1/order_book/?pair= &limit=1000".format(pair=pair)) except Exception as e: print(e) print("Пара , время работы функции: ".format(pair=pair, t=time.time()-local_start_time)) global_start_time = time.time() for pair in pairs: get_rates(pair) print('Общее время работы '.format(s=time.time()-global_start_time))

На получение 10 пар ушло 3.5 секунды, данные первой пары соответственно устарели на три секунды.
Многопоточность
В питоне из коробки идет модуль threading. Именно он отвечает за (условно) параллельное исполнение кода. Почему условное – расскажу ниже.
Для того, что бы создать поток, нужно два действия
- Запланировать его
- Запустить
Подготовим скрипт и прогоним его:
import time import requests import threading pairs = ['BTC_LTC', 'BTC_ETH', 'BTC_USD', 'BTC_EUR', 'BTC_PLN', 'BCH_BTC', 'EOS_BTC', 'EOS_USD', 'BCH_RUB', 'BCH_ETH'] def get_rates(pair): local_start_time = time.time() try: requests.get("https://api.exmo.com/v1/order_book/?pair= &limit=1000".format(pair=pair)) except Exception as e: print(e) print("Пара , время работы функции: ".format(pair=pair, t=time.time()-local_start_time)) global_start_time = time.time() threads = [] for pair in pairs: # Подготавливаем потоки, складываем их в массив threads.append(threading.Thread(target=get_rates, args=(pair,))) # Запускаем каждый поток for thread in threads: thread.start() # Ждем завершения каждого потока for thread in threads: thread.join() print('Общее время работы '.format(s=time.time()-global_start_time))

А теперь рассмотрим детально и сделаем выводы.
Во-первых, можно увидеть, что пары обрабатывались не в том порядке, что мы вызывали.
Во-вторых, общее время работы уменьшилось (а у отдельно взятых пар увеличилось). Так же самое устаревшее время составило 2.47 секунды. В целом результат позитивный, а теперь заглянем под капот.
Когда вы передаете что-либо модулю threading, это равнозначно тому, что вы вкладываете что-то в руку Шиве с заглавной картинки поста. В данном случае в блоке кода
threads = [] for pair in pairs: # Подготавливаем потоки, складываем их в массив threads.append(threading.Thread(target=get_rates, args=(pair,)))
Я создал, грубо говоря отложенное задание threading.Thread(target=get_rates, args=(pair,)), и поместил его в массив threads. По аналогии, я соорудил гарпун и дал его Шиве, а тот схватил его в свободную руку.
В этом блоке кода
# Запускаем каждый поток for thread in threads: thread.start()
я запустил каждое задание – сказал Шиве запулить каждый гарпун в сторону каждой акулы. Шива запустил и уже начал тянуть обратно.
Я мог бы на этом закончить скрипт, тогда все запущенные потоки продолжали бы работать в фоне, но я хочу дождаться окончания каждого потока и посмотреть, что там получилось. Поэтому я добавляю третий блок кода:
# Ждем завершения каждого потока for thread in threads: thread.join()
Я как бы заявил – буду стоять тут и ждать, пока ты все гарпуны не вытащишь.
В итоге каждая команда запустилась независимо от других, и они начали соперничать за сетевую карту, процессор, оперативную память и т.п. – и порядок выполнения начали определять операционная система, процессор и GIL – такая глобальная штука в питоне, которая разруливает запущенные процессы.
После этого блока кода ничего не выполнится до тех пор, пока каждый поток не будет окончательно завершен. Основной процесс будет просто ждать, регулярно проверяя состояние дочерних потоков.
Вы можете использовать эти блоки как шаблон, меняться будет только название функции и передаваемые параметры. Внимание – если передаваемый параметр один, после него все равно должна стоять запятая, как в примере.
Получаем курсы с разных бирж
Давайте для закрепления сделаем что-то более-менее полезное – например, будем сравнивать курсы одной и той же пары на разных биржах. Для универсальности возьмем ETH_BTC – такая пара есть везде 🙂
Т.к. на каждой бирже своё API, то под каждую биржу создадим свою функцию. Так же нам понадобится глобальный объект (словарь, в данном случае), куда каждый поток будет складывать полученные данные.
Так же меняется принцип работы – в каждой функции свой бесконечный цикл, таким образом система получается следующая:
- Основной поток создает три дочерних потока, каждый из которых бесконечно получает последнюю цену со своей биржи, и складывает в глобальный словарь.
- Так же создается отдельный поток, который бесконечно выводит текущие содержимое глобального словаря.
Вот такой примерно результат работы:

А вот, собственно, код:
import time import requests import threading c1 = 'ETH' c2 = 'BTC' # Глобальный словарь, куда каждый поток складывает полученную информацию stock_rates = 'exmo':0, 'binance':0, 'bittrex': 0> # Получить последнюю цену с Эксмо def get_exmo_rates(pair): while True: try: stock_rates['exmo'] = requests.get("https://api.exmo.com/v1/ticker/".format(pair=pair)).json()[pair]['last_trade'] except Exception as e: print(e) time.sleep(0.5) # Получить последнюю цену с Binance def get_binance_rates(pair): while True: try: stock_rates['binance'] = requests.get("https://api.binance.com/api/v3/ticker/price?symbol= ".format(pair=pair)).json()['price'] except Exception as e: print(e) time.sleep(0.5) # Получить последнюю цену с Bittrex def get_bittrex_rates(pair): while True: try: stock_rates['bittrex'] = requests.get("https://bittrex.com/api/v1.1/public/getticker?market= ".format(pair=pair)).json()['result']['Last'] except Exception as e: print(e) time.sleep(0.5) def show_results(): while True: print(stock_rates) time.sleep(1) global_start_time = time.time() threads = [] # Подготавливаем потоки, складываем их в массив exmo_thread = threading.Thread(target=get_exmo_rates, args=(c1+'_'+c2,)) binance_thread = threading.Thread(target=get_binance_rates, args=(c1+c2,)) bittrex_thread = threading.Thread(target=get_bittrex_rates, args=(c2+'-'+c1,)) show_results_thread = threading.Thread(target=show_results) threads.append(exmo_thread) threads.append(binance_thread) threads.append(bittrex_thread) threads.append(show_results_thread) # Запускаем каждый поток for thread in threads: thread.start() # Ждем завершения каждого потока for thread in threads: thread.join()
Я немного по-другому создал потоки, что бы было нагляднее.
Заключение
Многие задачи можно решить и без многопоточности, например, запуская код последовательно или запуская разные экземпляры скриптов – и иногда так будет даже лучше. А иногда без многопототочности сложно.
Так же не стоить рассчитывать на неё как на панацею – если намечаются серьезные вычисления, и будет сильно задействован процессор, то многопоточность может наоборот сильно замедлить выполнение.
А вот если большей частью проходят операции ввода и вывода (запросы по сети, работа с оперативной памятью, чтение/запись с диска) то процессы могут очень хорошо организоваться и прирост производительности будет серьезный.
В общем, это еще одна фишка, которая может в какой-то момент очень круто пригодиться, об этом стоит знать, я считаю 🙂
Не забудьте рассказать друзьям об этой статье.
Чтобы поддержать ресурс Bablofil достаточно просто поделиться с друзьями этой статьей в социальных сетях. Каждый репост — это самая высокая оценка качества материала. Спасибо, что читаете этот блог.
Python в три ручья: работаем с потоками (часть 1)
В каких случаях вам нужна многопоточность, как реализовать её на Python и что нужно знать о глобальной блокировке GIL.
04 мая 2018 8 минут 272700

Автор статьи
Мария Лисянская

Автор статьи
Мария Лисянская
https://gbcdn.mrgcdn.ru/uploads/post/1582/og_cover_image/b8f59e927c66e03053113b8036b95103

Из этой статьи вы узнаете, как с Python выполнять несколько операций одновременно и распределять нагрузку между ядрами процессора, какие особенности языка учитывать. Но главное — поймете, когда многопоточность в Python нужна, а когда только мешает.
Небольшое предупреждение для тех, кто впервые слышит о параллельных вычислениях. Что такое поток и чем он отличается от процесса, мы выяснили в статье «Внутри процесса: многопоточность и пинг-понг mutex’ом». Тогда мы приводили примеры на Java, но теоретические основы многопоточности верны и для Python. Совпадают, в том числе, механизмы синхронизации потоков: семафоры, взаимные исключения (mutex), условия, события. Поэтому сегодня сделаем акцент на особенностях Python, его механизмах и инструментах, связанных с многопоточностью.
Организовать параллельные вычисления в Python без внешних библиотек можно с помощью модулей:
- threading — для управления потоками.
- queue — для организации очередей.
- multiprocessing — для управления процессами.
Пока нас интересует только первый пункт списка.
Как создавать потоки в Python
Метод 1 — «функциональный»
Для работы с потоками из модуля threading импортируем класс Thread. В начале кода пишем:
from threading import Thread
После этого нам будет доступна функция Thread() — с ней легко создавать потоки. Синтаксис такой:
variable = Thread(target=function_name, args=(arg1, arg2,))
Первый параметр target — это «целевая» функция, которая определяет поведение потока и создаётся заранее. Следом идёт список аргументов. Если судьбу аргументов (например, кто будет делимым, а кто делителем в уравнении) определяет их позиция, их записывают как args=(x,y). Если же вам нужны аргументы в виде пар «ключ-значение», используйте запись вида kwargs=.
Ради удобства отладки можно также дать новому потоку имя. Для этого среди параметров функции прописывают name=«Имя потока». По умолчанию name хранит значение null. А ещё потоки можно группировать с помощью параметра group, который по умолчанию — None.
За дело! Пусть два потока параллельно выводят каждый в свой файл заданное число строк. Для начала нам понадобится функция, которая выполнит задуманный нами сценарий. Аргументами целевой функции будут число строк и имя текстового файла для записи.
#coding: UTF-8 from threading import Thread def prescript(thefile, num): with open(thefile, 'w') as f: for i in range(num): if num > 500: f.write('МногоБукв\n') else: f.write('МалоБукв\n') thread1 = Thread(target=prescript, args=('f1.txt', 200,)) thread2 = Thread(target=prescript, args=('f2.txt', 1000,)) thread1.start() thread2.start() thread1.join() thread2.join()
Что start() запускает ранее созданный поток, вы уже догадались. Метод join() останавливает поток, когда тот выполнит свои задачи. Ведь нужно закрыть открытые файлы и освободить занятые ресурсы. Это называется «Уходя, гасите свет». Завершать потоки в предсказуемый момент и явно — надёжнее, чем снаружи и неизвестно когда. Меньше риск, что вмешаются случайные факторы. В качестве параметра в скобках можно указать, на сколько секунд блокировать поток перед продолжением его работы.
Метод 2 — «классовый»
Для потока со сложным поведением обычно пишут отдельный класс, который наследуют от Thread из модуля threading. В этом случае программу действий потока прописывают в методе run() созданного класса. Ту же петрушку мы видели и в Java.
#coding: UTF-8 import threading class MyThread(threading.Thread): def __init__(self, num): super().__init__(self, name="threddy" + num) self.num = num def run(self): print ("Thread ", self.num), thread1 = MyThread("1") thread2 = MyThread("2") thread1.start() thread2.start() thread1.join() thread2.join()
Стандартные методы работы с потоками

Чтобы управлять потоками, нужно следить, как они себя ведут. И для этого в threading есть специальные методы:
current_thread() — смотрим, какой поток вызвал функцию;
active_count() — считаем работающие в данный момент экземпляры класса Thread;
enumerate() — получаем список работающих потоков.
Ещё можно управлять потоком через методы класса:
is_alive() — спрашиваем поток: «Жив ещё, курилка?» — получаем true или false;
getName() — узнаём имя потока;
setName(any_name) — даём потоку имя;
У каждого потока, пока он работает, есть уникальный идентификационный номер, который хранится в переменной ident.
thread1.start() print(thread1.ident)
Отсрочить операции в вызываемых потоком функциях можно с помощью таймера. В инициализаторе объектов класса Timer всего два аргумента — время ожидания в секундах и функция, которую нужно в итоге выполнить:
import threading print ("Waiting. ") def timer_test(): print ("The timer has done its job!") tim = threading.Timer(5.0, timer_test) tim.start()
Таймер можно один раз создать, а затем запускать в разных частях кода.
Потусторонние потоки
Обычно Python-приложение не завершается, пока работает хоть один его поток. Но есть особые потоки, которые не мешают закрытию программы и останавливается вместе с ней. Их называют демонами (daemons). Проверить, является ли поток демоном, можно методом isDaemon(). Если является, метод вернёт истину.
Назначить поток демоном можно при создании — через параметр “daemon=True” или аргумент в инициализаторе класса.
thread0 = Thread(target=target_func, kwargs=10>, daemon=True)
Не поздно демонизировать и уже существующий поток методом setDaemon(daemonic).
Всё бы ничего, но это даже не верхушка айсберга, потому что прямо сейчас нас ждут великие открытия.
Приключение начинается. У древнего шлюза

Питон слывёт дружелюбным и простым в общении, но есть у него причуды. Нельзя просто взять и воспользоваться всеми преимуществами многопоточности в Python! Дорогу вам преградит огромный шлюз… Даже так — глобальный шлюз (Global Interpreter Lock, он же GIL), который ограничивает многопоточность на уровне интерпретатора. Технически, это один на всех mutex, созданный по умолчанию. Такого нет ни в C, ни в Java.
Задача шлюза — пропускать потоки строго по одному, чтоб не летали наперегонки, как печально известные стритрейсеры, и не создавали угрозу работе интерпретатора.
Без шлюза потоки подрезали бы друг друга, чтобы первыми добраться до памяти, но это еще не всё. Они имеют обыкновение внезапно засыпать за рулём! Операционная система не спрашивает, вовремя или невовремя — просто усыпляет их в ей одной известный момент. Из-за этого неупорядоченные потоки могут неожиданно перехватывать друг у друга инициативу в работе с общими ресурсами.
Дезориентированный спросонок поток, который видит перед собой совсем не ту ситуацию, при которой засыпал, рискует разбиться и повалить интерпретатор, либо попасть в тупиковую ситуацию (deadlock). Например, перед сном Поток 1 начал работу со списком, а после пробуждения не нашёл в этом списке элементов, т.к. их удалил или перезаписал Поток 2.
Чтобы такого не было, GIL в предсказуемый момент (по умолчанию раз в 5 миллисекунд для Python 3.2+) командует отработавшему потоку: «СПАААТЬ!» — тот отключается и не мешает проезжать следующему желающему. Даже если желающего нет, блокировщик всё равно подождёт, прежде чем вернуться к предыдущему активному потоку.

Благодаря шлюзу однопоточные приложения работают быстро, а потоки не конфликтуют. Но, к сожалению, многопоточные программы при таком подходе выполняются медленнее — слишком много времени уходит на регулировку «дорожного движения». А значит обработка графики, расчет математических моделей и поиск по большим массивам данных c GIL идут неприемлемо долго.
В статье «Understanding Python GIL»технический директор компании Gaglers Inc. и разработчик со стажем Chetan Giridhar приводит такой пример:
from datetime import datetime import threading def factorial(number): fact = 1 for n in range(1, number+1): fact *= n return fact number = 100000 thread = threading.Thread(target=factorial, args=(number,)) startTime = datetime.now() thread.start() thread.join() endTime = datetime.now() print "Время выполнения: ", endTime - startTime
Код вычисляет факториал числа 100 000 и показывает, сколько времени ушло у машины на эту задачу. При тестировании на одном ядре и с одним потоком вычисления заняли 3,4 секунды. Тогда Четан создал и запустил второй поток. Расчет факториала на двух ядрах длился 6,2 секунды. А ведь по логике скорость вычислений не должна была существенно измениться! Повторите этот эксперимент на своей машине и посмотрите, насколько медленнее будет решена задача, если вы добавите thread2. Я получила замедление ровно вдвое.
Глобальный шлюз — наследие времён, когда программисты боролись за достойную реализацию многозадачности и у них не очень получалось. Но зачем он сегодня, когда есть много- и очень многоядерные процессоры? Как объяснил Гвидо ван Россум, без GIL не будут нормально работать C-расширения для Python. Ещё упадёт производительность однопоточных приложений: Python 3 станет медленнее, чем Python 2, а это никому не нужно.
«Нормальные герои всегда идут в обход»

Шлюз можно временно отключить. Для этого интерпретатор Python нужно отвлечь вызовом функции из внешней библиотеки или обращением к операционной системе. Например, шлюз выключится на время сохранения или открытия файла. Помните наш пример с записью строк в файлы? Как только вызванная функция возвратит управление коду Python или интерфейсу Python C API, GIL снова включается.
Как вариант, для параллельных вычислений можно использовать процессы, которые работают изолированно и неподвластны GIL. Но это большая отдельная тема. Сейчас нам важнее найти решение для многопоточности.
Если вы собираетесь использовать Python для сложных научных расчётов, обойти скоростную проблему GIL помогут библиотеки Numba, NumPy, SciPy и др. Опишу некоторые из них в двух словах, чтобы вы поняли, стоит ли разведывать это направление дальше.
Numba для математики
Numba — динамически, «на лету» компилирует Python-код, превращая его в машинный код для исполнения на CPU и GPU. Такая технология компиляции называется JIT — “Just in time”. Она помогает оптимизировать производительность программ за счет ускорения работы циклов и компиляции функций при первом запуске.
Суть в том, что вы ставите аннотации (декораторы) в узких местах кода, где вам нужно ускорить работу функций.
Для математических расчётов библиотеку удобно использовать в связке c NumPy. Допустим, нужно сложить одномерные массивы — элемент за элементом.
def arr_sum (x , y): result_arr = nupmy.empty_like ( x) for i in range (len (x)) : result_arr [i ] = x[i ] + y[i ] return result_arr
Метод nupmy.empty_like() принимает массив и возвращает (но не инициализирует!) другой — соответствующий исходному по форме и типу. Чтобы ускорить выполнение кода, импортируем класс jit из модуля numba и добавляем в начало кода аннотацию @jit:
from numba import jit @jit def arr_sum(x,y):
Это скромное дополнение способно ускорить выполнение операции более чем в 100 раз! Если интересно, посмотрите замеры скорости математических расчётов при использовании разных библиотек для Python.
PyCUDA и Numba для графики
В графических вычислениях Numba тоже кое-что может. Она умеет работать с программной моделью CUDA, чтобы визуализировать научные данные и работу алгоритмов, выдавать информацию о GPU и др. Подробнее о том, как работают графический процессор и CUDA — здесь. И снова мы встретимся с многопоточностью.
При работе с многомерными массивами в CUDA, чтобы понять, какой поток сейчас работает с элементами массива, нужно отследить, кто и когда вызывает функцию ядра. Например, поток может определять свою позицию в сетке блоков и рассчитать соответствующий элемент массива:
from numba import cuda @cuda.jit def call_for_kernel(io_arr): # Идентификатор потока в одномерном блоке thread_x = cuda.threadIdx.x # Идентификатор блока в одномерной сетке thread_y = cuda.blockIdx.x # Число потоков на блок (т.е. ширина блока) block_width = cuda.blockDim.x # Находим положение в массиве t_position = thread_x + thread_y * block_width if t_position io_arr.size: # Убеждаемся, что не вышли за границы массива io_arr[ t_position] *= 2 # Считаем
Главный плюс этого кода даже не в скорости исполнения, а в прозрачности и простоте. Снова сошлюсь на Хабр, где есть сравнение скорости GPU-расчетов при использовании Numba, PyCUDA и эталонного С CUDA. Небольшой спойлер: PyCUDA позволяет достичь скорости вычислений, сопоставимой с Cи, а Numba подходит для небольших задач.
Когда многопоточность в Python оправдана
Стоит ли преодолевать связанные c GIL сложности и тратить время на реализацию многопоточности? Вот примеры ситуаций, когда многопоточность несёт с собой больше плюсов, чем минусов.
- Для длительных и несвязанных друг с другом операций ввода-вывода. Например, нужно обрабатывать ворох разрозненных запросов с большой задержкой на ожидание. В режиме «живой очереди» это долго — лучше распараллелить задачу.
- Вычисления занимают более миллисекунды и вы хотите сэкономить время за счёт их параллельного выполнения. Если операции укладываются в 1 мс, многопоточность не оправдает себя из-за высоких накладных расходов.
- Число потоков не превышает количество ядер. В противном случае параллельной работы всех потоков не получается и мы больше теряем, чем выигрываем.
Когда лучше с одним потоком

- При взаимозависимых вычислениях. Считать что-то в одном потоке и передавать для дальнейшей обработки второму — плохая идея. Возникает лишняя зависимость, которая приводит к снижению производительности, а в случае ошибки — к ступору и краху программы.
- При работе через GIL. Это мы уже выяснили выше.
- Когда важна хорошая переносимость на разных устройствах. Правильно подобрать число потоков для машины пользователя — задача не из легких. Если вы пишете под известное вам «железо», всё можно решить тестированием. Если же нет — понадобится дополнительно создавать гибкую систему подстройки под аппаратную часть, что потребует времени и умения.
Анонс — взаимные блокировки в Python
Самое смешное, что по умолчанию GIL защищает только интерпретатор и не предохраняет наш код от взаимных блокировок (deadlock) и других логических ошибок синхронизации. Поэтому разводить потоки по углам, как и в Java, нужно принудительно — с помощью блокирующих механизмов. Об этом и о не упомянутых в статье компонентах модуля threading мы поговорим в следующий раз.
