Многопоточность
Base.Threads.@threadsМакрос
Threads.@threads [schedule] for ... end
Макрос для выполнения цикла for параллельно. Пространство итераций распределяется по крупным задачам. Данную политику можно указать с помощью аргумента schedule. Выполнение цикла ожидает завершения всех итераций.
См. также: @spawn и pmap в Distributed.
Расширенная справка
Семантика
За исключением более строгих гарантий, установленных опцией планирования, цикл, выполняемый макросом @threads, имеет следующую семантику.
Макрос @threads выполняет тело цикла в неопределённом порядке и потенциально параллельно. Он не определяет точное назначение задач и потоков. Назначение может быть различным для каждого выполнения. Код тела цикла (включая любой код, который вызывается из него) не должен делать предположений о распределении итераций по задачам или потоку, в котором они выполняются. Тело цикла для каждой итерации должно иметь возможность прогрессировать независимо от других итераций и быть свободным от гонок данных. Таким образом, некорректные синхронизации между итерациями могут привести к тупику, а несинхронизированные обращения к памяти могут привести к неопределённому поведению.
Например, вышеуказанные условия подразумевают:
- Блокировка, взятая в итерации, должна быть освобождена в той же итерации.
- Обмен данными между итерациями с помощью блокирующих примитивов, таких как
Channel, некорректен. - Запись только в места, не общие для итераций (если не используется блокировка или атомарная операция).
- За исключением использования расписания
:static, значениеthreadid()может измениться даже внутри одной итерации. См.Task Migration.
Планировщики
Без аргумента планировщика точное планирование не определено и зависит от версий Julia. В настоящее время используется :dynamic, когда планировщик не указан.
Аргумент schedule доступен начиная с Julia 1.5.
:dynamic (по умолчанию)
Планировщик :dynamic динамически распределяет итерации по доступным рабочим потокам. Текущая реализация предполагает, что нагрузка для каждой итерации равномерна. Однако это предположение может быть удалено в будущем.
Эта опция планирования является лишь подсказкой для основного механизма выполнения. Тем не менее, можно ожидать некоторых свойств. Количество Task используется планировщиком :dynamic ограничено небольшой константой, кратной количеству доступных рабочих потоков (Threads.threadpoolsize()). Каждая задача обрабатывает непрерывные области пространства итераций. Таким образом, @threads :dynamic for x in xs; f(x); end обычно более эффективен, чем @sync for x in xs; @spawn f(x); end, если length(xs) значительно больше количества рабочих потоков, а время выполнения f(x) относительно меньше стоимости запуска и синхронизации задачи (обычно менее 10 микросекунд).
Опция :dynamic для аргумента schedule доступна и является значением по умолчанию начиная с Julia 1.8.
:static
Планировщик :static создаёт одну задачу на поток и равномерно распределяет итерации между ними, назначая каждую задачу конкретному потоку. В частности, значение threadid() гарантированно остаётся постоянным в пределах одной итерации. Указание :static является ошибкой, если используется внутри другого цикла @threads или из потока, отличного от 1.
Планирование :static существует для поддержки перехода кода, написанного до Julia 1.3. В новых функциях библиотек планирование :static не рекомендуется, так как функции, использующие эту опцию, не могут быть вызваны из произвольных рабочих потоков.
Пример
Для иллюстрации различных стратегий планирования рассмотрим следующую функцию busywait, содержащую неблокирующий цикл с таймером, который работает заданное количество секунд.
julia> function busywait(seconds)
tstart = time_ns()
while (time_ns() - tstart) / 1e9 < seconds
end
end
julia> @time begin
Threads.@spawn busywait(5)
Threads.@threads :static for i in 1:Threads.threadpoolsize()
busywait(1)
end
end
6.003001 seconds (16.33 k allocations: 899.255 KiB, 0.25% compilation time)
julia> @time begin
Threads.@spawn busywait(5)
Threads.@threads :dynamic for i in 1:Threads.threadpoolsize()
busywait(1)
end
end
2.012056 seconds (16.05 k allocations: 883.919 KiB, 0.66% compilation time)
Пример :dynamic занимает 2 секунды, так как один из не занятых потоков может выполнить две итерации по 1 секунде, чтобы завершить цикл.
Base.Threads.foreachФункция
Threads.foreach(f, channel::Channel;
schedule::Threads.AbstractSchedule=Threads.FairSchedule(),
ntasks=Threads.threadpoolsize())
Аналогично foreach(f, channel), но итерация по channel и вызовы к f распределяются по ntasks задачам, созданным Threads.@spawn. Эта функция ожидает завершения всех созданных задач перед возвратом.
Если schedule isa FairSchedule, Threads.foreach попытается создать задачи таким образом, чтобы планировщик Julia мог более свободно распределять задачи по потокам. Этот подход, как правило, имеет более высокую накладные расходы на элемент, но может работать лучше, чем StaticSchedule при одновременной работе с другими многопоточными задачами.
Если schedule isa StaticSchedule, Threads.foreach создаст задачи таким образом, чтобы накладные расходы на элемент были ниже, чем у FairSchedule, но это менее способствует балансировке нагрузки. Таким образом, этот подход может быть более подходящим для мелкозернистых, равномерно распределённых задач, но может работать хуже, чем FairSchedule при одновременной работе с другими многопоточными задачами.
Примеры
julia> n = 20
julia> c = Channel{Int}(ch -> foreach(i -> put!(ch, i), 1:n), 1)
julia> d = Channel{Int}(n) do ch
f = i -> put!(ch, i^2)
Threads.foreach(f, c)
end
julia> collect(d)
collect(d) = [1, 4, 9, 16, 25, 36, 49, 64, 81, 100, 121, 144, 169, 196, 225, 256, 289, 324, 361, 400]
Эта функция требует Julia 1.6 или более поздней версии.
Base.Threads.@spawnМакрос
Threads.@spawn [:default|:interactive] expr
Создать Task и schedule его для выполнения в любом доступном потоке в указанном пуле потоков (:default если не указано). Задача назначается потоку, как только он станет доступным. Для ожидания завершения задачи вызовите wait на результате этого макроса или вызовите fetch для ожидания и получения его возвращаемого значения.
Значения могут быть интерполированы в @spawn с помощью $, который копирует значение непосредственно в созданный базовый закрытый метод. Это позволяет вам вставлять значение переменной, изолируя асинхронный код от изменений значения переменной в текущей задаче.
Поток, в котором выполняется задача, может измениться, если задача приостановится, поэтому threadid() не следует рассматривать как постоянную для задачи. См. Task Migration, а также более широкую документацию по многопоточности многопоточность. См. также главу о пулах потоков.
Этот макрос доступен начиная с Julia 1.3.
Интерполяция значений с помощью $ доступна начиная с Julia 1.4.
Пул потоков может быть задан начиная с Julia 1.9.
Примеры
julia> t() = println("Hello from ", Threads.threadid());
julia> tasks = fetch.([Threads.@spawn t() for i in 1:4]);
Hello from 1
Hello from 1
Hello from 3
Hello from 4
исходный код
Base.Threads.threadidФункция
Threads.threadid() -> Int
Получить идентификатор текущего потока выполнения. Главный поток имеет идентификатор 1.
Примеры
julia> Threads.threadid()
1
julia> Threads.@threads for i in 1:4
println(Threads.threadid())
end
4
2
5
4
Поток, в котором выполняется задача, может измениться, если задача приостановится, что известно как Task Migration. По этой причине в большинстве случаев использование threadid() для индексации, например, в векторе буферов или объектов состояния, небезопасно.
Base.Threads.maxthreadidФункция
Threads.maxthreadid() -> Int
Получить нижнюю границу количества потоков (во всех пулах потоков), доступных процессу Julia, с помощью семантики атомарного захвата. Результат всегда будет больше или равен threadid(), а также threadid(task) для любой задачи, которую вы смогли наблюдать до вызова maxthreadid.
Base.Threads.nthreadsФункция
Threads.nthreads(:default | :interactive) -> Int
Получить текущее количество потоков в указанном пуле потоков. Потоки в :interactive имеют идентификаторы 1:nthreads(:interactive), а потоки в :default имеют идентификаторы в nthreads(:interactive) .+ (1:nthreads(:default)).
См. также BLAS.get_num_threads и BLAS.set_num_threads в стандартной библиотеке LinearAlgebra, а также nprocs() в стандартной библиотеке Distributed и Threads.maxthreadid().
Base.Threads.threadpoolФункция
Threads.threadpool(tid = threadid()) -> Symbol
Возвращает пул потоков указанного потока; либо :default, :interactive, или :foreign.
Base.Threads.nthreadpoolsФункция
Threads.nthreadpools() -> Int
Возвращает количество пулов потоков, которые в настоящее время настроены.
исходный код
Base.Threads.threadpoolsizeФункция
Threads.threadpoolsize(pool::Symbol = :default) -> Int
Получает количество потоков, доступных для основного пула потоков (или для указанного пула потоков).
См. также: BLAS.get_num_threads и BLAS.set_num_threads в стандартной библиотеке LinearAlgebra и nprocs() в стандартной библиотеке Distributed.
Base.Threads.ngcthreadsФункция
Threads.ngcthreads() -> Int
Возвращает количество потоков GC, которые в настоящее время настроены. Это включает в себя как потоки маркировки, так и потоки одновременной очистки.
исходный кодСм. также Многопоточность.
Атомарные операции
atomicКлючевое слово
Небезопасные операции с указателями совместимы с загрузкой и сохранением указателей, объявленных с помощью типа _Atomic и std::atomic соответственно в C11 и C++23. Может быть выброшено исключение, если нет поддержки атомарной загрузки типа Julia T.
См. также: unsafe_load, unsafe_modify!, unsafe_replace!, unsafe_store!, unsafe_swap!
Base.@atomicМакрос
@atomic var @atomic order ex
Отметить var или ex как выполняемые атомарно, если ex является поддерживаемым выражением. Если order не указано, по умолчанию используется :sequentially_consistent.
@atomic a.b.x = new @atomic a.b.x += addend @atomic :release a.b.x = new @atomic :acquire_release a.b.x += addend
Выполнить операцию сохранения, выраженную справа, атомарно и вернуть новое значение.
С =, эта операция преобразуется в вызов setproperty!(a.b, :x, new). С любым оператором эта операция преобразуется в вызов modifyproperty!(a.b, :x, +, addend)[2].
@atomic a.b.x max arg2 @atomic a.b.x + arg2 @atomic max(a.b.x, arg2) @atomic :acquire_release max(a.b.x, arg2) @atomic :acquire_release a.b.x + arg2 @atomic :acquire_release a.b.x max arg2
Выполнить бинарную операцию, выраженную справа, атомарно. Сохранить результат в поле в первом аргументе и вернуть значения (old, new).
Эта операция преобразуется в вызов modifyproperty!(a.b, :x, func, arg2).
См. раздел Атомарные операции по полям в руководстве для получения более подробной информации.
Примеры
julia> mutable struct Atomic{T}; @atomic x::T; end
julia> a = Atomic(1)
Atomic{Int64}(1)
julia> @atomic a.x # fetch field x of a, with sequential consistency
1
julia> @atomic :sequentially_consistent a.x = 2 # set field x of a, with sequential consistency
2
julia> @atomic a.x += 1 # increment field x of a, with sequential consistency
3
julia> @atomic a.x + 1 # increment field x of a, with sequential consistency
3 => 4
julia> @atomic a.x # fetch field x of a, with sequential consistency
4
julia> @atomic max(a.x, 10) # change field x of a to the max value, with sequential consistency
4 => 10
julia> @atomic a.x max 5 # again change field x of a to the max value, with sequential consistency
10 => 10
Для этой функции требуется как минимум Julia 1.7.
Base.@atomicswapМакрос
@atomicswap a.b.x = new @atomicswap :sequentially_consistent a.b.x = new
Сохраняет new в a.b.x и возвращает старое значение a.b.x.
Эта операция преобразуется в вызов swapproperty!(a.b, :x, new).
См. раздел Атомарные операции по полям в руководстве для получения более подробной информации.
Примеры
julia> mutable struct Atomic{T}; @atomic x::T; end
julia> a = Atomic(1)
Atomic{Int64}(1)
julia> @atomicswap a.x = 2+2 # replace field x of a with 4, with sequential consistency
1
julia> @atomic a.x # fetch field x of a, with sequential consistency
4
Для этой функции требуется как минимум Julia 1.7.
Base.@atomicreplaceМакрос
@atomicreplace a.b.x expected => desired @atomicreplace :sequentially_consistent a.b.x expected => desired @atomicreplace :sequentially_consistent :monotonic a.b.x expected => desired
Выполнить условную замену, выраженную парой атомарно, вернув значения (old, success::Bool). Где success указывает, была ли замена завершена.
Эта операция преобразуется в вызов replaceproperty!(a.b, :x, expected, desired).
См. раздел Атомарные операции по полям в руководстве для получения более подробной информации.
Примеры
julia> mutable struct Atomic{T}; @atomic x::T; end
julia> a = Atomic(1)
Atomic{Int64}(1)
julia> @atomicreplace a.x 1 => 2 # replace field x of a with 2 if it was 1, with sequential consistency
(old = 1, success = true)
julia> @atomic a.x # fetch field x of a, with sequential consistency
2
julia> @atomicreplace a.x 1 => 2 # replace field x of a with 2 if it was 1, with sequential consistency
(old = 2, success = false)
julia> xchg = 2 => 0; # replace field x of a with 0 if it was 2, with sequential consistency
julia> @atomicreplace a.x xchg
(old = 2, success = true)
julia> @atomic a.x # fetch field x of a, with sequential consistency
0
Для этой функции требуется как минимум Julia 1.7.
Следующие API довольно примитивны и, вероятно, будут представлены через оболочку типа unsafe_*.
Core.Intrinsics.atomic_pointerref(pointer::Ptr{T}, order::Symbol) --> T
Core.Intrinsics.atomic_pointerset(pointer::Ptr{T}, new::T, order::Symbol) --> pointer
Core.Intrinsics.atomic_pointerswap(pointer::Ptr{T}, new::T, order::Symbol) --> old
Core.Intrinsics.atomic_pointermodify(pointer::Ptr{T}, function::(old::T,arg::S)->T, arg::S, order::Symbol) --> old
Core.Intrinsics.atomic_pointerreplace(pointer::Ptr{T}, expected::Any, new::T, success_order::Symbol, failure_order::Symbol) --> (old, cmp)
Следующие API устарели, хотя поддержка их, вероятно, сохранится на несколько релизов.
Base.Threads.AtomicТип
Threads.Atomic{T}
Содержит ссылку на объект типа T, гарантируя, что к нему обращаются только атомарно, т.е. безопасным для многопоточности способом.
Только некоторые "простые" типы могут использоваться атомарно, а именно примитивные логические, целочисленные и типы с плавающей точкой. Это Bool, Int8...Int128, UInt8...UInt128, и Float16...Float64.
Новые атомарные объекты могут быть созданы из значений, не являющихся атомарными; если ни один не указан, атомарный объект инициализируется нулём.
К атомарным объектам можно обращаться с помощью обозначения []:
Примеры
julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)
julia> x[] = 1
1
julia> x[]
1
Атомарные операции используют префикс atomic_ , такой как atomic_add!, atomic_xchg! и т.д.
Base.Threads.atomic_cas!Функция
Threads.atomic_cas!(x::Atomic{T}, cmp::T, newval::T) where T
Атомарно сравнивает и устанавливает x
Атомарно сравнивает значение в x со значением cmp. Если они равны, записывает newval в x. В противном случае, оставляет x неизменным. Возвращает старое значение в x. Сравнивая возвращаемое значение с cmp (через ===) можно узнать, было ли x модифицировано и теперь содержит новое значение newval.
Для получения более подробной информации см. инструкцию LLVM cmpxchg.
Эта функция может использоваться для реализации транзакционных семантик. Перед транзакцией записывается значение в x. После транзакции новое значение хранится только в том случае, если x не был изменен тем временем.
Примеры
julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)
julia> Threads.atomic_cas!(x, 4, 2);
julia> x
Base.Threads.Atomic{Int64}(3)
julia> Threads.atomic_cas!(x, 3, 2);
julia> x
Base.Threads.Atomic{Int64}(2)
исходный код
Base.Threads.atomic_xchg!Функция
Threads.atomic_xchg!(x::Atomic{T}, newval::T) where T
Атомарно обменивает значение в x
Атомарно обменивает значение в x со значением newval. Возвращает старое значение.
Для получения более подробной информации см. инструкцию LLVM atomicrmw xchg.
Примеры
julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)
julia> Threads.atomic_xchg!(x, 2)
3
julia> x[]
2
исходный код
Base.Threads.atomic_add!Функция
Threads.atomic_add!(x::Atomic{T}, val::T) where T <: ArithmeticTypes
Атомарно добавляет val к x
Выполняет x[] += val атомарно. Возвращает старое значение. Не определено для Atomic{Bool}.
Для получения более подробной информации см. инструкцию LLVM atomicrmw add.
Примеры
julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)
julia> Threads.atomic_add!(x, 2)
3
julia> x[]
5
исходный код
Base.Threads.atomic_sub!Функция
Threads.atomic_sub!(x::Atomic{T}, val::T) where T <: ArithmeticTypes
Атомарно вычитает val из x
Выполняет x[] -= val атомарно. Возвращает старое значение. Не определено для Atomic{Bool}.
Для получения более подробной информации см. инструкцию LLVM atomicrmw sub.
Примеры
julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)
julia> Threads.atomic_sub!(x, 2)
3
julia> x[]
1
исходный код
Base.Threads.atomic_and!Функция
Threads.atomic_and!(x::Atomic{T}, val::T) where T
Атомарное побитовое И x с val
Выполняется x[] &= val атомарно. Возвращает значение old.
Для получения более подробной информации см. инструкцию LLVM atomicrmw and.
Примеры
julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)
julia> Threads.atomic_and!(x, 2)
3
julia> x[]
2
исходный код
Base.Threads.atomic_nand!Функция
Threads.atomic_nand!(x::Atomic{T}, val::T) where T
Атомарное побитовое ИЛИ-НЕ (не-и) x с val
Выполняется x[] = ~(x[] & val) атомарно. Возвращает значение old.
Для получения более подробной информации см. инструкцию LLVM atomicrmw nand.
Примеры
julia> x = Threads.Atomic{Int}(3)
Base.Threads.Atomic{Int64}(3)
julia> Threads.atomic_nand!(x, 2)
3
julia> x[]
-3
исходный код
Base.Threads.atomic_or!Функция
Threads.atomic_or!(x::Atomic{T}, val::T) where T
Атомарное побитовое ИЛИ x с val
Выполняется x[] |= val атомарно. Возвращает значение old.
Для получения более подробной информации см. инструкцию LLVM atomicrmw or.
Примеры
julia> x = Threads.Atomic{Int}(5)
Base.Threads.Atomic{Int64}(5)
julia> Threads.atomic_or!(x, 7)
5
julia> x[]
7
исходный код
Base.Threads.atomic_xor!Функция
Threads.atomic_xor!(x::Atomic{T}, val::T) where T
Атомарное побитовое ИСКЛЮЧАЮЩЕЕ ИЛИ (XOR) x с val
Выполняется x[] $= val атомарно. Возвращает значение old.
Для получения более подробной информации см. инструкцию LLVM atomicrmw xor.
Примеры
julia> x = Threads.Atomic{Int}(5)
Base.Threads.Atomic{Int64}(5)
julia> Threads.atomic_xor!(x, 7)
5
julia> x[]
2
исходный код
Base.Threads.atomic_max!Функция
Threads.atomic_max!(x::Atomic{T}, val::T) where T
Атомарно сохраняет максимальное значение из x и val в x
Выполняется x[] = max(x[], val) атомарно. Возвращает значение old.
Для получения более подробной информации см. инструкцию LLVM atomicrmw max.
Примеры
julia> x = Threads.Atomic{Int}(5)
Base.Threads.Atomic{Int64}(5)
julia> Threads.atomic_max!(x, 7)
5
julia> x[]
7
исходный код
Base.Threads.atomic_min!Функция
Threads.atomic_min!(x::Atomic{T}, val::T) where T
Атомарно сохраняет минимальное значение из x и val в x
Выполняется x[] = min(x[], val) атомарно. Возвращает значение old.
Для получения более подробной информации см. инструкцию LLVM atomicrmw min.
Примеры
julia> x = Threads.Atomic{Int}(7)
Base.Threads.Atomic{Int64}(7)
julia> Threads.atomic_min!(x, 5)
7
julia> x[]
5
исходный код
Base.Threads.atomic_fenceФункция
Threads.atomic_fence()
Вставка барьера последовательной согласованности памяти
Вставляет барьер памяти с семантикой последовательной согласованности. В некоторых алгоритмах это необходимо, т.е. когда порядка приобретения/освобождения недостаточно.
Вероятно, эта операция очень ресурсоемкая. Учитывая, что все другие атомарные операции в Julia уже имеют семантику приобретения/освобождения, явные барьеры в большинстве случаев не нужны.
Для получения более подробной информации см. инструкцию LLVM fence.
ccall с использованием пула потоков libuv (Экспериментальная функция)
Base.@threadcallМакрос
@threadcall((cfunc, clib), rettype, (argtypes...), argvals...)
Макрос @threadcall вызывается так же, как ccall, но выполняет работу в другом потоке. Это полезно, когда вы хотите вызвать блокирующую C-функцию, не блокируя текущий julia поток. Одновременность ограничена размером пула потоков libuv, который по умолчанию составляет 4 потока, но его можно увеличить, установив переменную окружения UV_THREADPOOL_SIZE и перезапустив процесс julia.
Обратите внимание, что вызываемая функция никогда не должна вызывать обратно в Julia.
исходный кодПримитивы низкоуровневой синхронизации
Эти строительные блоки используются для создания обычных объектов синхронизации.
Base.Threads.SpinLockТип
SpinLock()
Создает нерекурсивную блокировку вращения с проверкой и установкой. Рекурсивное использование приведет к тупиковой ситуации. Этот тип блокировки следует использовать только для кода, выполняющегося за короткое время и не блокирующего (например, выполняющего ввод-вывод). В общем случае, следует использовать ReentrantLock.
Каждый lock должен соответствовать unlock. Если !islocked(lck::SpinLock) активен, trylock(lck) выполняется успешно, если нет других задач, пытающихся захватить блокировку «одновременно».
Блокировки вращения с проверкой и проверкой и установкой работают быстрее до примерно 30 конкурирующих потоков. Если у вас больше конкурирующих потоков, следует рассмотреть другие подходы к синхронизации.
исходный код
© 2009–2024 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.10/base/multi-threading/