Control.Concurrent
| Copyright | (c) The University of Glasgow 2001 |
|---|---|
| License | BSD-style (see the file libraries/base/LICENSE) |
| Maintainer | libraries@haskell.org |
| Stability | experimental |
| Portability | non-portable (concurrency) |
| Safe Haskell | Trustworthy |
| Language | Haskell2010 |
Описание
Общий интерфейс для набора полезных абстракций конкурентности.
Конкурентный Haskell
Расширение конкурентности для Haskell описано в статье Concurrent Haskell http://www.haskell.org/ghc/docs/papers/concurrent-haskell.ps.gz.
Конкурентность является "легковесной", что означает, что накладные расходы на создание потоков и переключение контекста чрезвычайно низки. Планирование потоков Haskell выполняется внутри системы времени выполнения Haskell и не использует какие-либо пакеты потоков, предоставляемые операционной системой.
Однако, если вы хотите взаимодействовать с внешней библиотекой, которая ожидает, что ваша программа будет использовать пакет потоков, предоставляемый операционной системой, вы можете сделать это, используя forkOS вместо forkIO.
Потоки Haskell могут взаимодействовать через MVars, своего рода синхронизированную изменяемую переменную (см. Control.Concurrent.MVar). Несколько распространенных абстракций конкурентности могут быть построены из MVars, и они предоставляются библиотекой Control.Concurrent. В GHC потоки также могут взаимодействовать через исключения.
Основные операции конкурентности
ThreadId является абстрактным типом, представляющим дескриптор потока. ThreadId является экземпляром Eq, Ord и Show, где экземпляр Ord реализует произвольный полный порядок над ThreadIds. Экземпляр Show позволяет преобразовать ThreadId с произвольным значением в строковую форму; отображение значения ThreadId иногда полезно при отладке или диагностике поведения конкурентной программы.
Примечание: в GHC, если у вас есть ThreadId, вы по существу имеете указатель на сам поток. Это означает, что сам поток не может быть собран сборщиком мусора, пока вы не отбросите ThreadId. Надеемся, что этот недостаток будет исправлен позже.
Экземпляры
| Eq ThreadId | Since: base-4.2.0.0 |
| Ord ThreadId | Since: base-4.2.0.0 |
Определено в GHC.Conc.Sync | |
| Show ThreadId | Since: base-4.2.0.0 |
myThreadId :: IO ThreadId Source
Возвращает ThreadId вызывающего потока (только GHC).
forkIO :: IO () -> IO ThreadId Source
Создает новый поток для выполнения вычисления IO, переданного в качестве первого аргумента, и возвращает ThreadId вновь созданного потока.
Новый поток будет легковесным, несвязанным потоком. Не гарантируется, что внешние вызовы, сделанные этим потоком, будут выполнены каким-либо конкретным потоком ОС; если вам нужно, чтобы внешние вызовы выполнялись конкретным потоком ОС, используйте forkOS вместо этого.
Новый поток наследует замаскированное состояние родителя (см. mask).
Созданный поток имеет обработчик исключений, который игнорирует исключения BlockedIndefinitelyOnMVar, BlockedIndefinitelyOnSTM, и ThreadKilled, а все остальные исключения передает обработчику необработанных исключений.
forkFinally :: IO a -> (Either SomeException a -> IO ()) -> IO ThreadId Source
Создает поток и вызывает предоставленную функцию, когда поток завершается, с исключением или возвращаемым значением. Функция вызывается с замаскированными асинхронными исключениями.
forkFinally action and_then =
mask $ \restore ->
forkIO $ try (restore action) >>= and_then
Эта функция полезна для информирования родительского процесса о завершении дочернего процесса, например.
Since: base-4.6.0.0
forkIOWithUnmask :: ((forall a. IO a -> IO a) -> IO ()) -> IO ThreadId Source
Подобно forkIO, но дочернему потоку передается функция, которая может использоваться для размаскирования асинхронных исключений. Эта функция обычно используется следующим образом
... mask_ $ forkIOWithUnmask $ \unmask ->
catch (unmask ...) handler
чтобы обработчик исключений в дочернем потоке был установлен с замаскированными асинхронными исключениями, в то время как основная часть работы дочернего потока выполняется в размаскированном состоянии.
Обратите внимание, что функция размаскирования, переданная дочернему потоку, должна использоваться только в этом потоке; поведение не определено, если она вызывается в другом потоке.
Since: base-4.4.0.0
killThread :: ThreadId -> IO () Source
killThread вызывает исключение ThreadKilled в заданном потоке (только GHC).
killThread tid = throwTo tid ThreadKilled
throwTo :: Exception e => ThreadId -> e -> IO () Source
throwTo вызывает произвольное исключение в целевом потоке (только GHC).
Доставка исключения синхронизируется между исходным и целевым потоком: throwTo не возвращается, пока исключение не будет вызвано в целевом потоке. Таким образом, вызывающий поток может быть уверен, что целевой поток получил исключение. Доставка исключений также атомарна по отношению к другим исключениям. Атомарность полезна при работе с гонками: например, если два потока могут убить друг друга, гарантируется, что только один из потоков сможет убить другой.
Любая работа, которую выполнял целевой поток, когда было вызвано исключение, не теряется: вычисление приостанавливается до тех пор, пока это не потребуется другому потоку.
Если целевой поток в настоящее время выполняет вызов внешней функции, то исключение не будет вызвано (и, следовательно, throwTo не вернётся) до завершения вызова. Это верно независимо от того, находится ли вызов внутри mask или нет. Однако в GHC вызов внешней функции может быть помечен как interruptible, в этом случае throwTo заставит RTS попытаться вызвать возврат; см. документацию GHC для получения дополнительных сведений.
Важно отметить: поведение throwTo отличается от описанного в статье «Асинхронные исключения в Haskell» (http://research.microsoft.com/~simonpj/Papers/asynch-exns.htm). В статье throwTo является асинхронной; но реализация библиотеки использует более синхронный дизайн, в котором throwTo не возвращается, пока исключение не будет получено целевым потоком. Компромисс обсуждается в разделе 9 статьи. Как и любая блокирующая операция, throwTo поэтому прерывиста (см. раздел 5.3 статьи). Однако в отличие от других прерывистых операций, throwTo всегда прерывиста, даже если она фактически не блокируется.
Нет гарантии, что исключение будет доставлено немедленно, хотя среда выполнения постарается гарантировать, что произвольные задержки не появятся. В GHC исключение может быть вызвано только тогда, когда поток достигает безопасной точки, где безопасной точкой является точка, где происходит выделение памяти. Некоторые циклы не выполняют выделения памяти внутри цикла и поэтому не могут быть прерваны throwTo.
Если целью throwTo является вызывающий поток, то поведение аналогично throwIO, за исключением того, что исключение выбрасывается как асинхронное исключение. Это означает, что если есть охватывающее чистое вычисление, что будет иметь место, если текущая операция ввода-вывода находится внутри unsafePerformIO или unsafeInterleaveIO, то это вычисление не заменяется исключением, а приостанавливается так, как если бы оно получило асинхронное исключение.
Обратите внимание, что если throwTo вызывается с текущим потоком в качестве цели, исключение будет вызвано, даже если поток в данный момент находится внутри mask или uninterruptibleMask.
Потоки с аффинностью
forkOn :: Int -> IO () -> IO ThreadId Source
Подобно forkIO, но позволяет указать, на какой возможности должен выполняться поток. В отличие от потока forkIO, поток, созданный с помощью forkOn, будет оставаться на той же возможности на протяжении всего своего жизненного цикла (потоки forkIO могут перемещаться между возможностями в соответствии с политикой планирования). forkOn полезно для переопределения политики планирования, когда заранее известно, как лучше распределить потоки.
Аргумент Int задаёт номер возможности (см. getNumCapabilities). Обычно возможности соответствуют физическим процессорам, но точное поведение зависит от реализации. Значение, переданное forkOn, интерпретируется по модулю общего числа возможностей, как возвращается getNumCapabilities.
Примечание GHC: количество возможностей задаётся параметром +RTS -N при запуске программы. Возможности можно закрепить за фактическими процессорными ядрами с помощью +RTS -qa, если это поддерживает основная операционная система, хотя на практике это обычно не требуется (и может фактически снизить производительность в некоторых случаях — рекомендуется проводить эксперименты).
Since: base-4.4.0.0
forkOnWithUnmask :: Int -> ((forall a. IO a -> IO a) -> IO ()) -> IO ThreadId Source
Подобно forkIOWithUnmask, но дочерний поток закрепляется за заданным процессором, как и с forkOn.
Since: base-4.4.0.0
getNumCapabilities :: IO Int Source
Возвращает число потоков Haskell, которые могут выполняться одновременно (на разных физических процессорах) в любой момент времени. Чтобы изменить это значение, используйте setNumCapabilities.
Since: base-4.4.0.0
setNumCapabilities :: Int -> IO () Source
Устанавливает число потоков Haskell, которые могут выполняться одновременно (на разных физических процессорах) в любой момент времени. Число, переданное forkOn, интерпретируется по модулю этого значения. Начальное значение задается флагом времени выполнения +RTS -N.
Это также число потоков, которые будут участвовать в параллельном сборке мусора. Сильно рекомендуется, чтобы число возможностей не было больше числа физических процессорных ядер, и часто может быть полезно оставить одно или несколько ядер свободными для избежания конфликтов с другими процессами в машине.
Since: base-4.5.0.0
threadCapability :: ThreadId -> IO (Int, Bool) Source
Возвращает номер возможности, на которой в данный момент выполняется поток, и логическое значение, указывающее, заблокирован ли поток к этой возможности. Поток заблокирован к возможности, если он был создан с помощью forkOn.
Since: base-4.4.0.0
Планирование
Планирование может быть либо прерывистым, либо кооперативным, в зависимости от реализации Concurrent Haskell (см. ниже информацию, относящуюся к конкретным компиляторам). В кооперативной системе переключение контекста происходит только при использовании одной из примитивов, определенных в этом модуле. Это означает, что программы, такие как:
main = forkIO (write 'a') >> write 'b'
where write c = putChar c >> write c
будут печатать либо aaaaaaaaaaaaaa... или bbbbbbbbbbbb..., а не какое-то случайное чередование a и b. На практике кооперативная многозадачность достаточна для написания простых графических пользовательских интерфейсов.
Действие yield позволяет (принуждает, в реализации кооперативной многозадачности) переключение контекста на любой другой в данный момент запущенный поток (если таковые имеются), и это иногда полезно при реализации абстракций параллельности.
Блокировка
Разные реализации Haskell имеют разные характеристики относительно операций, которые блокируют все потоки.
Использование GHC без опции -threaded, все внешние вызовы будут блокировать все остальные потоки Haskell в системе, хотя операции ввода-вывода не будут. С опцией -threaded, будут блокировать все остальные потоки только внешние вызовы с атрибутом unsafe.
Ожидание
threadDelay :: Int -> IO () Source
Приостанавливает текущий поток на заданное количество микросекунд (только GHC).
Нет гарантии, что поток будет перепланирован незамедлительно по истечении задержки, но поток никогда не продолжит выполнение раньше указанного.
threadWaitRead :: Fd -> IO () Source
Блокирует текущий поток до тех пор, пока данные не станут доступны для чтения в заданном дескрипторе файла (только GHC).
Это выбросит IOError, если дескриптор файла был закрыт, пока этот поток был заблокирован. Для безопасного закрытия дескриптора файла, который использовался с threadWaitRead, используйте closeFdWith.
threadWaitWrite :: Fd -> IO () Source
Блокирует текущий поток до тех пор, пока данные не смогут быть записаны в заданный дескриптор файла (только GHC).
Это выбросит IOError, если дескриптор файла был закрыт, пока этот поток был заблокирован. Для безопасного закрытия дескриптора файла, который использовался с threadWaitWrite, используйте closeFdWith.
threadWaitReadSTM :: Fd -> IO (STM (), IO ()) Source
Возвращает действие STM, которое может быть использовано для ожидания данных для чтения из дескриптора файла. Второе возвращаемое значение — действие IO, которое можно использовать для отмены интереса к дескриптору файла.
Since: base-4.7.0.0
threadWaitWriteSTM :: Fd -> IO (STM (), IO ()) Source
Возвращает действие STM, которое может быть использовано для ожидания, пока данные не смогут быть записаны в дескриптор файла. Второе возвращаемое значение — действие IO, которое можно использовать для отмены интереса к дескриптору файла.
Since: base-4.7.0.0
Абстракции связи
module Control.Concurrent.MVar
module Control.Concurrent.Chan
module Control.Concurrent.QSem
module Control.Concurrent.QSemN
Связанные потоки
Поддержка нескольких потоков операционной системы и связанных потоков, как описано ниже, в настоящее время доступна в системе выполнения GHC только при использовании опции -threaded при компоновке.
Другие системы Haskell в настоящее время не поддерживают несколько потоков операционной системы.
Связанный поток — это поток Haskell, который связан с потоком операционной системы. Хотя связанный поток по-прежнему планируется системой выполнения Haskell, поток операционной системы обрабатывает все внешние вызовы, сделанные связанным потоком.
Для внешней библиотеки связанный поток будет выглядеть точно так же, как обычный поток операционной системы, созданный с помощью функций ОС, таких как pthread_create или CreateThread.
Связанные потоки могут быть созданы с помощью функции forkOS ниже. Все внешние экспортированные функции выполняются в связанном потоке (связанном с потоком ОС, который вызвал функцию). Кроме того, действие main каждой программы Haskell выполняется в связанном потоке.
Почему это нужно? Потому что если внешняя библиотека вызывается из потока, созданного с помощью forkIO, она не будет иметь доступа к локальному состоянию потока — переменным состояния, которые имеют определенные значения для каждого потока ОС (см. POSIX's pthread_key_create или Win32's TlsAlloc). Следовательно, некоторые библиотеки (OpenGL, например) не будут работать из потока, созданного с помощью forkIO. Они работают нормально в потоках, созданных с помощью forkOS или при вызове из main или из связанного foreign export.
С точки зрения производительности, связанные потоки (связанные) намного дороже, чем несвязанные потоки, поскольку связанный поток привязан к определенному потоку ОС, а несвязанный поток может выполняться любым потоком ОС. Переключение контекста между связанным потоком и несвязанным потоком многократно дороже, чем между двумя связанными потоками.
Обратите особое внимание, что основной поток программы (поток, выполняющий Main.main) всегда является связанным потоком, поэтому для хорошей производительности параллельности необходимо обеспечить, чтобы основной поток не выполнял многократные связи с другими потоками в системе. Как правило, это подразумевает создание дочерних потоков для выполнения задач с использованием forkIO, и ожидание результатов в главном потоке.
rtsSupportsBoundThreads :: Bool Source
True если связанные потоки поддерживаются. Если rtsSupportsBoundThreads равно False, isCurrentThreadBound всегда вернет False и forkOS и runInBoundThread завершатся с ошибкой.
forkOS :: IO () -> IO ThreadId Source
Как forkIO, это запускает новый поток для выполнения вычисления IO, переданного в качестве первого аргумента, и возвращает ThreadId нового созданного потока.
Однако forkOS создает связанный поток, что необходимо, если вам нужно вызвать внешние (не-Haskell) библиотеки, которые используют состояние, локальное для потока, например, OpenGL (см. Control.Concurrent).
Использование forkOS вместо forkIO не влияет на поведение планирования системы выполнения Haskell. Распространенное заблуждение, что нужно использовать forkOS вместо forkIO для предотвращения блокировки всех потоков Haskell при выполнении внешнего вызова; это не так. Для возможности выполнения внешних вызовов без блокировки всех потоков Haskell (с GHC), достаточно использовать опцию -threaded при компоновке программы и убедиться, что внешний импорт не помечен как unsafe.
forkOSWithUnmask :: ((forall a. IO a -> IO a) -> IO ()) -> IO ThreadId Source
Как forkIOWithUnmask, но дочерний поток является связанным потоком, как и в случае forkOS.
isCurrentThreadBound :: IO Bool Source
Возвращает True если вызывающий поток связан, то есть если безопасно использовать внешние библиотеки, которые полагаются на состояние, локальное для потока, из вызывающего потока.
runInBoundThread :: IO a -> IO a Source
Запустите вычисление IO, переданное в качестве первого аргумента. Если вызывающая нить не является привязанной, временно создается привязанная нить. runInBoundThread не завершается, пока не завершится вычисление IO.
Вы можете обернуть серию вызовов внешних функций, которые зависят от состояния, локального для потока, с помощью runInBoundThread, чтобы использовать их, не зная, является ли текущая нить привязанной.
runInUnboundThread :: IO a -> IO a Исходный код
Запустите вычисление IO, переданное в качестве первого аргумента. Если вызывающая нить привязанная, временно создается непривязанная нить с помощью forkIO. runInBoundThread не завершается, пока не завершится вычисление IO.
Используйте эту функцию только в редких случаях, когда вы действительно наблюдали снижение производительности из-за использования привязанных нитей. Программа, которой не нужно, чтобы ее основная нить была привязана и которая активно использует конкурентность (например, веб-сервер), может обернуть свое действие main в runInUnboundThread.
Обратите внимание, что исключения, которые выбрасываются в текущую нить, в свою очередь, выбрасываются в нить, которая выполняет данное вычисление. Это гарантирует, что всегда есть возможность убить разветвленную нить.
Слабые ссылки на идентификаторы потоков
mkWeakThreadId :: ThreadId -> IO (Weak ThreadId) Исходный код
Создайте слабую ссылку на идентификатор потока. Это может быть важно, если вы хотите сохранить ссылку на идентификатор потока, одновременно позволяя нити получать исключения семейства BlockedIndefinitely (например, BlockedIndefinitelyOnMVar). Сохранение обычной ссылки ThreadId предотвратит доставку исключений BlockedIndefinitely, потому что ссылка может быть использована как целевой объект для throwTo в любое время, что разблокирует нить.
Сохранение слабой ссылки, с другой стороны, не помешает нити получать исключения BlockedIndefinitely. Исключение всё ещё может быть выброшено в нить со слабой ссылкой, но вызывающий код должен сначала использовать deRefWeak, чтобы определить, существует ли нить.
С версии: base-4.6.0.0
Реализация конкурентности в GHC
Этот раздел описывает особенности, специфичные для реализации Concurrent Haskell в GHC.
Потоки Haskell и потоки операционной системы
В GHC потоки, созданные с помощью forkIO, являются лёгкими потоками и управляются полностью средой выполнения GHC. Как правило, потоки Haskell на порядок или два эффективнее (по времени и ресурсам) потоков операционной системы.
Недостатком использования лёгких потоков является то, что только один из них может выполняться в данный момент, поэтому, если один поток заблокируется в вызове внешней функции, например, другие потоки не смогут продолжить работу. Среда выполнения GHC обходит эту проблему, используя полные потоки ОС при необходимости. Когда программа собирается с опцией -threaded (для подключения к многопотоковой версии среды выполнения), нить, производящая вызов safe внешней функции, не будет блокировать другие потоки в системе; другой поток ОС возьмёт на себя выполнение потоков Haskell до момента возвращения исходного вызова. Среда выполнения поддерживает пул таких рабочих потоков, чтобы несколько потоков Haskell могли одновременно участвовать во внешних вызовах.
Библиотека System.IO управляет мультиплексированием своим способом. В системах Windows она использует вызовы safe внешних функций для обеспечения того, чтобы потоки, выполняющие операции ввода-вывода, не блокировали всю среду выполнения, в то время как в системах Unix все текущие заблокированные запросы ввода-вывода обрабатываются единственным потоком (менеджером потоков ввода-вывода) с помощью механизма, такого как epoll или kqueue, в зависимости от того, что предоставляется хостовой операционной системой.
Среда выполнения запустит поток Haskell с использованием любого доступного рабочего потока ОС. Если вам нужен контроль над тем, какой именно поток ОС используется для запуска данного потока Haskell, возможно, потому, что вам нужно вызвать внешнюю библиотеку, которая использует состояние, локальное для потока ОС, то вам нужны привязанные потоки (см. Control.Concurrent).
Если вы не используете опцию -threaded, среда выполнения не использует несколько потоков ОС. Вызовы внешних функций будут блокировать все другие работающие потоки Haskell, пока вызов не вернётся. Библиотека System.IO всё ещё реализует мультиплексирование, так что может быть несколько потоков, выполняющих операции ввода-вывода, и это обрабатывается внутри среды выполнения с помощью select.
Завершение программы
В автономной программе GHC для завершения процесса требуется завершение только главной нити. Таким образом, все другие разветвленные потоки просто завершатся одновременно с главной нитью (терминология для такого поведения — «демоновые потоки»).
Если вы хотите, чтобы программа ждала завершения дочерних потоков перед выходом, вам нужно запрограммировать это самостоятельно. Простой механизм — каждый дочерний поток записывает в MVar при завершении, а главная нить ждёт всех MVar перед выходом:
myForkIO :: IO () -> IO (MVar ())
myForkIO io = do
mvar <- newEmptyMVar
forkFinally io (\_ -> putMVar mvar ())
return mvar
Обратите внимание, что мы используем forkFinally для обеспечения записи в MVar, даже если поток умирает или убивается по какой-либо причине.
Более надёжный метод — сохранять глобальный список всех дочерних потоков, за которыми следует ожидание по окончании программы:
children :: MVar [MVar ()]
children = unsafePerformIO (newMVar [])
waitForChildren :: IO ()
waitForChildren = do
cs <- takeMVar children
case cs of
[] -> return ()
m:ms -> do
putMVar children ms
takeMVar m
waitForChildren
forkChild :: IO () -> IO ThreadId
forkChild io = do
mvar <- newEmptyMVar
childs <- takeMVar children
putMVar children (mvar:childs)
forkFinally io (\_ -> putMVar mvar ())
main =
later waitForChildren $
...
Принцип главной нити также применим к вызовам Haskell извне, используя foreign export. Когда вызываемая функция foreign export вызывается, она запускает новую главную нить и возвращает значение, когда эта главная нить завершается. Если вызов приводит к созданию новых потоков, они могут оставаться в системе после возвращения из foreign export функции.
Прерывание
GHC реализует прерывание задач: выполнение потоков переплетается случайным образом. Более конкретно, поток может быть прерван, когда он выделяет память, что, к сожалению, означает, что жёсткие циклы, не выполняющие выделение, склонны блокировать другие потоки (это происходит только в патологических сценариях кода для бенчмарков).
Таймер перепланирования работает по умолчанию с интервалом 20 мс, но это можно изменить, используя опцию -i<n> RTS. После «титика» перепланирования работающий поток прерывается как можно скорее.
Ещё одна заметка: пример aaaa bbbb может работать не очень хорошо в GHC (см. Перепланирование выше) из-за блокировки на Handle. Только один поток может удерживать блокировку на Handle в любой момент, поэтому если происходит перепланирование, пока поток удерживает блокировку, другой поток не сможет запуститься. Результат заключается в том, что переход от aaaa к bbbbb происходит редко. Это можно улучшить, уменьшив период перепланирования. У нас также есть патч, который вызывает перепланирование всякий раз, когда поток, ждущий блокировки, просыпается, но мы не нашли ему применения, кроме этого примера :-)
Тупик
GHC пытается обнаружить тупиковые ситуации у потоков с помощью сборщика мусора. Поток, который недоступен (не может быть найден, следуя указателям от живых объектов), должен быть тупиковым, и в этом случае в поток посылается исключение. Исключение является BlockedIndefinitelyOnMVar, BlockedIndefinitelyOnSTM, NonTermination или Deadlock, в зависимости от способа застревания потока.
Обратите внимание, что эта функция предназначена для отладки и не должна использоваться для корректной работы вашей программы. Нет гарантии, что сборщик мусора будет достаточно точным, чтобы обнаружить ваш тупик, и нет гарантии, что сборщик мусора будет выполняться достаточно быстро. В основном, те же оговорки, что и для финализаторов, применяются к обнаружению тупика.
Существует тонкое взаимодействие между обнаружением тупика и финализаторами (созданными с помощью newForeignPtr или функций в System.Mem.Weak): если поток заблокирован, ожидая выполнения финализатора, то поток будет считаться заблокированным и получит исключение. Так что предпочтительно не делать этого, но если у вас нет альтернативы, то можно предотвратить, чтобы поток не считался тупиковым, создав слабую ссылку на него. Не забудьте освободить слабую ссылку позже с помощью freeStablePtr.
© The University of Glasgow and others
Licensed under a BSD-style license (see top of the page).
https://downloads.haskell.org/~ghc/8.10.2/docs/html/libraries/base-4.14.1.0/Control-Concurrent.html