Spec-Zone.ru › Julia 1.8

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

Base.Threads.@threadsМакрос

Threads.@threads [schedule] for ... end

Макрос для выполнения цикла for параллельно. Пространство итераций распределяется по крупнозернистым задачам. Данная политика может быть задана аргументом schedule. Выполнение цикла ждёт оценки всех итераций.

См. также: @spawn и pmap в Distributed.

Подробная справка

Семантика

За исключением случаев, когда более сильные гарантии задаются параметром планирования, цикл, выполняемый макросом @threads, имеет следующую семантику.

Макрос @threads выполняет тело цикла в неопределённом порядке и потенциально одновременно. Он не указывает точное распределение задач и потоков исполнителей. Распределение может отличаться для каждого выполнения. Тело цикла (включая любой код, вызываемый из него) не должно делать предположений о распределении итераций по задачам или потокам исполнителей, в которых они выполняются. Тело цикла для каждой итерации должно быть способно к прогрессу независимо от других итераций и свободно от гонок данных. Таким образом, неверные синхронизации между итерациями могут привести к тупику, а несинхронизированные обращения к памяти — к неопределённому поведению.

Например, вышеуказанные условия подразумевают, что:

  • Блокировка, взятая в итерации, должна быть освобождена в той же итерации.
  • Обмен данными между итерациями с помощью блокирующих примитивов, таких как Channels, некорректен.
  • Записи должны производиться только в локациях, не общих для итераций (если не используется блокировка или атомарная операция).
  • Значение threadid() может меняться даже в пределах одной итерации.

Планировщики

Без аргумента планировщика точное планирование не определено и варьируется между версиями Julia. В настоящее время используется :dynamic, когда планировщик не указан.

Аргумент schedule доступен начиная с Julia 1.5.

:dynamic (по умолчанию)

Планировщик :dynamic динамически выполняет итерации на доступных потоках исполнителей. Текущая реализация предполагает, что нагрузка для каждой итерации равномерна. Однако это предположение может быть удалено в будущем.

Этот параметр планирования является лишь подсказкой для механизма выполнения. Однако можно ожидать нескольких свойств. Количество Tasks, используемых планировщиком :dynamic, ограничено небольшой константой, умноженной на количество доступных потоков исполнителей (nthreads()). Каждая задача обрабатывает смежные области пространства итераций. Таким образом, @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.nthreads()
                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.nthreads()
                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.nthreads())

Аналогично foreach(f, channel), но итерация по channel и вызовы f разбиваются на ntasks задач, запущенных Threads.@spawn. Эта функция подождёт завершения всех внутренних задач, прежде чем вернуть результат.

Если schedule isa FairSchedule, Threads.foreach попытается запустить задачи таким образом, чтобы планировщик Julia мог более свободно распределять задания по потокам. Этот подход, как правило, имеет более высокую накладные расходы на элемент, но может работать лучше, чем StaticSchedule при одновременной работе с другими многопоточными задачами.

Если schedule isa StaticSchedule, Threads.foreach запустит задачи таким образом, чтобы накладные расходы на элемент были ниже, чем у FairSchedule, но при этом он хуже подходит для распределения нагрузки. Таким образом, этот подход может быть более подходящим для мелкозернистых однородных задач, но может работать хуже, чем FairSchedule при одновременной работе с другими многопоточными задачами.

Для этой функции требуется Julia 1.6 или новее.

исходный код

Base.Threads.@spawnМакрос

Threads.@spawn expr

Создаёт Task и планирует его для выполнения в любом доступном потоке. Задача назначается потоку после его освобождения. Чтобы подождать завершения задачи, вызовите wait на результате этого макроса или вызовите fetch, чтобы подождать и получить его возвращаемое значение.

Значения могут быть интерполированы в @spawn с помощью $, что напрямую копирует значение в созданный базовый замыкание. Это позволяет вставить значение переменной, изолируя асинхронный код от изменений значения переменной в текущей задаче.

См. главу руководства по многопоточности для важных замечаний.

Этот макрос доступен начиная с Julia 1.3.

Интерполяция значений с помощью $ доступна начиная с Julia 1.4.

исходный код

Base.Threads.threadidФункция

Threads.threadid()

Получить идентификатор текущего потока выполнения. Главный поток имеет идентификатор 1.

исходный код

Base.Threads.nthreadsФункция

Threads.nthreads()

Получить количество потоков, доступных для процесса Julia. Это верхняя граница, включая threadid().

См. также: BLAS.get_num_threads и BLAS.set_num_threads в стандартной библиотеке LinearAlgebra и nprocs() в стандартной библиотеке Distributed.

исходный код

См. также Многопоточность.

Атомарные операции

Base.@atomicМакрос

@atomic var
@atomic order ex

Отметить var или ex как выполняемые атомарно, если ex является поддерживаемым выражением.

@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).

Дополнительные сведения см. в разделе «Атомарные операции по полю» («Per-field atomics») руководства по ссылке здесь.

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 1, 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, гарантируя, что к нему обращаются только атомарно, т. е. безопасным для потоков способом.

Для атомарного использования могут быть использованы только определённые «простые» типы, а именно примитивные типы boolean, integer и float. Это 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
исходный код

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

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

Атомарное побитовое ИСКЛЮЧАЮЩЕЕ ИЛИ 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 с использованием пула потоков (Экспериментально)

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.

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

исходный код

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

Spec-Zone.ru

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