Spec-Zone.ru › Julia 1.9

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

Посетите эту статью блога для ознакомления с функциями многопоточности Julia.

Запуск Julia с несколькими потоками

По умолчанию Julia запускается с одним потоком выполнения. Это можно проверить, используя команду Threads.nthreads():

julia> Threads.nthreads()
1

Количество потоков выполнения контролируется либо с помощью аргумента командной строки -t/--threads, либо с помощью переменной окружения JULIA_NUM_THREADS. Если оба указаны, то -t/--threads имеет приоритет.

Количество потоков может быть указано либо как целое число (--threads=4) либо как auto (--threads=auto), где auto пытается вывести полезное значение по умолчанию для количества потоков (подробнее см. Параметры командной строки).

Аргумент командной строки -t/--threads требует по крайней мере Julia 1.5. В более старых версиях необходимо использовать переменную окружения.

Использование auto в качестве значения переменной окружения JULIA_NUM_THREADS требует по крайней мере Julia 1.7. В более старых версиях это значение игнорируется.

Давайте запустим Julia с 4 потоками:

$ julia --threads 4

Давайте проверим, что у нас есть 4 потока.

julia> Threads.nthreads()
4

Но мы сейчас находимся в главном потоке. Чтобы проверить, используем функцию Threads.threadid

julia> Threads.threadid()
1

Если вы предпочитаете использовать переменную окружения, вы можете установить её следующим образом в Bash (Linux/macOS):

export JULIA_NUM_THREADS=4

C shell на Linux/macOS, CMD на Windows:

set JULIA_NUM_THREADS=4

Powershell на Windows:

$env:JULIA_NUM_THREADS=4

Обратите внимание, что это необходимо сделать перед запуском Julia.

Количество потоков, указанное с помощью -t/--threads, передаётся в процессы-рабочие, которые создаются с помощью аргументов командной строки -p/--procs или --machine-file. Например, julia -p2 -t2 запускает 1 основной процесс с 2-мя рабочими процессами, и все три процесса имеют включенные 2 потока. Для более тонкого управления потоками рабочих процессов используйте addprocs и передайте -t/--threads как exeflags.

Пулы потоков

Когда потоки программы заняты множеством задач, задачи могут испытывать задержки, что может негативно повлиять на отзывчивость и интерактивность программы. Для решения этой проблемы вы можете указать, что задача интерактивна, когда вы Threads.@spawn её:

using Base.Threads
@spawn :interactive f()

Интерактивные задачи должны избегать выполнения операций с высокой задержкой, а если это задачи с длительным временем выполнения, должны часто уступать управление.

Julia может запускаться с одним или несколькими потоками, выделенными для выполнения интерактивных задач:

$ julia --threads 3,1

Аналогично можно использовать переменную окружения JULIA_NUM_THREADS.

export JULIA_NUM_THREADS=3,1

Это запускает Julia с 3 потоками в пуле потоков :default и 1 потоком в пуле потоков :interactive:

julia> using Base.Threads

julia> nthreads()
4

julia> nthreadpools()
2

julia> threadpool()
:default

julia> nthreads(:interactive)
1

Любое или оба числа могут быть заменены словом auto, что заставит Julia выбрать разумное значение по умолчанию.

Обмен данными и синхронизация

Хотя потоки Julia могут обмениваться данными через общую память, очень сложно написать корректный и свободный от гонок многопоточный код. Сетевые каналы Julia (Channel) являются потокобезопасными и могут использоваться для безопасного обмена данными.

Отсутствие гонок данных

Вы полностью отвечаете за обеспечение того, чтобы ваша программа была свободной от гонок данных, и никакие обещания, сделанные здесь, не могут быть гарантированы, если вы не соблюдаете это требование. Результаты, которые вы наблюдаете, могут быть очень неинтуитивными.

Лучший способ добиться этого — получить блокировку вокруг любого доступа к данным, который может быть наблюдаем из нескольких потоков. Например, в большинстве случаев вы должны использовать следующую модель кода:

julia> lock(lk) do
           use(a)
       end

