Spec-Zone.ru › Julia 1.10

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

Посетите эту статью блога для ознакомления с функциями многопоточности 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.

Несколько потоков сборщика мусора

Сборщик мусора (GC) может использовать несколько потоков. Количество используемых потоков равно либо половине количества вычислительных потоков-воркеров, либо настраивается с помощью аргумента командной строки --gcthreads или переменной среды JULIA_NUM_GC_THREADS.

Аргумент командной строки --gcthreads требует как минимум Julia 1.10.

Потоковые пулы

Когда потоки программы заняты множеством задач, задачи могут испытывать задержки, что может негативно сказаться на отзывчивости и интерактивности программы. Чтобы решить эту проблему, вы можете указать, что задача является интерактивной, когда вы 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> nthreadpools()
2

julia> threadpool() # the main thread is in the interactive thread pool
:interactive

julia> nthreads(:default)
3

julia> nthreads(:interactive)
1

julia> nthreads()
3

Версия nthreads без аргументов возвращает количество потоков в пуле по умолчанию.

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

Любое из этих чисел или оба могут быть заменены словом 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(), аналог C yield.)

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

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

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

Предостережения

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

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

  • Типы коллекций Base требуют ручного блокирования, если они используются одновременно несколькими потоками, где хотя бы один поток изменяет коллекцию (типичные примеры включают push! массивов или вставка элементов в Dict).
  • Планировщик, используемый @spawn , является недетерминированным и на него нельзя полагаться.
  • Задача с вычислениями, не требующая выделения памяти, может препятствовать выполнению сборки мусора в других потоках, которые выделяют память. В этих случаях может потребоваться вставить явный вызов GC.safepoint() для запуска сборки мусора. Это ограничение будет устранено в будущем.
  • Избегайте выполнения операций верхнего уровня, например 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–2024 Jeff Bezanson, Stefan Karpinski, Viral B. Shah, and other contributors
Licensed under the MIT License.
https://docs.julialang.org/en/v1.10/manual/multi-threading/

Spec-Zone.ru

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