Перейти к содержимому

Потоки, процессы и async

У питониста для одновременной работы три инструмента: threading, multiprocessing и asyncio. Прямых аналогов в Mojo 1.1 нет ни у одного, и причины у всех трёх разные. Эта глава объясняет, что есть вместо них, где кончаются возможности языка и какие ловушки ждут тех, кто пришёл из Python.

PythonЗачем он тамMojo 1.1
threadingожидание ввода-вывода; считать мешает GILпотоков «вручную» нет; есть пул потоков: parallelize из пакета max
concurrent.futurespool.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 есть, но нестабилен, запускать задачи нечем

В 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 Atomic
from 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 parallelize
from 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 parallelize
from 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] безопасна: ячейки разные, а размер списка не меняется.

Все три способа дают верный ответ. Разница — в скорости. Замер: сумма десяти миллионов дешёвых значений, лучшее из пяти прогонов:

bench_counter.mojo
"""Четыре способа посчитать сумму в несколько потоков. Вывод у каждой машины свой.
Запуск: mojo build bench_counter.mojo && ./bench_counter 10000000
"""
from std.atomic import Atomic
from std.benchmark import keep
from std.runtime import parallelism_level
from std.sys import argv
from std.time import perf_counter_ns
from std.utils import BlockingScopedLock, BlockingSpinLock
from 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 и блокировки хороши для редких событий — посчитать найденное, добавить результат в список. Горячий цикл должен работать с локальными данными.

В пакете 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() дожидается конца и отдаёт код возврата.

processes.mojo
# Три внешние программы работают одновременно; ждём каждую и читаем код возврата.
from std.os import Process
from 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 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.

Примеры проверены на Mojo 1.1.0

Тексты курса — CC BY-NC-SA 4.0, код примеров — Apache 2.0