julia> begin
           lock(lk)
           try
               use(a)
           finally
               unlock(lk)
           end
       end

где lk — это блокировка (например, ReentrantLock()), а a — данные.

Кроме того, Julia не безопасна с точки зрения памяти в случае гонок данных. Будьте очень осторожны при чтении любых данных, если другой поток может их записать! Вместо этого всегда используйте приведенную выше модель с блокировкой при изменении данных (например, при присваивании глобальной или переменной замкнутого блока), к которым обращаются другие потоки.

Thread 1:
global b = false
global a = rand()
global b = true

Thread 2:
while !b; end
bad_read1(a) # it is NOT safe to access `a` here!

Thread 3:
while !@isdefined(a); end
bad_read2(a) # it is NOT safe to access `a` here

Макрос @threads

Давайте рассмотрим простой пример с нашими собственными потоками. Создадим массив нулей:

julia> a = zeros(10)
10-element Vector{Float64}:
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0
 0.0

Давайте обработаем этот массив одновременно с помощью 4-х потоков. Каждый поток будет записывать свой идентификатор потока в каждый элемент.

Julia поддерживает параллельные циклы с помощью макроса Threads.@threads. Этот макрос ставится перед циклом for, чтобы указать Julia, что этот участок кода — многопоточный:

julia> Threads.@threads for i = 1:10
           a[i] = Threads.threadid()
       end

Пространство итераций делится между потоками, после чего каждый поток записывает свой идентификатор потока в назначенные ему элементы:

julia> a
10-element Vector{Float64}:
 1.0
 1.0
 1.0
 2.0
 2.0
 2.0
 3.0
 3.0
 4.0
 4.0

Обратите внимание, что Threads.@threads не имеет необязательного параметра редукции, как @distributed.

Использование @threads без гонок данных

Рассмотрим пример с простой суммой:

julia> function sum_single(a)
           s = 0
           for i in a
               s += i
           end
           s
       end
sum_single (generic function with 1 method)

julia> sum_single(1:1_000_000)
500000500000

Простая сумма @threads создаёт гонку данных, так как несколько потоков читают и записывают s одновременно.

julia> function sum_multi_bad(a)
           s = 0
           Threads.@threads for i in a
               s += i
           end
           s
       end
sum_multi_bad (generic function with 1 method)

julia> sum_multi_bad(1:1_000_000)
70140554652

Обратите внимание, что результат не 500000500000, как должен быть, и, скорее всего, будет меняться при каждом выполнении.

Для исправления этого можно использовать буферы, специфичные для задачи, чтобы разбить сумму на участки, свободные от гонок. Здесь sum_single повторно используется, со своим внутренним буфером s, а вектор a разбивается на nthreads() фрагмента для параллельной обработки с помощью nthreads() @spawn-задач.

julia> function sum_multi_good(a)
           chunks = Iterators.partition(a, length(a) ÷ Threads.nthreads())
           tasks = map(chunks) do chunk
               Threads.@spawn sum_single(chunk)
           end
           chunk_sums = fetch.(tasks)
           return sum_single(chunk_sums)
       end
sum_multi_good (generic function with 1 method)

julia> sum_multi_good(1:1_000_000)
500000500000

!!! Примечание Буферы не должны управляться на основе threadid(), то есть buffers = zeros(Threads.nthreads()), так как конкурирующие задачи могут уступать управление, что означает, что несколько конкурирующих задач могут использовать один и тот же буфер в заданном потоке, создавая риск гонок данных. Кроме того, когда доступно более одного потока, задачи могут менять потоки в точках уступки управления, что известно как миграция задач.

Другой вариант — использование атомарных операций над переменными, общими для задач/потоков, что может быть более эффективным в зависимости от характеристик операций.

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

Julia поддерживает доступ и изменение значений атомарно, то есть безопасным для потоков способом, чтобы избежать гонок. Значение (которое должно быть примитивного типа) может быть обернуто как Threads.Atomic, чтобы указать, что к нему нужно обращаться таким образом. Вот пример:

julia> i = Threads.Atomic{Int}(0);

