Потоки, процессы и async
У питониста для одновременной работы три инструмента: threading,
multiprocessing и asyncio. Прямых аналогов в Mojo 1.1 нет ни у одного,
и причины у всех трёх разные. Эта глава объясняет, что есть вместо них,
где кончаются возможности языка и какие ловушки ждут тех, кто
пришёл из Python.
| Python | Зачем он там | Mojo 1.1 |
|---|---|---|
threading | ожидание ввода-вывода; считать мешает GIL | потоков «вручную» нет; есть пул потоков: parallelize из пакета max |
concurrent.futures | pool.map(f, items) | ближе всего parallelize и sync_parallelize |
threading.Lock | защита общих данных | BlockingSpinLock из std.utils |
| — | в Python хватает GIL | атомарные операции: Atomic из std.atomic |
multiprocessing | обойти GIL | не нужен: GIL нет; внешние программы — Process.run и subprocess.run |
asyncio | тысячи одновременных соединений | в разработке: async def есть, но нестабилен, запускать задачи нечем |
Потоки без GIL
Заголовок раздела «Потоки без GIL»В Python потоки по очереди держат GIL (global interpreter lock): байт-код
в каждый момент исполняет только один из них. Для счёта потоки поэтому
бесполезны, и ради нескольких ядер приходится запускать процессы.
Зато GIL многое прощает: counter += 1 из двух потоков редко
ломается на практике.
В Mojo GIL нет. Потоки работают по-настоящему одновременно, каждый
на своём ядре, и multiprocessing для счёта не нужен. Цена у этого
одна: всё, что GIL прощал, теперь ломается.
Основной инструмент мы уже видели в главе
«Векторизация и параллелизм»:
parallelize(work, n) из пакета max вызывает work(0) … work(n - 1)
на пуле потоков и возвращается, когда все вызовы закончились. Вызовы
не раздаются по одному: n номеров делятся на столько подряд идущих
кусков, сколько потоков в пуле, — parallelism_level() из std.runtime.
Отдельной функции «запустить поток», как threading.Thread, в Mojo 1.1
нет — ни в стандартной библиотеке, ни в max. Нет и очередей, каналов,
событий. Модель одна: разделить работу, раздать потокам, дождаться
всех (fork-join). Для вычислений этого обычно хватает, а для
сервера, который ждёт сотни соединений, — нет.
Гонка данных: компилятор не спасёт
Заголовок раздела «Гонка данных: компилятор не спасёт»Возьмём самый наивный код: миллион вызовов из пула потоков увеличивают
общий счётчик. Замыкание захватывает его с mut, как в главе
«Замыкания и лямбды»:
from max.algorithm import parallelize
def main(): var total = 0
def work(i: Int) {mut total}: total += 1
parallelize(work, 1_000_000) print(total)928064
Должно быть 1 000 000. Пять запусков на двух ядрах дали 928 064,
536 622, 909 769, 953 996 и 851 678. Когда машина занята другой работой,
потоки иногда отрабатывают по очереди, и выходит ровно 1 000 000, — но
гонка от этого никуда не девается. Это гонка данных (data race):
total += 1 — три шага (прочитать, прибавить, записать), и два потока
читают одно и то же старое значение, после чего одно из прибавлений
теряется.
Главное здесь не число, а то, что код собрался без единого
предупреждения. В Rust такое не скомпилируется: правила заимствования
и трейты Send и Sync не дадут отдать изменяемую ссылку в несколько
потоков. В Mojo 1.1
аналогичной проверки нет. Захват {mut total} одинаково разрешён
и для обычного замыкания, и для такого, которое будут вызывать
из разных потоков.
С коллекциями ещё хуже. Если несколько потоков одновременно делают
append в один список, они вместе перевыделяют его память:
from max.algorithm import parallelize
def collect() -> List[Int]: var found = List[Int]()
def work(i: Int) {mut found}: if i % 3 == 0: found.append(i) # гонка: так делать нельзя
parallelize(work, 1_000_000) return found^Такой код обычно падает, и каждый раз по-своему: Segmentation fault
в List._realloc, сообщение аллокатора Possible double free detected,
неверное число элементов, а изредка и верный ответ. Об этом и предупреждала глава
«Указатели и память»: ошибки памяти
не обязаны проявляться одинаково.
Три способа сделать правильно
Заголовок раздела «Три способа сделать правильно»Атомарные операции
Заголовок раздела «Атомарные операции»Atomic из std.atomic выполняет «прочитать-изменить-записать» одной
неделимой инструкцией процессора:
from std.atomic import Atomicfrom max.algorithm import parallelize
def main(): var total = Atomic[Int](0)
def work(i: Int) {mut total}: total += 1 # то же, что total.fetch_add(1)
parallelize(work, 1_000_000) print(total.load())1000000
Ответ всегда верный. Atomic работает с числами — Int, Int64,
UInt32, Float64 и другими скалярными числовыми типами, но не с Bool,
структурами и строками. Флаг «уже найдено» заводите как Atomic[Int].
Кроме += и -= у него есть load, store, fetch_add, fetch_sub,
compare_exchange, max и min.
Блокировка
Заголовок раздела «Блокировка»Когда общих данных больше одного числа — например, список, — нужна
блокировка. В std.utils она называется BlockingSpinLock, а
BlockingScopedLock захватывает её на время блока with:
from max.algorithm import parallelizefrom std.utils import BlockingScopedLock, BlockingSpinLock
def main(): var found = List[Int]() var lock = BlockingSpinLock()
def work(i: Int) {mut found, mut lock}: if i % 100_000 == 0: with BlockingScopedLock(lock): # внутри — только один поток found.append(i)
parallelize(work, 1_000_000) sort(found) # порядок добавления зависит от потоков print(found)[0, 100000, 200000, 300000, 400000, 500000, 600000, 700000, 800000, 900000]
Это спин-блокировка (spin lock): поток, который ждёт, сначала
крутится в цикле и проверяет, не освободилась ли она, потом уступает
ядро и засыпает короткими паузами, около миллисекунды, снова проверяя.
Для коротких и редких участков это быстро. Но разбудить ждущий поток
сразу, как только блокировку отпустили, как это делает мьютекс
операционной системы, она не умеет. Такого мьютекса, как threading.Lock,
в стандартной библиотеке 1.1 нет.
Свой кусок каждому потоку
Заголовок раздела «Свой кусок каждому потоку»Лучший способ — не делить данные вовсе. Каждый поток считает свою часть в локальную переменную и в конце один раз пишет результат в свою ячейку списка:
from max.algorithm import parallelizefrom std.runtime import parallelism_level
def main(): var n = 1_000_000 var workers = parallelism_level() var partial = List[Int](length=workers, fill=0)
def work(w: Int) {mut partial, imm n, imm workers}: var s = 0 for i in range(w * n // workers, (w + 1) * n // workers): s += i % 7 partial[w] = s # у каждого потока своя ячейка
parallelize(work, workers) var total = 0 for p in partial: total += p print(total)2999997
Здесь parallelize получает не миллион номеров, а столько, сколько
потоков в пуле, и каждый вызов сам обходит свой диапазон. Запись
partial[w] безопасна: ячейки разные, а размер списка не меняется.
Сколько это стоит
Заголовок раздела «Сколько это стоит»Все три способа дают верный ответ. Разница — в скорости. Замер: сумма десяти миллионов дешёвых значений, лучшее из пяти прогонов:
"""Четыре способа посчитать сумму в несколько потоков. Вывод у каждой машины свой.
Запуск: mojo build bench_counter.mojo && ./bench_counter 10000000"""
from std.atomic import Atomicfrom std.benchmark import keepfrom std.runtime import parallelism_levelfrom std.sys import argvfrom std.time import perf_counter_nsfrom std.utils import BlockingScopedLock, BlockingSpinLockfrom max.algorithm import parallelize
def value(i: Int) -> Int: """Работа над одним элементом: дешёвая, но не сворачиваемая.""" return (i * 2654435761) % 1000
def one_thread(n: Int) -> Int: var total = 0 for i in range(n): total += value(i) return total
def atomic_each(n: Int) -> Int: var total = Atomic[Int](0)
def work(i: Int) {mut total}: total += value(i)
parallelize(work, n) return total.load()
def lock_each(n: Int) -> Int: var total = 0 var lock = BlockingSpinLock()
def work(i: Int) {mut total, mut lock}: var v = value(i) with BlockingScopedLock(lock): total += v
parallelize(work, n) return total
def partial_sums(n: Int) -> Int: var workers = parallelism_level() var partial = List[Int](length=workers, fill=0)
def work(w: Int) {mut partial, imm n, imm workers}: var s = 0 for i in range(w * n // workers, (w + 1) * n // workers): s += value(i) partial[w] = s # у каждого потока своя ячейка
parallelize(work, workers) var total = 0 for p in partial: total += p return total
def best_ms[f: def(Int) thin -> Int](n: Int, expected: Int) raises -> Float64: """Лучшее время из пяти запусков; заодно проверяем ответ.""" var best = Float64.MAX for _ in range(5): var t0 = perf_counter_ns() var result = f(n) var t1 = perf_counter_ns() keep(result) if result != expected: raise Error(String("неверный ответ: ", result, " вместо ", expected)) best = min(best, Float64(t1 - t0) / 1e6) return best
def main() raises: var n = Int(argv()[1]) if len(argv()) > 1 else 10_000_000 var expected = one_thread(n) print("элементов:", n, " потоков:", parallelism_level()) print("один поток ", round(best_ms[one_thread](n, expected), 2), "мс") print("атомик ", round(best_ms[atomic_each](n, expected), 2), "мс") print("блокировка ", round(best_ms[lock_each](n, expected), 2), "мс") print("частичные суммы ", round(best_ms[partial_sums](n, expected), 2), "мс")| Способ | Время | Относительно одного потока |
|---|---|---|
| один поток | 14,7–14,9 мс | 1 |
| атомик | 62–203 мс | в 4–14 раз медленнее |
| блокировка | 0,77–1,2 с | в 50–80 раз медленнее |
| частичные суммы | 7,6–11,5 мс | в 1,3–1,9 раза быстрее |
Замерено на двух ядрах (Intel Xeon 2,8 ГГц, облачный контейнер), несколько запусков программы на свободной машине. Разброс у атомика и блокировки большой: всё решает то, как часто потоки сталкиваются на одной ячейке, а это зависит от планировщика. Если машина занята чем-то ещё, частичные суммы перестают выигрывать вовсе: второе ядро программе не достаётся.
Вывод жёсткий: атомик и блокировка на каждый элемент съедают весь
выигрыш от потоков. Два ядра, которые непрерывно спорят за одну ячейку
памяти, работают медленнее одного. Atomic и блокировки хороши для
редких событий — посчитать найденное, добавить результат в список. Горячий
цикл должен работать с локальными данными.
parallelize и sync_parallelize
Заголовок раздела «parallelize и sync_parallelize»В пакете max две функции запуска:
parallelize(work, n) | sync_parallelize(work, n) | |
|---|---|---|
| как раздаёт | n номеров делит на куски по числу потоков | каждый номер — отдельная задача |
| замыкание | не может бросать исключения | может, но исключение обрывает программу |
| на сколько частей делит | parallelism_level() или третий аргумент | n задач |
| сколько потоков работает | сколько в пуле, третий аргумент этого не меняет | сколько в пуле |
Про исключения стоит запомнить отдельно. Ошибка внутри задачи
не долетит до вызывающего кода, и try вокруг sync_parallelize
не поможет. Компилятор даже предупредит, что except недостижим
('except' logic is unreachable, try doesn't raise an exception),
а при запуске программа оборвётся:
ABORT: max/mojo/max/algorithm/backend/cpu/parallelize.mojo:69:22: задача 3 сломаласьПоэтому ошибки внутри потока обрабатывайте там же, а наружу передавайте результат — например, флаг в своей ячейке списка.
Процессы
Заголовок раздела «Процессы»Раз GIL нет, multiprocessing для счёта не нужен: потоки и так
займут все ядра. А вот запустить другую программу иногда нужно.
Для этого в std.os есть Process: Process.run запускает программу
и не ждёт её, wait() дожидается конца и отдаёт код возврата.
# Три внешние программы работают одновременно; ждём каждую и читаем код возврата.from std.os import Processfrom std.time import perf_counter
def main() raises: var t0 = perf_counter() var jobs = List[Process]() for _ in range(3): jobs.append(Process.run("sleep", ["0.3"])) # запуск не ждёт конца
# List[Process] не перебрать через for: Process нельзя копировать for i in range(len(jobs)): var status = jobs[i].wait() print("задача", i, "код:", status.exit_code.value()) print("вместе, а не по очереди:", perf_counter() - t0 < 0.8)
var failed = Process.run("sh", ["-c", "exit 3"]) print("sh вернул:", failed.wait().exit_code.value())задача 0 код: 0 задача 1 код: 0 задача 2 код: 0 вместе, а не по очереди: True sh вернул: 3
Три программы по 0,3 секунды отработали одновременно. У Process есть
ещё poll() (проверить, не закончилась ли, без ожидания), kill(),
interrupt() и hangup(). Путь к программе ищется в PATH, как
в оболочке. Если объект Process уничтожается раньше, чем программа
закончилась, Mojo дождётся её конца: «забыть» про запущенный процесс
не выйдет.
Чего нет:
- перехватить вывод программы, запущенной через
Process.run, нельзя — для этого естьrunизstd.subprocess, но он ждёт конца команды (см. «Стандартная библиотека на каждый день»); - запустить в другом процессе функцию Mojo, как
Pool.mapвmultiprocessing, нельзя. Процесс — это отдельная программа, и данные между процессами передаются через файлы или аргументы командной строки.
Process.run работает на Linux и macOS.
Async: пока в разработке
Заголовок раздела «Async: пока в разработке»async def и await в Mojo есть как ключевые слова, но пользоваться
ими для одновременной работы в 1.1 нельзя. Вот что говорят об этом
сами разработчики:
- страница о стабильности:
асинхронная система не достроена,
asyncиawaitсчитаются нестабильными и могут измениться; - заметки к выпуску 1.1, раздел «Async APIs»:
типы корутин (
Coroutineи соседние) и API запуска задач изstd.runtime.asyncrtубраны из публичного доступа: на них начали строить код, хотя никаких гарантий стабильности у этого API нет; - дорожная карта: полноценный
async, встроенный в систему типов и модель памяти, стоит во втором этапе развития языка и помечен как ещё не начатый.
Код с async при этом компилируется. Но await сразу выполняет
корутину и ждёт её конца, так что всё идёт по очереди:
async def add(a: Int, b: Int) -> Int: return a + b
def main(): var x = await add(1, 2) # обычный вызов, никакой одновременности print(x)3
Запустить две корутины одновременно, как asyncio.gather, в 1.1 нечем.
А если забыть await, как это бывает в Python, компилятор напомнит:
error: 'c' abandoned without being explicitly destroyed: type 'Coroutine' does not conform to 'Deinitable' and must be explicitly destroyed
Вызов async-функции без await возвращает корутину — объект, который
ещё ничего не сделал. Выбросить его молча Mojo не даёт.
Добавьте await: var c = await add(1, 2). А лучше не пишите
async в коде, который должен жить дольше одного выпуска Mojo.
Что делать, если нужен сервер или тысячи одновременных соединений?
Честный ответ для Mojo 1.1 — эту часть лучше оставить Python
(asyncio) и вызывать из неё Mojo для тяжёлых вычислений, как в главе
«Вызов Mojo из Python». Сетевого кода
в стандартной библиотеке тоже нет.
Почему в курсе нет отдельных разделов про потоки и async
Заголовок раздела «Почему в курсе нет отдельных разделов про потоки и async»Потому что учить нечему, кроме того, что есть в этой главе. Mojo 1.1
даёт для CPU одну модель — пул потоков с parallelize, атомики
и спин-блокировку. Её мы разобрали здесь и в главе
«Векторизация и параллелизм».
А async официально нестабилен: глава о нём устарела бы с ближайшим
выпуском, а курс не должен учить тому, что вот-вот поменяется. Когда
в языке появятся полноценные async, потоки или каналы, появятся
и главы о них.
🎯 Проверь себя
Почему в Mojo не нужен multiprocessing для счёта на нескольких ядрах?
В Python процессы нужны, чтобы обойти GIL. В Mojo GIL нет: потоки
из parallelize работают одновременно, каждый на своём ядре, и делят
память — ничего не нужно пересылать между процессами.
Код `total += 1` внутри parallelize с захватом {mut total} скомпилировался. Значит, он правильный?
Нет. Компилятор Mojo 1.1 не проверяет, что данные безопасно делить между потоками. Это гонка данных: ответ может быть меньше ожидаемого и обычно разный от запуска к запуску.
Atomic даёт верный ответ. Почему бы не использовать его всегда?
Потому что потоки непрерывно спорят за одну ячейку памяти. На замере атомик на каждый элемент оказался в 4–14 раз медленнее одного потока. Быстро — когда каждый поток считает свою часть локально и один раз пишет её в свою ячейку.
Задача внутри sync_parallelize бросила исключение. Что увидит вызывающий код?
Ничего: программа завершится с ABORT. Исключение из задачи
не передаётся наружу, поэтому ошибки нужно обрабатывать внутри задачи.
Можно ли в Mojo 1.1 выполнить две async-функции одновременно?
Нет. async def и await компилируются, но await сразу выполняет
корутину и ждёт её конца. API для запуска задач в 1.1 убрали из публичного
доступа, а сама асинхронность официально нестабильна.
Раздел пройден
Заголовок раздела «Раздел пройден»Шесть глав — путь от «как измерить» до «как ускорить»:
| Глава | Чему учит |
|---|---|
| Как честно мерить скорость | не обманывать себя |
| SIMD с нуля | считать по многу чисел за раз |
| Векторизация и параллелизм | поручить это библиотеке и ядрам |
| Указатели и память | управлять памятью, когда контейнеров мало |
| Замыкания и лямбды | передавать поведение как значение |
| Потоки, процессы и async | делить работу между потоками и не сломать данные |
Главный вывод раздела — не про приёмы, а про порядок действий: сначала замер, потом оптимизация. Мы не раз видели, как правдоподобное объяснение рассыпалось от одного прогона: «векторизация даёт больше ширины вектора», «параллелизм всегда ускоряет», «атомик — дешёвый способ сделать правильно». Измеряйте.
Что дальше
Заголовок раздела «Что дальше»Дальше — как Mojo работает с чужим кодом, начиная с вызова Python из Mojo.
Тексты курса — CC BY-NC-SA 4.0, код примеров — Apache 2.0