HORIZON HASKELLDocslts/ghc-9.10.x248f8f02026-10-05Search names, modules, packages, or :: a typeCtrl K

GHC 9.10.3 · lts/ghc-9.10.x · 248f8f0 · 2026-10-05

Modulestreamly-0.10.1Haskell2010

Streamly.Internal.Data.Stream.Parallel

Deprecated. Please use Streamly.Internal.Data.Stream.Concurrent instead.

To run examples in this module:

Example3 expressions
import qualified Streamly.Prelude as Streamimport Control.Concurrent (threadDelay):{ delay n = do     threadDelay (n * 1000000)   -- sleep for n seconds     putStrLn (show n ++ " sec") -- print "n sec"     return n                    -- IO Int:}
  • 2 types
  • 9 values
  • Packagestreamly-0.10.1
  • Exports11
  • LanguageHaskell2010
  • LicenceBSD-3-Clause
  • SourceParallel.hs

Parallel Stream Type

3 declarations
newtypenewtype ParallelT (m :: Type -> Type) a
#

For ParallelT streams:

(<>) = Streamly.Prelude.parallel
(>>=) = flip . Streamly.Prelude.concatMapWith Streamly.Prelude.parallel

See Streamly.Prelude.AsyncT, ParallelT is similar except that all iterations are strictly concurrent while in AsyncT it depends on the consumer demand and available threads. See parallel for more details.

Since: 0.1.0 (Streamly)

Since: 0.7.0 (maxBuffer applies to ParallelT streams)

Constructors

Instances10IsStream, MonadReader, MonadState, Monad, Functor, Applicative, …
typetype Parallel = ParallelT IO
#

A parallely composing IO stream of elements of type a. See ParallelT documentation for more details.

Since: 0.2.0 (Streamly)

valueconsM :: MonadAsync m => m a -> ParallelT m a -> ParallelT m a
#

XXX we can implement it more efficienty by directly implementing instead of combining streams using parallel.

Merge Concurrently

3 declarations

Evaluate Concurrently

2 declarations

Tap Concurrently

2 declarations
valuetapAsyncK
  1. :: MonadAsync m
  2. => StreamK m a -> m b
  3. -> StreamK m a
  4. -> StreamK m a
#

Redirect a copy of the stream to a supplied fold and run it concurrently in an independent thread. The fold may buffer some elements. The buffer size is determined by the prevailing Streamly.Prelude.maxBuffer setting.

              StreamK m a -> m b
                      |
-----stream m a ---------------stream m a-----

> S.drain $ S.tapAsync (S.mapM_ print) (S.enumerateFromTo 1 2)
1
2

Exceptions from the concurrently running fold are propagated to the current computation. Note that, because of buffering in the fold, exceptions may be delayed and may not correspond to the current element being processed in the parent stream, but we guarantee that before the parent stream stops the tap finishes and all exceptions from it are drained.

Compare with tap.

Pre-release

Callbacks

1 declaration
valuenewCallbackStream :: MonadAsync m => m (a -> m (), StreamK m a)
#

Generates a callback and a stream pair. The callback returned is used to queue values to the stream. The stream is infinite, there is no way for the callback to indicate that it is done now.

Pre-release