Spec-Zone.ru › Julia 1.9

Многопоточность

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.

исходный код

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

Получить текущее количество потоков в указанном пуле потоков. Потоки по умолчанию имеют номера идентификаторов 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.

исходный код

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.@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 атомарно. Возвращает старое значение.

Для получения дополнительной информации см. инструкцию 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) атомарно. Возвращает старое значение.

Для получения дополнительной информации см. инструкцию 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 атомарно. Возвращает старое значение.

Для получения дополнительной информации см. инструкцию LLVM atomicrmw or.

Примеры

julia> x = Threads.Atomic{Int}(5)
Base.Threads.Atomic{Int64}(5)

julia> Threads.atomic_or!(x, 7)
5

julia> x[]
7
исходный код
END_OF_DOCUMENT_MARKER

Base.Threads.atomic_xor!Функция

Threads.atomic_xor!(x::Atomic{T}, val::T) where T

Атомарное побитовое исключающее ИЛИ (XOR) x с val

Выполняет x[] $= val атомарно. Возвращает старое значение.

Для получения дополнительной информации см. инструкцию 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) атомарно. Возвращает старое значение.

Для получения дополнительной информации см. инструкцию 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) атомарно. Возвращает старое значение.

Для получения дополнительной информации см. инструкцию 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()

Создаёт нерекурсивную блокировку вращения (spin lock), использующую механизм "проверка-проверка-установка". Рекурсивное использование приведёт к тупику. Этот тип блокировки следует использовать только для кода, который выполняется быстро и не блокирует выполнение (например, для операций ввода-вывода). В общем случае следует использовать ReentrantLock.

Каждый вызов lock должен быть сопоставлен с вызовом unlock. Если !islocked(lck::SpinLock) возвращает true, то trylock(lck) завершится успешно, если нет других задач, пытающихся получить блокировку "одновременно".

Блокировки вращения с механизмом "проверка-проверка-установка" наиболее эффективны при до 30-ти конкурирующих потоков. Если у вас больше конкурирующих потоков, следует рассмотреть другие методы синхронизации.

исходный код

© 2009–2023 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.9/base/multi-threading/

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API