Control.Concurrent
| Авторские права | (c) Университет Глазго 2001 |
|---|---|
| Лицензия | BSD-стиль (см. файл libraries/base/LICENSE) |
| Поддержка | libraries@haskell.org |
| Устойчивость | экспериментальная |
| Переносимость | непереносимо (конкурентность) |
| Безопасный Haskell | Надежный |
| Язык | Haskell2010 |
Содержание
Описание
Общий интерфейс для набора полезных абстракций конкурентности.
Конкурентный Haskell
Расширение конкурентности для Haskell описано в статье Конкурентный Haskell http://www.haskell.org/ghc/docs/papers/concurrent-haskell.ps.gz.
Конкурентность «лёгковесная», что означает, что накладные расходы как при создании потоков, так и при переключении контекстов чрезвычайно низкие. Планирование потоков Haskell выполняется внутри системы выполнения Haskell и не использует пакеты потоков операционной системы.
Однако, если вам нужно взаимодействовать с внешней библиотекой, которая ожидает, что ваша программа будет использовать пакет потоков операционной системы, вы можете сделать это, используя forkOS вместо forkIO.
Потоки Haskell могут обмениваться данными через MVar — вид синхронизированной изменяемой переменной (см. Control.Concurrent.MVar). Несколько общих абстракций конкурентности могут быть построены из MVar и они предоставляются библиотекой Control.Concurrent. В GHC потоки также могут обмениваться данными через исключения.
Основные операции конкурентности
ThreadId — это абстрактный тип, представляющий дескриптор потока. ThreadId является экземпляром Eq, Ord и Show, где экземпляр Ord реализует произвочный общий порядок над ThreadId. Экземпляр Show позволяет преобразовать ThreadId произвольного типа в строку; отображать значение ThreadId бывает полезно при отладке или диагностике поведения конкурирующей программы.
Примечание: в GHC, если у вас есть ThreadId, у вас по сути есть указатель на сам поток. Это означает, что сам поток не может быть удалён из памяти, пока вы не освободите ThreadId. Эта особенность, к сожалению, будет исправлена в будущем.
myThreadId :: IO ThreadId Источник
Возвращает идентификатор ThreadId вызывающего потока (только GHC).
forkIO :: IO () -> IO ThreadId Источник
Создаёт новый поток для выполнения IO вычисления, переданного в качестве первого аргумента, и возвращает идентификатор ThreadId вновь созданного потока.
Новый поток будет лёгковесным, *непривязанным* потоком. Внешние вызовы, выполняемые этим потоком, не гарантируют выполнения с помощью определённого потока ОС; если вам нужны внешние вызовы, выполняемые определённым потоком ОС, используйте forkOS вместо этого.
Новый поток наследует *замаскированное* состояние родительского (см. mask).
У вновь созданного потока есть обработчик исключений, который игнорирует исключения BlockedIndefinitelyOnMVar, BlockedIndefinitelyOnSTM и ThreadKilled, и передаёт все другие исключения обработчику необработанных исключений.
forkFinally :: IO a -> (Either SomeException a -> IO ()) -> IO ThreadId Источник
Создаёт поток и вызывает переданную функцию, когда поток завершается с исключением или возвращаемым значением. Функция вызывается с замаскированными асинхронными исключениями.
forkFinally action and_then =
mask $ \restore ->
forkIO $ try (restore action) >>= and_then
Эта функция полезна для информирования родителя о завершении дочернего потока, например.
С версии: 4.6.0.0
forkIOWithUnmask :: ((forall a. IO a -> IO a) -> IO ()) -> IO ThreadId Источник
Подобно forkIO, но дочернему потоку передаётся функция, которая может использоваться для размаскирования асинхронных исключений. Эта функция обычно используется следующим образом:
... mask_ $ forkIOWithUnmask $ \unmask ->
catch (unmask ...) handler
чтобы обработчик исключений в дочернем потоке был установлен с замаскированными асинхронными исключениями, а основное тело дочернего потока выполнялось в размаскированном состоянии.
Обратите внимание, что функция размаскирования, переданная дочернему потоку, должна использоваться только в этом потоке; поведение неопределено, если она вызывается в другом потоке.
С версии: 4.4.0.0
killThread :: ThreadId -> IO () Источник
killThread генерирует исключение ThreadKilled в заданном потоке (только GHC).
killThread tid = throwTo tid ThreadKilled
throwTo :: Exception e => ThreadId -> e -> IO () Источник
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, за исключением того, что исключение генерируется как асинхронное исключение. Это означает, что если есть окружающее чистое вычисление, что будет иметь место, если текущая операция IO находится внутри unsafePerformIO или unsafeInterleaveIO, это вычисление не заменяется исключением, а приостанавливается так, как если бы оно получило асинхронное исключение.
Обратите внимание, что если throwTo вызывается с текущей нитью в качестве целевой, исключение будет брошено даже если нить в настоящее время находится внутри mask или uninterruptibleMask.
Потоки с привязкой
forkOn :: Int -> IO () -> IO ThreadId Источник
Подобно forkIO, но позволяет указать, на какой ресурс должна выполняться нить. В отличие от нити forkIO, нить, созданная с помощью forkOn, будет оставаться на том же ресурсе в течение всего своего жизненного цикла (нити forkIO могут перемещаться между ресурсами в соответствии с политикой планирования). forkOn полезна для переопределения политики планирования, когда заранее известно, как лучше распределить потоки.
Аргумент Int задаёт номер ресурса (см. getNumCapabilities). Обычно ресурсы соответствуют физическим процессорам, но точное поведение зависит от реализации. Значение, переданное в forkOn, интерпретируется по модулю общего количества ресурсов, возвращаемого функцией getNumCapabilities.
Примечание GHC: количество ресурсов задаётся опцией +RTS -N при запуске программы. Ресурсы могут быть закреплены за реальными ядрами процессора с помощью +RTS -qa, если это поддерживается операционной системой, хотя на практике это обычно не требуется (и может даже ухудшить производительность в некоторых случаях — рекомендуется экспериментирование).
С версии: 4.4.0.0
forkOnWithUnmask :: Int -> ((forall a. IO a -> IO a) -> IO ()) -> IO ThreadId Источник
Подобно forkIOWithUnmask, но дочерняя нить закреплена за указанным процессором, как в случае с forkOn.
С версии: 4.4.0.0
getNumCapabilities :: IO Int Источник
Возвращает количество потоков Haskell, которые могут выполняться одновременно (на отдельных физических процессорах) в любой момент времени. Для изменения этого значения используйте setNumCapabilities.
С версии: 4.4.0.0
setNumCapabilities :: Int -> IO () Источник
Устанавливает количество потоков Haskell, которые могут выполняться одновременно (на отдельных физических процессорах) в любой момент времени. Число, переданное forkOn, интерпретируется по модулю этого значения. Начальное значение задаётся флагом времени выполнения +RTS -N.
Это также количество потоков, которые будут участвовать в параллельной сборке мусора. Сильно рекомендуется, чтобы количество ресурсов не превышало количество физических ядер процессора, и часто полезно оставлять одно или несколько ядер свободными, чтобы избежать конфликтов с другими процессами в системе.
С версии: 4.5.0.0
threadCapability :: ThreadId -> IO (Int, Bool) Источник
Возвращает номер ресурса, на котором в данный момент выполняется поток, и булево значение, указывающее, заблокирован ли поток на этом ресурсе. Поток заблокирован на ресурсе, если он был создан с помощью forkOn.
С версии: 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 () Источник
Приостанавливает текущую нить на заданное количество микросекунд (только GHC).
Нет гарантии, что поток будет возобновлён сразу после истечения срока ожидания, но поток никогда не продолжит выполнение раньше, чем задано.
threadWaitRead :: Fd -> IO () Источник
Блокирует текущий поток до тех пор, пока данные не станут доступными для чтения из заданного дескриптора файла (только GHC).
Это вызовет IOError, если дескриптор файла был закрыт, пока этот поток был заблокирован. Для безопасного закрытия дескриптора файла, который использовался с threadWaitRead, используйте closeFdWith.
threadWaitWrite :: Fd -> IO () Источник
Блокирует текущий поток до тех пор, пока данные не станут доступными для записи в заданный дескриптор файла (только GHC).
Это вызовет IOError, если дескриптор файла был закрыт, пока этот поток был заблокирован. Для безопасного закрытия дескриптора файла, который использовался с threadWaitWrite, используйте closeFdWith.
threadWaitReadSTM :: Fd -> IO (STM (), IO ()) Источник
Возвращает действие STM, которое может быть использовано для ожидания данных для чтения из дескриптора файла. Второе возвращаемое значение — действие IO, которое может быть использовано для отмены интереса к дескриптору файла.
С версии: 4.7.0.0
threadWaitWriteSTM :: Fd -> IO (STM (), IO ()) Источник
Возвращает действие STM, которое может быть использовано для ожидания, пока данные не станут доступны для записи в дескриптор файла. Второе возвращаемое значение — действие IO, которое может быть использовано для отмены интереса к дескриптору файла.
С версии: 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.
С точки зрения производительности, потоки forkOS (также известные как привязанные) намного дороже, чем потоки forkIO (также известные как свободные), поскольку поток forkOS привязан к конкретному потоку ОС, тогда как поток forkIO может быть запущен любым потоком ОС. Переключение контекста между потоком forkOS и потоком forkIO во много раз дороже, чем между двумя потоками forkIO.
Обратите особое внимание на то, что основной поток программы (поток, выполняющий 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 runtime. Распространённое заблуждение состоит в том, что вам нужно использовать forkOS вместо forkIO для предотвращения блокировки всех потоков Haskell при выполнении внешнего вызова; это не так. Чтобы позволить внешним вызовам выполняться без блокировки всех потоков Haskell (с GHC), необходимо только использовать опцию -threaded при линковке вашей программы и убедиться, что внешний импорт не помечен как unsafe.
isCurrentThreadBound :: IO Bool Source
Возвращает True, если вызывающий поток является привязанным, то есть если безопасно использовать внешние библиотеки, которые полагаются на состояние, локальное для потока, из вызывающего потока.
runInBoundThread :: IO a -> IO a Source
Выполняет IO вычисление, переданное в качестве первого аргумента. Если вызывающий поток не является привязанным, временно создаётся привязанный поток. runInBoundThread не завершается до тех пор, пока IO вычисление не завершится.
Вы можете обернуть серию вызовов внешних функций, которые полагаются на состояние, локальное для потока, с помощью runInBoundThread, чтобы использовать их, не зная, является ли текущий поток привязанным.
runInUnboundThread :: IO a -> IO a Source
Выполняет IO вычисление, переданное в качестве первого аргумента. Если вызывающий поток является привязанным, временно создаётся свободный поток с помощью forkIO. runInBoundThread не завершается до тех пор, пока IO вычисление не завершится.
Используйте эту функцию только в редких случаях, когда вы действительно наблюдали потерю производительности из-за использования привязанных потоков. Программа, которой не нужен привязанный основной поток и которая интенсивно использует конкурецию (например, веб-сервер), может обернуть своё main действие в runInUnboundThread.
Обратите внимание, что исключения, которые выбрасываются в текущий поток, выбрасываются, в свою очередь, в поток, выполняющий заданное вычисление. Это гарантирует, что всегда есть способ убить порождённый поток.
Слабые ссылки на ThreadIds
mkWeakThreadId :: ThreadId -> IO (Weak ThreadId) Source
Создаёт слабую ссылку на ThreadId. Это может быть важно, если вы хотите сохранить ссылку на ThreadId и всё ещё хотите разрешить потоку получать исключения семейства BlockedIndefinitely (например, BlockedIndefinitelyOnMVar). Сохранение обычной ThreadId ссылки предотвратит доставку исключений BlockedIndefinitely, потому что ссылка может быть использована в качестве цели throwTo в любое время, что разблокирует поток.
Сохранение Weak ThreadId, с другой стороны, не предотвратит получение потоком исключений BlockedIndefinitely. Вы всё ещё можете выбросить исключение в Weak ThreadId, но вызывающая сторона должна сначала использовать deRefWeak для определения того, существует ли поток по-прежнему.
Since: 4.6.0.0
Реализация конкуретности GHC
В этом разделе описываются функции, специфичные для реализации Concurrent Haskell в GHC.
Потоки Haskell и потоки операционной системы
В GHC потоки, созданные с помощью forkIO, являются лёгкими потоками и управляются исключительно Haskell runtime. Обычно потоки Haskell на порядок или два более эффективны (по времени и по памяти) по сравнению с потоками операционной системы.
Недостатком использования лёгких потоков является то, что только один из них может работать одновременно, поэтому, если один поток заблокируется во внешнем вызове, например, другие потоки не смогут продолжить работу. Haskell runtime обходит это, используя полноценные потоки ОС при необходимости. Когда программа построена с опцией -threaded (для линковки с многопоточной версией runtime), поток, производящий safe внешний вызов, не будет блокировать другие потоки в системе; другой поток ОС возьмет на себя управление потоками Haskell до момента возврата из исходного вызова. Runtime поддерживает пул этих потоков-рабочих, чтобы несколько потоков Haskell могли участвовать во внешних вызовах одновременно.
Библиотека System.IO управляет мультиплексированием своим собственным способом. На системах Windows она использует safe внешние вызовы, чтобы гарантировать, что потоки, выполняющие операции ввода-вывода, не блокируют весь runtime, а на системах Unix все текущие заблокированные запросы ввода-вывода управляются одним потоком (поток менеджера ввода-вывода) с помощью механизма, такого как epoll или kqueue, в зависимости от того, что предоставляет операционная система.
Runtime будет запускать поток Haskell с помощью любого из доступных потоков ОС-рабочих. Если вам нужен контроль над конкретным потоком ОС, который используется для выполнения данного потока Haskell, возможно, потому, что вам нужно вызвать внешнюю библиотеку, которая использует состояние, локальное для потока ОС, тогда вам нужны привязанные потоки (см. Control.Concurrent).
Если вы не используете опцию -threaded, runtime не использует несколько потоков ОС. Внешние вызовы будут блокировать все другие работающие потоки Haskell до возврата из вызова. Библиотека System.IO всё ещё выполняет мультиплексирование, так что может быть несколько потоков, выполняющих операции ввода-вывода, и это обрабатывается runtime внутри с помощью 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): если поток заблокирован в ожидании выполнения финализатора, то поток будет считаться заблокированным и получит исключение. Поэтому, желательно этого не делать, но если у вас нет альтернативы, то можно предотвратить списание потока как заблокированного, создав StablePtr, указывающий на него. Не забудьте освободить StablePtr позже с помощью freeStablePtr.
© The University of Glasgow and others
Licensed under a BSD-style license (see top of the page).
https://downloads.haskell.org/~ghc/7.10.3/docs/html/libraries/base-4.8.2.0/Control-Concurrent.html