julia> ids = zeros(4);

julia> old_is = zeros(4);

julia> Threads.@threads for id in 1:4
           old_is[id] = Threads.atomic_add!(i, id)
           ids[id] = id
       end

julia> old_is
4-element Vector{Float64}:
 0.0
 1.0
 7.0
 3.0

julia> i[]
 10

julia> ids
4-element Vector{Float64}:
 1.0
 2.0
 3.0
 4.0

Если бы мы попытались выполнить сложение без атомарной метки, мы могли бы получить неправильный ответ из-за гонки. Вот пример того, что бы произошло, если бы мы не избежали гонки:

julia> using Base.Threads

julia> Threads.nthreads()
4

julia> acc = Ref(0)
Base.RefValue{Int64}(0)

julia> @threads for i in 1:1000
          acc[] += 1
       end

julia> acc[]
926

julia> acc = Atomic{Int64}(0)
Atomic{Int64}(0)

julia> @threads for i in 1:1000
          atomic_add!(acc, 1)
       end

julia> acc[]
1000

Атомарные значения полей

Мы также можем использовать атомарные значения на более мелком уровне, используя макросы @atomic, @atomicswap и @atomicreplace.

Подробные сведения о модели памяти и других деталях дизайна описаны в Манифесте атомарных значений Julia, который позже будет опубликован в официальном виде.

Любое поле в объявлении структуры может быть помечено @atomic, а затем любая запись должна быть отмечена @atomic и должна использовать одно из определённых атомарных упорядочений (:monotonic, :acquire, :release, :acquire_release, или :sequentially_consistent). Любое чтение атомарного поля также может быть снабжено атомарным ограничением упорядочения, или будет выполняться с монотонным (релаксированным) упорядочением, если оно не указано.

Атомарные значения полей требуют по крайней мере Julia 1.7.

Побочные эффекты и изменяемые аргументы функций

При использовании многопоточности необходимо быть внимательными, когда используются функции, которые не являются чистыми, так как мы можем получить неверный результат. Например, функции, имена которых заканчиваются на !, по соглашению изменяют свои аргументы и, следовательно, не являются чистыми.

@threadcall

Внешние библиотеки, такие как те, к которым обращаются через ccall, создают проблему для механизма ввода-вывода на основе задач Julia. Если библиотека C выполняет блокирующую операцию, это предотвращает планировщик Julia от выполнения других задач до тех пор, пока вызов не вернётся. (Исключения составляют вызовы в пользовательский код C, который вызывает обратно в Julia, что может привести к уступке управления, или код C, который вызывает jl_yield(), эквивалент yield в C.)

Макрос @threadcall предоставляет способ избежать приостановки выполнения в такой ситуации. Он планирует выполнение функции C в отдельном потоке. Для этого используется пул потоков с размером по умолчанию 4. Размер пула потоков контролируется переменной среды UV_THREADPOOL_SIZE. Во время ожидания свободного потока и во время выполнения функции, когда поток становится доступным, запрашиваемая задача (в главном цикле событий Julia) уступает другим задачам. Обратите внимание, что @threadcall не возвращается до тех пор, пока выполнение не завершится. С точки зрения пользователя, это, следовательно, блокирующий вызов, подобный другим API Julia.

Очень важно, чтобы вызываемая функция не вызывала Julia, так как это приведёт к ошибке сегментирования.

@threadcall может быть удалена/изменена в будущих версиях Julia.

Предупреждения

В настоящее время большинство операций в среде выполнения Julia и стандартных библиотеках могут использоваться в потокобезопасном режиме, если код пользователя не содержит гонок данных. Однако в некоторых областях продолжается работа по стабилизации поддержки потоков. Многопоточная программирование имеет много внутренних трудностей, и если программа, использующая потоки, демонстрирует необычное или нежелательное поведение (например, аварии или загадочные результаты), то обычно следует в первую очередь подозревать взаимодействие потоков.

Есть несколько конкретных ограничений и предупреждений, о которых следует знать при использовании потоков в Julia:

  • Типы коллекций Base требуют ручного блокирования, если они используются одновременно несколькими потоками, где по крайней мере один поток изменяет коллекцию (общие примеры включают push! для массивов или вставку элементов в Dict).
  • График, используемый @spawn , является недетерминированным и на нём нельзя полагаться.
  • Задачами с интенсивными вычислениями, не требующими выделения памяти, могут препятствовать выполнению сбора мусора в других потоках, которые выделяют память. В таких случаях может потребоваться вставить вызов GC.safepoint() , чтобы позволить GC работать. Это ограничение будет устранено в будущем.
  • Избегайте выполнения операций верхнего уровня, например include, или eval определения типов, методов и модулей параллельно.
  • Следует учитывать, что финализаторы, зарегистрированные библиотекой, могут нарушиться, если включены потоки. Для широкого и уверенного принятия потоков в экосистеме может потребоваться определённая переходная работа. Смотрите следующую секцию для получения дополнительных подробностей.

Миграция задач

После запуска задачи в определённом потоке она может переместиться в другой поток, если задача уступит.

Такие задачи могут быть запущены с помощью @spawn или @threads, хотя параметр планирования :static для @threads замораживает идентификатор потока.

Это означает, что в большинстве случаев threadid() не следует рассматривать как постоянную величину в рамках задачи и, следовательно, её нельзя использовать для индексирования вектора буферов или объектов с состоянием.

Миграция задач была введена в Julia 1.7. До этого задачи всегда оставались в том же потоке, в котором они были запущены.

Безопасное использование финализаторов

Поскольку финализаторы могут прерывать любой код, они должны быть очень осторожны в том, как они взаимодействуют с любым глобальным состоянием. К сожалению, основная причина использования финализаторов заключается в обновлении глобального состояния (чистая функция обычно довольно бесполезна в качестве финализатора). Это приводит нас к некоторому парадоксу.

  1. В однопоточном режиме код может вызвать внутреннюю функцию C jl_gc_enable_finalizers , чтобы предотвратить планирование финализаторов внутри критической области. Внутренне это используется в некоторых функциях (например, в наших C-блокировках) для предотвращения рекурсии при выполнении определённых операций (инкрементальная загрузка пакетов, генерация кода и т. д.). Комбинация блокировки и этого флага может использоваться для обеспечения безопасности финализаторов.

  2. Вторая стратегия, используемая Base в нескольких местах, заключается в явном откладывании финализатора до тех пор, пока он не сможет получить блокировку без рекурсии. Следующий пример демонстрирует, как эту стратегию можно применить к Distributed.finalize_ref:

    function finalize_ref(r::AbstractRemoteRef)
        if r.where > 0 # Check if the finalizer is already run
            if islocked(client_refs) || !trylock(client_refs)
                # delay finalizer for later if we aren't free to acquire the lock
                finalizer(finalize_ref, r)
                return nothing
            end
            try # `lock` should always be followed by `try`
                if r.where > 0 # Must check again here
                    # Do actual cleanup here
                    r.where = 0
                end
            finally
                unlock(client_refs)
            end
        end
        nothing
    end
  3. Связанная третья стратегия заключается в использовании очереди без ожидания. В настоящее время в Base не реализована бесблокировочная очередь, но Base.IntrusiveLinkedListSynchronized{T} подходит. Это часто может быть хорошей стратегией для использования в коде с циклами событий. Например, эта стратегия используется Gtk.jl для управления подсчётом ссылок жизненного цикла. В этом подходе мы не выполняем никаких явных операций внутри finalizer, а вместо этого добавляем его в очередь для выполнения в более безопасное время. На самом деле, планировщик задач Julia уже использует это, поэтому определение финализатора как x -> @spawn do_cleanup(x) является одним примером такого подхода. Однако обратите внимание, что это не контролирует, в каком потоке do_cleanup выполняется, поэтому do_cleanup всё ещё должен получить блокировку. Это не обязательно должно выполняться, если вы реализуете свою собственную очередь, так как вы можете явно получить очередь только из своего потока.

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

Spec-Zone.ru

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