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.Data.Stream.Prelude

For upgrading to streamly-0.9.0+ please read the Streamly-0.9.0 upgrade guide. Also, see the Streamly.Data.Stream.MkType module for direct replacement of stream types that have been removed in 0.9.0.

All Stream related combinators including the streamly-core Streamly.Data.Stream module, concurrency, time and lifted exception operations. For more pre-release operations also see Streamly.Internal.Data.Stream.Prelude module.

  • 5 types
  • 1 class
  • 128 values
  • Packagestreamly-0.10.1
  • Exports134
  • LanguageHaskell2010
  • LicenceBSD-3-Clause
  • SourcePrelude.hs

Streamly.Data.Stream

90 declarations

All Streamly.Data.Stream combinators are re-exported via this module. For more pre-release combinators also see Streamly.Internal.Data.Stream module.

valuerepeat :: Monad m => a -> Stream m a
#

Generate an infinite stream by repeating a pure value.

Example1 expression
repeat x = Stream.repeatM (pure x)
datadata Stream (m :: Type -> Type) a
#

A stream consists of a step function that generates the next step given a current state, and the current state.

Instances9Functor, Foldable, IsList, Eq, Ord, Read, …
valueonException :: MonadCatch m => m b -> Stream m a -> Stream m a
#

Run the action m b if the stream evaluation is aborted due to an exception. The exception is not caught, simply rethrown.

Observes exceptions only in the stream generation, and not in stream consumers.

Inhibits stream fusion

valuebracketIO
  1. :: (MonadIO m, MonadCatch m)
  2. => IO b
  3. -> b -> IO c
  4. -> b -> Stream m a
  5. -> Stream m a
#

Run the alloc action IO b with async exceptions disabled but keeping blocking operations interruptible (see mask). Use the output b of the IO action as input to the function b -> Stream m a to generate an output stream.

b is usually a resource under the IO monad, e.g. a file handle, that requires a cleanup after use. The cleanup action b -> IO c, runs whenever (1) the stream ends normally, (2) due to a sync or async exception or, (3) if it gets garbage collected after a partial lazy evaluation. The exception is not caught, it is rethrown.

bracketIO only guarantees that the cleanup action runs, and it runs with async exceptions enabled. The action must ensure that it can successfully cleanup the resource in the face of sync or async exceptions.

When the stream ends normally or on a sync exception, cleanup action runs immediately in the current thread context, whereas in other cases it runs in the GC context, therefore, cleanup may be delayed until the GC gets to run. An example where GC based cleanup happens is when a stream is being folded but the fold terminates without draining the entire stream or if the consumer of the stream encounters an exception.

Observes exceptions only in the stream generation, and not in stream consumers.

See also: bracketUnsafe

Inhibits stream fusion

valuetoList :: Monad m => Stream m a -> m [a]
#

Definitions:

Example2 expressions
toList = Stream.foldr (:) []toList = Stream.fold Fold.toList

Convert a stream into a list in the underlying monad. The list can be consumed lazily in a lazy monad (e.g. Identity). In a strict monad (e.g. IO) the whole list is generated and buffered before it can be consumed.

Warning! working on large lists accumulated as buffers in memory could be very inefficient, consider using Streamly.Data.Array instead.

Note that this could a bit more efficient compared to Stream.fold Fold.toList, and it can fuse with pure list consumers.

valuenil :: Applicative m => Stream m a
#

A stream that terminates without producing any output or side effect.

Example1 expression
Stream.toList Stream.nil[]
valuecons :: Applicative m => a -> Stream m a -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of elements. Use the O(n) complexity StreamK.Streamly.Data.StreamK.cons unless you want to statically fuse just a few elements.

Fuse a pure value at the head of an existing stream::

Example2 expressions
s = 1 `Stream.cons` Stream.fromList [2,3]Stream.toList s[1,2,3]

Definition:

Example1 expression
cons x xs = return x `Stream.consM` xs
valueunfold :: Applicative m => Unfold m a b -> a -> Stream m b
#

Convert an Unfold into a stream by supplying it an input seed.

Example2 expressions
s = Stream.unfold Unfold.replicateM (3, putStrLn "hello")Stream.fold Fold.drain shellohellohello
valueunfoldr :: Monad m => (s -> Maybe (a, s)) -> s -> Stream m a
#

Build a stream by unfolding a pure step function step starting from a seed s. The step function returns the next element in the stream and the next seed value. When it is done it returns Nothing and the stream ends. For example,

Example1 expression
:{let f b =        if b > 2        then Nothing        else Just (b, b + 1)in Stream.toList $ Stream.unfoldr f 0:}[0,1,2]
valueunfoldrM :: Monad m => (s -> m (Maybe (a, s))) -> s -> Stream m a
#

Build a stream by unfolding a monadic step function starting from a seed. The step function returns the next element in the stream and the next seed value. When it is done it returns Nothing and the stream ends. For example,

Example1 expression
:{let f b =        if b > 2        then return Nothing        else return (Just (b, b + 1))in Stream.toList $ Stream.unfoldrM f 0:}[0,1,2]
valuefromPure :: Applicative m => a -> Stream m a
#

Create a singleton stream from a pure value.

Example3 expressions
fromPure a = a `Stream.cons` Stream.nilfromPure = purefromPure = Stream.fromEffect . pure
valuefromEffect :: Applicative m => m a -> Stream m a
#

Create a singleton stream from a monadic action.

Example2 expressions
fromEffect m = m `Stream.consM` Stream.nilfromEffect = Stream.sequence . Stream.fromPure
Example1 expression
Stream.fold Fold.drain $ Stream.fromEffect (putStrLn "hello")hello
valuerepeatM :: Monad m => m a -> Stream m a
#
Example1 expression
repeatM = Stream.sequence . Stream.repeat

Generate a stream by repeatedly executing a monadic action forever.

Example1 expression
:{repeatAction =       Stream.repeatM (threadDelay 1000000 >> print 1)     & Stream.take 10     & Stream.fold Fold.drain:}
valuereplicate :: Monad m => Int -> a -> Stream m a
#
Example2 expressions
replicate n = Stream.take n . Stream.repeatreplicate n x = Stream.replicateM n (pure x)

Generate a stream of length n by repeating a value n times.

valuereplicateM :: Monad m => Int -> m a -> Stream m a
#
Example1 expression
replicateM n = Stream.sequence . Stream.replicate n

Generate a stream by performing a monadic action n times.

classclass Enum a => Enumerable a where
#

Types that can be enumerated as a stream. The operations in this type class are equivalent to those in the Enum type class, except that these generate a stream instead of a list. Use the functions in Streamly.Internal.Data.Stream.Enumeration module to define new instances.

Methods

  • enumerateFrom :: Monad m => a -> Stream m a

    enumerateFrom from generates a stream starting with the element from, enumerating up to maxBound when the type is Bounded or generating an infinite stream when the type is not Bounded.

    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFrom (0 :: Int)[0,1,2,3]

    For Fractional types, enumeration is numerically stable. However, no overflow or underflow checks are performed.

    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFrom 1.1[1.1,2.1,3.1,4.1]
  • enumerateFromTo :: Monad m => a -> a -> Stream m a

    Generate a finite stream starting with the element from, enumerating the type up to the value to. If to is smaller than from then an empty stream is returned.

    Example1 expression
    Stream.toList $ Stream.enumerateFromTo 0 4[0,1,2,3,4]

    For Fractional types, the last element is equal to the specified to value after rounding to the nearest integral value.

    Example1 expression
    Stream.toList $ Stream.enumerateFromTo 1.1 4[1.1,2.1,3.1,4.1]
    Example1 expression
    Stream.toList $ Stream.enumerateFromTo 1.1 4.6[1.1,2.1,3.1,4.1,5.1]
  • enumerateFromThen :: Monad m => a -> a -> Stream m a

    enumerateFromThen from then generates a stream whose first element is from, the second element is then and the successive elements are in increments of then - from. Enumeration can occur downwards or upwards depending on whether then comes before or after from. For Bounded types the stream ends when maxBound is reached, for unbounded types it keeps enumerating infinitely.

    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 2[0,2,4,6]
    Example1 expression
    Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 (-2)[0,-2,-4,-6]
  • enumerateFromThenTo :: Monad m => a -> a -> a -> Stream m a

    enumerateFromThenTo from then to generates a finite stream whose first element is from, the second element is then and the successive elements are in increments of then - from up to to. Enumeration can occur downwards or upwards depending on whether then comes before or after from.

    Example1 expression
    Stream.toList $ Stream.enumerateFromThenTo 0 2 6[0,2,4,6]
    Example1 expression
    Stream.toList $ Stream.enumerateFromThenTo 0 (-2) (-6)[0,-2,-4,-6]
Instances21Enumerable, …
  • Enumerable IntegerDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable NaturalDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int16Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int32Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int64Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Int8Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word16Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word32Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word64Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable Word8Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable BoolDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable CharDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable DoubleDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable FloatDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable IntDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable OrderingDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable WordDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable ()Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Integral a => Enumerable (Ratio a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • Enumerable a => Enumerable (Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
  • HasResolution a => Enumerable (Fixed a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Generate
valueiterate :: Monad m => (a -> a) -> a -> Stream m a
#

Generate an infinite stream with x as the first element and each successive element derived by applying the function f on the previous element.

Example1 expression
Stream.toList $ Stream.take 5 $ Stream.iterate (+1) 1[1,2,3,4,5]
valueiterateM :: Monad m => (a -> m a) -> m a -> Stream m a
#

Generate an infinite stream with the first element generated by the action m and each successive element derived by applying the monadic function f on the previous element.

Example1 expression
:{Stream.iterateM (\x -> print x >> return (x + 1)) (return 0)    & Stream.take 3    & Stream.toList:}01[0,1,2]
valuefold :: Monad m => Fold m a b -> Stream m a -> m b
#

Fold a stream using the supplied left Fold and reducing the resulting expression strictly at each step. The behavior is similar to foldl'. A Fold can terminate early without consuming the full stream. See the documentation of individual Folds for termination behavior.

Definitions:

Example2 expressions
fold f = fmap fst . Stream.foldBreak ffold f = Stream.parse (Parser.fromFold f)

Example:

Example1 expression
Stream.fold Fold.sum (Stream.enumerateFromTo 1 100)5050
valueuncons :: Monad m => Stream m a -> m (Maybe (a, Stream m a))
#

Decompose a stream into its head and tail. If the stream is empty, returns Nothing. If the stream is non-empty, returns Just (a, ma), where a is the head of the stream and ma its tail.

Properties:

Example2 expressions
Nothing <- Stream.uncons Stream.nilJust ("a", t) <- Stream.uncons (Stream.cons "a" Stream.nil)

This can be used to consume the stream in an imperative manner one element at a time, as it just breaks down the stream into individual elements and we can loop over them as we deem fit. For example, this can be used to convert a streamly stream into other stream types.

All the folds in this module can be expressed in terms of uncons, however, this is generally less efficient than specific folds because it takes apart the stream one element at a time, therefore, does not take adavantage of stream fusion.

foldBreak is a more general way of consuming a stream piecemeal.

Example1 expression
:{uncons xs = do    r <- Stream.foldBreak Fold.one xs    return $ case r of        (Nothing, _) -> Nothing        (Just h, t) -> Just (h, t):}
valuefoldrM :: Monad m => (a -> m b -> m b) -> m b -> Stream m a -> m b
#

Right associative/lazy pull fold. foldrM build final stream constructs an output structure using the step function build. build is invoked with the next input element and the remaining (lazy) tail of the output structure. It builds a lazy output expression using the two. When the "tail structure" in the output expression is evaluated it calls build again thus lazily consuming the input stream until either the output expression built by build is free of the "tail" or the input is exhausted in which case final is used as the terminating case for the output structure. For more details see the description in the previous section.

Example, determine if any element is odd in a stream:

Example3 expressions
s = Stream.fromList (2:4:5:undefined)step x xs = if odd x then return True else xsStream.foldrM step (return False) sTrue
valuefoldr :: Monad m => (a -> b -> b) -> b -> Stream m a -> m b
#

Right fold, lazy for lazy monads and pure streams, and strict for strict monads.

Please avoid using this routine in strict monads like IO unless you need a strict right fold. This is provided only for use in lazy monads (e.g. Identity) or pure streams. Note that with this signature it is not possible to implement a lazy foldr when the monad m is strict. In that case it would be strict in its accumulator and therefore would necessarily consume all its input.

Example1 expression
foldr f z = Stream.foldrM (\a b -> f a <$> b) (return z)

Note: This is similar to Fold.foldr' (the right fold via left fold), but could be more efficient.

valueisPrefixOf :: (Monad m, Eq a) => Stream m a -> Stream m a -> m Bool
#

Returns True if the first stream is the same as or a prefix of the second. A stream is a prefix of itself.

Example1 expression
Stream.isPrefixOf (Stream.fromList "hello") (Stream.fromList "hello" :: Stream IO Char)True
valueisSubsequenceOf :: (Monad m, Eq a) => Stream m a -> Stream m a -> m Bool
#

Returns True if all the elements of the first stream occur, in order, in the second stream. The elements do not have to occur consecutively. A stream is a subsequence of itself.

Example1 expression
Stream.isSubsequenceOf (Stream.fromList "hlo") (Stream.fromList "hello" :: Stream IO Char)True
valuestripPrefix
  1. :: (Monad m, Eq a)
  2. => Stream m a
  3. -> Stream m a
  4. -> m (Maybe (Stream m a))
#

stripPrefix prefix input strips the prefix stream from the input stream if it is a prefix of input. Returns Nothing if the input does not start with the given prefix, stripped input otherwise. Returns Just nil when the prefix is the same as the input stream.

Space: O(1)

valuesequence :: Monad m => Stream m (m a) -> Stream m a
#
Example1 expression
sequence = Stream.mapM id

Replace the elements of a stream of monadic actions with the outputs of those actions.

Example2 expressions
s = Stream.fromList [putStr "a", putStr "b", putStrLn "c"]Stream.fold Fold.drain $ Stream.sequence sabc
valuemapM :: Monad m => (a -> m b) -> Stream m a -> Stream m b
#
Example1 expression
mapM f = Stream.sequence . fmap f

Apply a monadic function to each element of the stream and replace it with the output of the resulting action.

Example2 expressions
s = Stream.fromList ["a", "b", "c"]Stream.fold Fold.drain $ Stream.mapM putStr sabc
valuetrace :: Monad m => (a -> m b) -> Stream m a -> Stream m a
#

Apply a monadic function to each element flowing through the stream and discard the results.

Example2 expressions
s = Stream.enumerateFromTo 1 2Stream.fold Fold.drain $ Stream.trace print s12

Compare with tap.

valuetap :: Monad m => Fold m a b -> Stream m a -> Stream m a
#

Tap the data flowing through a stream into a Fold. For example, you may add a tap to log the contents flowing through the stream. The fold is used only for effects, its result is discarded.

                  Fold m a b
                      |
-----stream m a ---------------stream m a-----

Example2 expressions
s = Stream.enumerateFromTo 1 2Stream.fold Fold.drain $ Stream.tap (Fold.drainMapM print) s12

Compare with trace.

valuedelay :: MonadIO m => Double -> Stream m a -> Stream m a
#

Introduce a delay of specified seconds between elements of the stream.

Definition:

Example2 expressions
sleep n = liftIO $ threadDelay $ round $ n * 1000000delay = Stream.intersperseM_ . sleep

Example:

Example2 expressions
input = Stream.enumerateFromTo 1 3Stream.fold (Fold.drainMapM print) $ Stream.delay 1 input123
valuescan :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Strict left scan. Scan a stream using the given monadic fold.

Example2 expressions
s = Stream.fromList [1..10]Stream.fold Fold.toList $ Stream.takeWhile (< 10) $ Stream.scan Fold.sum s[0,1,3,6]

See also: usingStateT

valuepostscan :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Postscan a stream using the given monadic fold.

The following example extracts the input stream up to a point where the running average of elements is no more than 10:

Example4 expressions
import Data.Maybe (fromJust)let avg = Fold.teeWith (/) Fold.sum (fmap fromIntegral Fold.length)s = Stream.enumerateFromTo 1.0 100.0:{ Stream.fold Fold.toList  $ fmap (fromJust . fst)  $ Stream.takeWhile (\(_,x) -> x <= 10)  $ Stream.postscan (Fold.tee Fold.latest avg) s:}[1.0,2.0,3.0,4.0,5.0,6.0,7.0,8.0,9.0,10.0,11.0,12.0,13.0,14.0,15.0,16.0,17.0,18.0,19.0]
valuefilter :: Monad m => (a -> Bool) -> Stream m a -> Stream m a
#

Include only those elements that pass a predicate.

Example3 expressions
filter p = Stream.filterM (return . p)filter p = Stream.mapMaybe (\x -> if p x then Just x else Nothing)filter p = Stream.scanMaybe (Fold.filtering p)
valuefilterM :: Monad m => (a -> m Bool) -> Stream m a -> Stream m a
#

Same as filter but with a monadic predicate.

Example2 expressions
f p x = p x >>= \r -> return $ if r then Just x else NothingfilterM p = Stream.mapMaybeM (f p)
valuedropWhile :: Monad m => (a -> Bool) -> Stream m a -> Stream m a
#

Drop elements in the stream as long as the predicate succeeds and then take the rest of the stream.

valueinsertBy :: Monad m => (a -> a -> Ordering) -> a -> Stream m a -> Stream m a
#

insertBy cmp elem stream inserts elem before the first element in stream that is less than elem when compared using cmp.

Example1 expression
insertBy cmp x = Stream.mergeBy cmp (Stream.fromPure x)
Example2 expressions
input = Stream.fromList [1,3,5]Stream.fold Fold.toList $ Stream.insertBy compare 2 input[1,2,3,5]
valueintersperseM :: Monad m => m a -> Stream m a -> Stream m a
#

Insert an effect and its output before consuming an element of a stream except the first one.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.trace putChar $ Stream.intersperseM (putChar '.' >> return ',') inputh.,e.,l.,l.,o"h,e,l,l,o"

Be careful about the order of effects. In the above example we used trace after the intersperse, if we use it before the intersperse the output would be he.l.l.o."h,e,l,l,o".

Example1 expression
Stream.fold Fold.toList $ Stream.intersperseM (putChar '.' >> return ',') $ Stream.trace putChar inputhe.l.l.o."h,e,l,l,o"
valueintersperse :: Monad m => a -> Stream m a -> Stream m a
#

Insert a pure value between successive elements of a stream.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.intersperse ',' input"h,e,l,l,o"
valuereverse :: Monad m => Stream m a -> Stream m a
#

Returns the elements of the stream in reverse order. The stream must be finite. Note that this necessarily buffers the entire stream in memory.

Definition:

Example1 expression
reverse m = Stream.concatEffect $ Stream.fold Fold.toListRev m >>= return . Stream.fromList
valueindexed :: Monad m => Stream m a -> Stream m (Int, a)
#
Example4 expressions
f = Fold.foldl' (\(i, _) x -> (i + 1, x)) (-1,undefined)indexed = Stream.postscan findexed = Stream.zipWith (,) (Stream.enumerateFrom 0)indexedR n = fmap (\(i, a) -> (n - i, a)) . indexed

Pair each element in a stream with its index, starting from index 0.

Example1 expression
Stream.fold Fold.toList $ Stream.indexed $ Stream.fromList "hello"[(0,'h'),(1,'e'),(2,'l'),(3,'l'),(4,'o')]
valuemapMaybe :: Monad m => (a -> Maybe b) -> Stream m a -> Stream m b
#

Map a Maybe returning function to a stream, filter out the Nothing elements, and return a stream of values extracted from Just.

Equivalent to:

Example1 expression
mapMaybe f = Stream.catMaybes . fmap f
valuemapMaybeM :: Monad m => (a -> m (Maybe b)) -> Stream m a -> Stream m b
#

Like mapMaybe but maps a monadic function.

Equivalent to:

Example1 expression
mapMaybeM f = Stream.catMaybes . Stream.mapM f
Example1 expression
mapM f = Stream.mapMaybeM (\x -> Just <$> f x)
valuemergeBy
  1. :: Monad m
  2. => a -> a -> Ordering
  3. -> Stream m a
  4. -> Stream m a
  5. -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.mergeBy otherwise.

Merge two streams using a comparison function. The head elements of both the streams are compared and the smaller of the two elements is emitted, if both elements are equal then the element from the first stream is used first.

If the streams are sorted in ascending order, the resulting stream would also remain sorted in ascending order.

Example3 expressions
s1 = Stream.fromList [1,3,5]s2 = Stream.fromList [2,4,6,8]Stream.fold Fold.toList $ Stream.mergeBy compare s1 s2[1,2,3,4,5,6,8]
valuemergeByM
  1. :: Monad m
  2. => a -> a -> m Ordering
  3. -> Stream m a
  4. -> Stream m a
  5. -> Stream m a
#

Like mergeBy but with a monadic comparison function.

Example, to merge two streams randomly:

> randomly _ _ = randomIO >>= x -> return $ if x then LT else GT
> Stream.toList $ Stream.mergeByM randomly (Stream.fromList [1,1,1,1]) (Stream.fromList [2,2,2,2])
[2,1,2,2,2,1,1,1]

Example, merge two streams in a proportion of 2:1:

Example1 expression
:{do let s1 = Stream.fromList [1,1,1,1,1,1]     s2 = Stream.fromList [2,2,2] let proportionately m n = do      ref <- newIORef $ cycle $ Prelude.concat [Prelude.replicate m LT, Prelude.replicate n GT]      return $ \_ _ -> do         r <- readIORef ref         writeIORef ref $ Prelude.tail r         return $ Prelude.head r f <- proportionately 2 1 xs <- Stream.fold Fold.toList $ Stream.mergeByM f s1 s2 print xs:}[1,1,2,1,1,2,1,1,2]
valuezipWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m c
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.zipWith otherwise.

Stream a is evaluated first, followed by stream b, the resulting elements a and b are then zipped using the supplied zip function and the result c is yielded to the consumer.

If stream a or stream b ends, the zipped stream ends. If stream b ends first, the element a from previous evaluation of stream a is discarded.

Example3 expressions
s1 = Stream.fromList [1,2,3]s2 = Stream.fromList [4,5,6]Stream.fold Fold.toList $ Stream.zipWith (+) s1 s2[5,7,9]
valueunfoldMany :: Monad m => Unfold m a b -> Stream m a -> Stream m b
#

unfoldMany unfold stream uses unfold to map the input stream elements to streams and then flattens the generated streams into a single output stream.

Like concatMap but uses an Unfold for stream generation. Unlike concatMap this can fuse the Unfold code with the inner loop and therefore provide many times better performance.

valueintercalate :: Monad m => Unfold m b c -> b -> Stream m b -> Stream m c
#

intersperse followed by unfold and concat.

Example3 expressions
intercalate u a = Stream.unfoldMany u . Stream.intersperse aintersperse = Stream.intercalate Unfold.identityunwords = Stream.intercalate Unfold.fromList " "
Example2 expressions
input = Stream.fromList ["abc", "def", "ghi"]Stream.fold Fold.toList $ Stream.intercalate Unfold.fromList " " input"abc def ghi"
valueintercalateSuffix
  1. :: Monad m
  2. => Unfold m b c
  3. -> b
  4. -> Stream m b
  5. -> Stream m c
#

intersperseMSuffix followed by unfold and concat.

Example3 expressions
intercalateSuffix u a = Stream.unfoldMany u . Stream.intersperseMSuffix aintersperseMSuffix = Stream.intercalateSuffix Unfold.identityunlines = Stream.intercalateSuffix Unfold.fromList "\n"
Example2 expressions
input = Stream.fromList ["abc", "def", "ghi"]Stream.fold Fold.toList $ Stream.intercalateSuffix Unfold.fromList "\n" input"abc\ndef\nghi\n"
valueconcatMap :: Monad m => (a -> Stream m b) -> Stream m a -> Stream m b
#

Map a stream producing function on each element of the stream and then flatten the results into a single stream.

Example3 expressions
concatMap f = Stream.concatMapM (return . f)concatMap f = Stream.concat . fmap fconcatMap f = Stream.unfoldMany (Unfold.lmap f Unfold.fromStream)

See unfoldMany for a fusible alternative.

valueconcatMapM :: Monad m => (a -> m (Stream m b)) -> Stream m a -> Stream m b
#

Map a stream producing monadic function on each element of the stream and then flatten the results into a single stream. Since the stream generation function is monadic, unlike concatMap, it can produce an effect at the beginning of each iteration of the inner loop.

See unfoldMany for a fusible alternative.

valuefoldMany :: Monad m => Fold m a b -> Stream m a -> Stream m b
#

Apply a Fold repeatedly on a stream and emit the results in the output stream.

Definition:

Example1 expression
foldMany f = Stream.parseMany (Parser.fromFold f)

Example, empty stream:

Example3 expressions
f = Fold.take 2 Fold.sumfmany = Stream.fold Fold.toList . Stream.foldMany ffmany $ Stream.fromList [][]

Example, last fold empty:

Example1 expression
fmany $ Stream.fromList [1..4][3,7]

Example, last fold non-empty:

Example1 expression
fmany $ Stream.fromList [1..5][3,7,5]

Note that using a closed fold e.g. Fold.take 0, would result in an infinite stream on a non-empty input stream.

valuechunksOf :: (MonadIO m, Unbox a) => Int -> Stream m a -> Stream m (Array a)
#

chunksOf n stream groups the elements in the input stream into arrays of n elements each.

Same as the following but may be more efficient:

Example1 expression
chunksOf n = Stream.foldMany (Array.writeN n)

Pre-release

valuesplitOn :: Monad m => (a -> Bool) -> Fold m a b -> Stream m a -> Stream m b
#

Split on an infixed separator element, dropping the separator. The supplied Fold is applied on the split segments. Splits the stream on separator elements determined by the supplied predicate, separator is considered as infixed between two segments:

Example2 expressions
splitOn' p xs = Stream.fold Fold.toList $ Stream.splitOn p Fold.toList (Stream.fromList xs)splitOn' (== '.') "a.b"["a","b"]

An empty stream is folded to the default value of the fold:

Example1 expression
splitOn' (== '.') ""[""]

If one or both sides of the separator are missing then the empty segment on that side is folded to the default output of the fold:

Example1 expression
splitOn' (== '.') "."["",""]
Example1 expression
splitOn' (== '.') ".a"["","a"]
Example1 expression
splitOn' (== '.') "a."["a",""]
Example1 expression
splitOn' (== '.') "a..b"["a","","b"]

splitOn is an inverse of intercalating single element:

Stream.intercalate (Stream.fromPure '.') Unfold.fromList . Stream.splitOn (== '.') Fold.toList === id

Assuming the input stream does not contain the separator:

Stream.splitOn (== '.') Fold.toList . Stream.intercalate (Stream.fromPure '.') Unfold.fromList === id
valuewordsBy :: Monad m => (a -> Bool) -> Fold m a b -> Stream m a -> Stream m b
#

Split the stream after stripping leading, trailing, and repeated separators as per the fold supplied. Therefore, ".a..b." with . as the separator would be parsed as ["a","b"]. In other words, its like parsing words from whitespace separated text.

valuebefore :: Monad m => m b -> Stream m a -> Stream m a
#

Run the action m b before the stream yields its first element.

Same as the following but more efficient due to fusion:

Example2 expressions
before action xs = Stream.nilM action <> xsbefore action xs = Stream.concatMap (const xs) (Stream.fromEffect action)
valuehandle
  1. :: (MonadCatch m, Exception e)
  2. => e -> m (Stream m a)
  3. -> Stream m a
  4. -> Stream m a
#

When evaluating a stream if an exception occurs, stream evaluation aborts and the specified exception handler is run with the exception as argument. The exception is caught and handled unless the handler decides to rethrow it. Note that exception handling is not applied to the stream returned by the exception handler.

Observes exceptions only in the stream generation, and not in stream consumers.

Inhibits stream fusion

valuenilM :: Applicative m => m b -> Stream m a
#

A stream that terminates without producing any output, but produces a side effect.

Example1 expression
Stream.fold Fold.toList (Stream.nilM (print "nil"))"nil"[]

Pre-release

valueparse :: Monad m => Parser a m b -> Stream m a -> m (Either ParseError b)
#

Parse a stream using the supplied Parser.

Parsers (See Streamly.Internal.Data.Parser) are more powerful folds that add backtracking and error functionality to terminating folds. Unlike folds, parsers may not always result in a valid output, they may result in an error. For example:

Example1 expression
Stream.parse (Parser.takeEQ 1 Fold.drain) Stream.nilLeft (ParseError "takeEQ: Expecting exactly 1 elements, input terminated on 0")

Note: parse p is not the same as head . parseMany p on an empty stream.

valueintersperseM_ :: Monad m => m b -> Stream m a -> Stream m a
#

Insert a side effect before consuming an element of a stream except the first one.

Example2 expressions
input = Stream.fromList "hello"Stream.fold Fold.drain $ Stream.trace putChar $ Stream.intersperseM_ (putChar '.') inputh.e.l.l.o

Pre-release

valuecatMaybes :: Monad m => Stream m (Maybe a) -> Stream m a
#

In a stream of Maybes, discard Nothings and unwrap Justs.

Example2 expressions
catMaybes = Stream.mapMaybe idcatMaybes = fmap fromJust . Stream.filter isJust

Pre-release

valueappend :: Monad m => Stream m a -> Stream m a -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.append otherwise.

Fuses two streams sequentially, yielding all elements from the first stream, and then all elements from the second stream.

Example3 expressions
s1 = Stream.fromList [1,2]s2 = Stream.fromList [3,4]Stream.fold Fold.toList $ s1 `Stream.append` s2[1,2,3,4]
valueinterleave :: Monad m => Stream m a -> Stream m a -> Stream m a
#

WARNING! O(n^2) time complexity wrt number of streams. Suitable for statically fusing a small number of streams. Use the O(n) complexity StreamK.Streamly.Data.StreamK.interleave otherwise.

Interleaves two streams, yielding one element from each stream alternately. When one stream stops the rest of the other stream is used in the output stream.

valueparseMany
  1. :: Monad m
  2. => Parser a m b
  3. -> Stream m a
  4. -> Stream m (Either ParseError b)
#

Apply a Parser repeatedly on a stream and emit the parsed values in the output stream.

Example:

Example3 expressions
s = Stream.fromList [1..10]parser = Parser.takeBetween 0 2 Fold.sumStream.fold Fold.toList $ Stream.parseMany parser s[Right 3,Right 7,Right 11,Right 15,Right 19]

This is the streaming equivalent of the Streamly.Data.Parser.many parse combinator.

Known Issues: When the parser fails there is no way to get the remaining stream.

valuegroupsOf :: Monad m => Int -> Fold m a b -> Stream m a -> Stream m b
#

Group the input stream into groups of n elements each and then fold each group using the provided fold function.

groupsOf n f = foldMany (FL.take n f)
Example1 expression
Stream.toList $ Stream.groupsOf 2 Fold.sum (Stream.enumerateFromTo 1 10)[3,7,11,15,19]

This can be considered as an n-fold version of take where we apply take repeatedly on the leftover stream until the stream exhausts.

valuemorphInner :: Monad n => (forall x. m x -> n x) -> Stream m a -> Stream n a
#

Transform the inner monad of a stream using a natural transformation.

Example, generalize the inner monad from Identity to any other:

Example1 expression
generalizeInner = Stream.morphInner (return . runIdentity)

Also known as hoist.

valuecatEithers :: Monad m => Stream m (Either a a) -> Stream m a
#

Remove the either wrapper and flatten both lefts and as well as rights in the output stream.

Example1 expression
catEithers = fmap (either id id)

Pre-release

valuescanMaybe :: Monad m => Fold m a (Maybe b) -> Stream m a -> Stream m b
#

Use a filtering fold on a stream.

Example1 expression
scanMaybe f = Stream.catMaybes . Stream.postscan f
valueconcatEffect :: Monad m => m (Stream m a) -> Stream m a
#

Given a stream value in the underlying monad, lift and join the underlying monad with the stream monad.

Example2 expressions
concatEffect = Stream.concat . Stream.fromEffectconcatEffect eff = Stream.concatMapM (\() -> eff) (Stream.fromPure ())

See also: concat, sequence

valuecrossWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> Stream m a
  4. -> Stream m b
  5. -> Stream m c
#

Definition:

Example1 expression
crossWith f m1 m2 = fmap f m1 `Stream.crossApply` m2

Note that the second stream is evaluated multiple times.

valueafterIO :: MonadIO m => IO b -> Stream m a -> Stream m a
#

Run the action IO b whenever the stream is evaluated to completion, or if it is garbage collected after a partial lazy evaluation.

The semantics of the action IO b are similar to the semantics of cleanup action in bracketIO.

See also afterUnsafe

valuebracketIO3
  1. :: (MonadIO m, MonadCatch m)
  2. => IO b
  3. -> b -> IO c
  4. -> b -> IO d
  5. -> b -> IO e
  6. -> b -> Stream m a
  7. -> Stream m a
#

Like bracketIO but can use 3 separate cleanup actions depending on the mode of termination:

  1. When the stream stops normally

  2. When the stream is garbage collected

  3. When the stream encounters an exception

bracketIO3 before onStop onGC onException action runs action using the result of before. If the stream stops, onStop action is executed, if the stream is abandoned onGC is executed, if the stream encounters an exception onException is executed.

The exception is not caught, it is rethrown.

Inhibits stream fusion

Pre-release

valuefinallyIO :: (MonadIO m, MonadCatch m) => IO b -> Stream m a -> Stream m a
#

Run the action IO b whenever the stream stream stops normally, aborts due to an exception or if it is garbage collected after a partial lazy evaluation.

The semantics of running the action IO b are similar to the cleanup action semantics described in bracketIO.

Example1 expression
finallyIO release = Stream.bracketIO (return ()) (const release)

See also finallyUnsafe

Inhibits stream fusion

valuefoldBreak :: Monad m => Fold m a b -> Stream m a -> m (b, Stream m a)
#

Like fold but also returns the remaining stream. The resulting stream would be Stream.nil if the stream finished before the fold.

Concurrent Operations

0 declarations

Channels

At a lower level, concurrency is implemented using channels that support concurrent evaluation of streams. We create a channel, and add one or more streams to it. The channel evaluates multiple streams concurrently and then generates a single output stream from the results. How the streams are combined depends on the configuration of the channel.

Primitives

There are only a few fundamental abstractions for concurrency, parEval, parConcatMap, and parConcatIterate, all concurrency combinators can be expressed in terms of these.

parEval evaluates a stream as a whole asynchronously with respect to the consumer of the stream. A worker thread evaluates multiple elements of the stream ahead of time and buffers the results; the consumer of the stream runs in another thread consuming the elements from the buffer, thus decoupling the production and consumption of the stream. parEval can be used to run different stages of a pipeline concurrently.

parConcatMap is used to evaluate multiple actions in a stream concurrently with respect to each other or to evaluate multiple streams concurrently and combine the results. A stream generator function is mapped to the input stream and all the generated streams are then evaluated concurrently, and the results are combined.

parConcatIterate is like parConcatMap but iterates a stream generator function recursively over the stream. This can be used to traverse trees or graphs.

Configuration

Concurrent combinators take a Config argument which controls the concurrent behavior. For example, maximum number of threads to be used (maxThreads) or the maxmimum size of the buffer (maxBuffer), or how the streams are scheduled with respect to each other (interleaved), or how the results are consumed (ordered).

Configuration is specified as Config -> Config modifier functions that can be composed together using function composition. For example, to specify the maximum threads we can use parConcatMap (maxThreads 10) if we also want to specify the maximum buffer we can compose the two options parConcatMap (maxThreads 10 . maxBuffer 100). To use default configuration use id as the config modifier e.g. parConcatMap id.

See the Configuration section and individual configuration options' documentation for the default behavior and default values of configuration parameters.

Scheduling

The most important configuration option is to control whether the output of the concurrent execution is consumed in the same order as the corresponding actions in the input stream or as soon as they arrive. The default is the latter, however, we can enforce the original order by using the ordered option.

Another important option controls whether the number of worker threads are automatically increased and decreased based on the consumption rate or threads are started as aggresively as possible until the maxThreads or maxBuffer limits are hit. The default is the former. However, the eager option can be enabled to use the latter behavior. When eager is on, even if the stream consumer thread blocks it does not make any impact on the scheduling of the available tasks.

Combinators

Using the few fundamental concurrency primitives we can implement all the usual streaming combinators with concurrent behavior. Combinators like unfoldrM, iterateM that are inherently serial can be evaluated concurrently with respect to the consumer pipeline using parEval. Combinators like zipWithM, mergeByM can also use parEval on the input streams to evaluate them concurrently before combining.

Combinators like repeatM, replicateM, fromListM, sequence, mapM in which all actions are independent of each other can be made concurrent using the parConcatMap operation.

A concurrent repeatM repeats an action using multiple concurrent executions of the action. Similarly, a concurrent mapM performs the mapped action in independent threads.

Some common concurrent combinators are provided in this module.

Types

Configuration

datadata Config
#

An abstract type for specifying the configuration parameters of a Channel. Use Config -> Config modifier functions to modify the default configuration. See the individual modifier documentation for default values.

Limits

valuemaxThreads :: Int -> Config -> Config
#

Specify the maximum number of threads that can be spawned by the channel. A value of 0 resets the thread limit to default, a negative value means there is no limit. The default value is 1500.

When the actions in a stream are IO bound, having blocking IO calls, this option can be used to control the maximum number of in-flight IO requests. When the actions are CPU bound this option can be used to control the amount of CPU used by the stream.

valuemaxBuffer :: Int -> Config -> Config
#

Specify the maximum size of the buffer for storing the results from concurrent computations. If the buffer becomes full we stop spawning more concurrent tasks until there is space in the buffer. A value of 0 resets the buffer size to default, a negative value means there is no limit. The default value is 1500.

CAUTION! using an unbounded maxBuffer value (i.e. a negative value) coupled with an unbounded maxThreads value is a recipe for disaster in presence of infinite streams, or very large streams. Especially, it must not be used when pure is used in ZipAsyncM streams as pure in applicative zip streams generates an infinite stream causing unbounded concurrent generation with no limit on the buffer or threads.

Rate Control

datadata Rate
#

Specifies the stream yield rate in yields per second (Hertz). We keep accumulating yield credits at rateGoal. At any point of time we allow only as many yields as we have accumulated as per rateGoal since the start of time. If the consumer or the producer is slower or faster, the actual rate may fall behind or exceed rateGoal. We try to recover the gap between the two by increasing or decreasing the pull rate from the producer. However, if the gap becomes more than rateBuffer we try to recover only as much as rateBuffer.

rateLow puts a bound on how low the instantaneous rate can go when recovering the rate gap. In other words, it determines the maximum yield latency. Similarly, rateHigh puts a bound on how high the instantaneous rate can go when recovering the rate gap. In other words, it determines the minimum yield latency. We reduce the latency by increasing concurrency, therefore we can say that it puts an upper bound on concurrency.

If the rateGoal is 0 or negative the stream never yields a value. If the rateBuffer is 0 or negative we do not attempt to recover.

Constructors

valuerate :: Maybe Rate -> Config -> Config
#

Specify the stream evaluation rate of a channel.

A Nothing value means there is no smart rate control, concurrent execution blocks only if maxThreads or maxBuffer is reached, or there are no more concurrent tasks to execute. This is the default.

When rate (throughput) is specified, concurrent production may be ramped up or down automatically to achieve the specified stream throughput. The specific behavior for different styles of Rate specifications is documented under Rate. The effective maximum production rate achieved by a channel is governed by:

  • The maxThreads limit

  • The maxBuffer limit

  • The maximum rate that the stream producer can achieve

  • The maximum rate that the stream consumer can achieve

Maximum production rate is given by:

rate = \frac{maxThreads}{latency}

If we know the average latency of the tasks we can set maxThreads accordingly.

valueavgRate :: Double -> Config -> Config
#

Same as rate (Just $ Rate (r/2) r (2*r) maxBound)

Specifies the average production rate of a stream in number of yields per second (i.e. Hertz). Concurrent production is ramped up or down automatically to achieve the specified average yield rate. The rate can go down to half of the specified rate on the lower side and double of the specified rate on the higher side.

valueminRate :: Double -> Config -> Config
#

Same as rate (Just $ Rate r r (2*r) maxBound)

Specifies the minimum rate at which the stream should yield values. As far as possible the yield rate would never be allowed to go below the specified rate, even though it may possibly go above it at times, the upper limit is double of the specified rate.

valuemaxRate :: Double -> Config -> Config
#

Same as rate (Just $ Rate (r/2) r r maxBound)

Specifies the maximum rate at which the stream should yield values. As far as possible the yield rate would never be allowed to go above the specified rate, even though it may possibly go below it at times, the lower limit is half of the specified rate. This can be useful in applications where certain resource usage must not be allowed to go beyond certain limits.

valueconstRate :: Double -> Config -> Config
#

Same as rate (Just $ Rate r r r 0)

Specifies a constant yield rate. If for some reason the actual rate goes above or below the specified rate we do not try to recover it by increasing or decreasing the rate in future. This can be useful in applications like graphics frame refresh where we need to maintain a constant refresh rate.

Stop behavior

datadata StopWhen
#

Specify when the Channel should stop.

Constructors

Scheduling behavior

valueeager :: Bool -> Config -> Config
#

By default, processing of output from the worker threads is given priority over dispatching new workers. More workers are dispatched only when there is no output to process. When eager is set to True, workers are dispatched aggresively as long as there is more work to do irrespective of whether there is output pending to be processed by the stream consumer. However, dispatching may stop if maxThreads or maxBuffer is reached.

Note: This option has no effect when rate has been specified.

Note: Not supported with interleaved.

valueordered :: Bool -> Config -> Config
#

When enabled the streams may be evaluated cocnurrently but the results are produced in the same sequence as a serial evaluation would produce.

Note: Not supported with interleaved.

valueinterleaved :: Bool -> Config -> Config
#

Interleave the streams fairly instead of prioritizing the left stream. This schedules all streams in a round robin fashion over limited number of threads.

Note: Can only be used on finite number of streams.

Note: Not supported with ordered.

Diagnostics

Combinators

Stream combinators using a concurrent channel.

Evaluate

Evaluate a stream as a whole concurrently with respect to the consumer of the stream.

valueparEval :: MonadAsync m => (Config -> Config) -> Stream m a -> Stream m a
#

parEval evaluates a stream as a whole asynchronously with respect to the consumer of the stream. A worker thread evaluates multiple elements of the stream ahead of time and buffers the results; the consumer of the stream runs in another thread consuming the elements from the buffer, thus decoupling the production and consumption of the stream. parEval can be used to run different stages of a pipeline concurrently.

It is important to note that parEval does not evaluate individual actions in the stream concurrently with respect to each other, it merely evaluates the stream serially but in a different thread than the consumer thread, thus the consumer and producer can run concurrently. See parMapM and parSequence to evaluate actions in the stream concurrently.

The evaluation requires only one thread as only one stream needs to be evaluated. Therefore, the concurrency options that are relevant to multiple streams do not apply here e.g. maxThreads, eager, interleaved, ordered, stopWhen options do not have any effect on parEval.

Useful idioms:

Example2 expressions
parUnfoldrM step = Stream.parEval id . Stream.unfoldrM stepparIterateM step = Stream.parEval id . Stream.iterateM step

Generate

Generate a stream by evaluating multiple actions concurrently.

valueparRepeatM :: MonadAsync m => (Config -> Config) -> m a -> Stream m a
#

Definition:

Example1 expression
parRepeatM cfg = Stream.parSequence cfg . Stream.repeat

Generate a stream by repeatedly executing a monadic action forever.

valueparReplicateM
  1. :: MonadAsync m
  2. => Config -> Config
  3. -> Int
  4. -> m a
  5. -> Stream m a
#

Generate a stream by concurrently performing a monadic action n times.

Definition:

Example1 expression
parReplicateM cfg n = Stream.parSequence cfg . Stream.replicate n

Example, parReplicateM in the following example executes all the replicated actions concurrently, thus taking only 1 second:

Example1 expression
Stream.fold Fold.drain $ Stream.parReplicateM id 10 $ delay 1...
valuefromCallback :: MonadAsync m => ((a -> m ()) -> m ()) -> Stream m a
#

fromCallback f creates an entangled pair of a callback and a stream i.e. whenever the callback is called a value appears in the stream. The function f is invoked with the callback as argument, and the stream is returned. f would store the callback for calling it later for generating values in the stream.

The callback queues a value to a concurrent channel associated with the stream. The stream can be evaluated safely in any thread.

Pre-release

Map

Map actions on a stream such that the mapped actions are evaluated concurrently with each other.

valueparMapM
  1. :: MonadAsync m
  2. => Config -> Config
  3. -> a -> m b
  4. -> Stream m a
  5. -> Stream m b
#

Definition:

Example1 expression
parMapM modifier f = Stream.parConcatMap modifier (Stream.fromEffect . f)

For example, the following finishes in 3 seconds (as opposed to 6 seconds) because all actions run in parallel. Even though results are available out of order they are ordered due to the config option:

Example2 expressions
f x = delay x >> return xStream.fold Fold.toList $ Stream.parMapM (Stream.ordered True) f $ Stream.fromList [3,2,1]1 sec2 sec3 sec[3,2,1]
valueparSequence
  1. :: MonadAsync m
  2. => Config -> Config
  3. -> Stream m (m a)
  4. -> Stream m a
#

Definition:

Example1 expression
parSequence modifier = Stream.parMapM modifier id

Useful idioms:

Example2 expressions
parFromListM = Stream.parSequence id . Stream.fromListparFromFoldableM = Stream.parSequence id . StreamK.toStream . StreamK.fromFoldable

Combine two

Combine two streams such that each stream as a whole is evaluated concurrently with respect to the other stream as well as the consumer of the resulting stream.

valueparZipWithM
  1. :: MonadAsync m
  2. => Config -> Config
  3. -> a -> b -> m c
  4. -> Stream m a
  5. -> Stream m b
  6. -> Stream m c
#

Evaluates the streams being zipped in separate threads than the consumer. The zip function is evaluated in the consumer thread.

Example1 expression
parZipWithM cfg f m1 m2 = Stream.zipWithM f (Stream.parEval cfg m1) (Stream.parEval cfg m2)

Multi-stream concurrency options won't apply here, see the notes in parEval.

If you want to evaluate the zip function as well in a separate thread, you can use a parEval on parZipWithM.

valueparZipWith
  1. :: MonadAsync m
  2. => Config -> Config
  3. -> a -> b -> c
  4. -> Stream m a
  5. -> Stream m b
  6. -> Stream m c
#
Example1 expression
parZipWith cfg f = Stream.parZipWithM cfg (\a b -> return $ f a b)
Example3 expressions
m1 = Stream.fromList [1,2,3]m2 = Stream.fromList [4,5,6]Stream.fold Fold.toList $ Stream.parZipWith id (,) m1 m2[(1,4),(2,5),(3,6)]
valueparMergeByM
  1. :: MonadAsync m
  2. => Config -> Config
  3. -> a -> a -> m Ordering
  4. -> Stream m a
  5. -> Stream m a
  6. -> Stream m a
#

Like mergeByM but evaluates both the streams concurrently.

Definition:

Example1 expression
parMergeByM cfg f m1 m2 = Stream.mergeByM f (Stream.parEval cfg m1) (Stream.parEval cfg m2)

List of streams

Shares a single channel across many streams.

Stream of streams

Apply

Concat

Shares a single channel across many streams.

valueparConcatMap
  1. :: MonadAsync m
  2. => Config -> Config
  3. -> a -> Stream m b
  4. -> Stream m a
  5. -> Stream m b
#

Map each element of the input to a stream and then concurrently evaluate and concatenate the resulting streams. Multiple streams may be evaluated concurrently but earlier streams are perferred. Output from the streams are used as they arrive.

Definition:

Example1 expression
parConcatMap modifier f stream = Stream.parConcat modifier $ fmap f stream

Examples:

Example1 expression
f cfg xs = Stream.fold Fold.toList $ Stream.parConcatMap cfg id $ Stream.fromList xs

The following streams finish in 4 seconds:

Example4 expressions
stream1 = Stream.fromEffect (delay 4)stream2 = Stream.fromEffect (delay 2)stream3 = Stream.fromEffect (delay 1)f id [stream1, stream2, stream3]1 sec2 sec4 sec[1,2,4]

Limiting threads to 2 schedules the third stream only after one of the first two has finished, releasing a thread:

Example1 expression
f (Stream.maxThreads 2) [stream1, stream2, stream3]...[2,1,4]

When used with a Single thread it behaves like serial concatMap:

Example1 expression
f (Stream.maxThreads 1) [stream1, stream2, stream3]...[4,2,1]
Example3 expressions
stream1 = Stream.fromList [1,2,3]stream2 = Stream.fromList [4,5,6]f (Stream.maxThreads 1) [stream1, stream2][1,2,3,4,5,6]

Schedule all streams in a round robin fashion over the available threads:

Example1 expression
f cfg xs = Stream.fold Fold.toList $ Stream.parConcatMap (Stream.interleaved True . cfg) id $ Stream.fromList xs
Example3 expressions
stream1 = Stream.fromList [1,2,3]stream2 = Stream.fromList [4,5,6]f (Stream.maxThreads 1) [stream1, stream2][1,4,2,5,3,6]

ConcatIterate

Observation

valueparTapCount
  1. :: MonadAsync m
  2. => a -> Bool
  3. -> Stream m Int -> m b
  4. -> Stream m a
  5. -> Stream m a
#

parTapCount predicate fold stream taps the count of those elements in the stream that pass the predicate. The resulting count stream is sent to a fold running concurrently in another thread.

For example, to print the count of elements processed every second:

Example4 expressions
rate = Stream.rollingMap2 (flip (-)) . Stream.delayPost 1report = Stream.fold (Fold.drainMapM print) . ratetap = Stream.parTapCount (const True) reportgo = Stream.fold Fold.drain $ tap $ Stream.enumerateFrom 0

Note: This may not work correctly on 32-bit machines because of Int overflow.

Pre-release

Time Related

0 declarations

Timers

valueinterject :: MonadAsync m => m a -> Double -> Stream m a -> Stream m a
#

Intersperse a monadic action into the input stream after every n seconds.

Definition:

Example1 expression
interject n f xs = Stream.parListEagerFst [xs, Stream.periodic f n]

Example:

Example3 expressions
s = Stream.fromList "hello"input = Stream.mapM (\x -> threadDelay 1000000 >> putChar x) sStream.fold Fold.drain $ Stream.interject (putChar ',') 1.05 inputh,e,l,l,o

Trimming

valuetakeInterval :: MonadAsync m => Double -> Stream m a -> Stream m a
#

takeInterval interval runs the stream only upto the specified time interval in seconds.

The interval starts when the stream is evaluated for the first time.

valuedropInterval :: MonadAsync m => Double -> Stream m a -> Stream m a
#

dropInterval interval drops all the stream elements that are generated before the specified interval in seconds has passed.

The interval begins when the stream is evaluated for the first time.

Chunking

valueintervalsOf
  1. :: MonadAsync m
  2. => Double
  3. -> Fold m a b
  4. -> Stream m a
  5. -> Stream m b
#

Group the input stream into windows of n second each and then fold each group using the provided fold function.

Example3 expressions
twoPerSec = Stream.parEval (Stream.constRate 2) $ Stream.enumerateFrom 1intervals = Stream.intervalsOf 1 Fold.toList twoPerSecStream.fold Fold.toList $ Stream.take 2 intervals[...,...]

Sampling

valuesampleIntervalEnd :: MonadAsync m => Double -> Stream m a -> Stream m a
#

Continuously evaluate the input stream and sample the last event in each time window of n seconds.

This is also known as throttle in some libraries.

Example1 expression
sampleIntervalEnd n = Stream.catMaybes . Stream.intervalsOf n Fold.latest
valuesampleIntervalStart :: MonadAsync m => Double -> Stream m a -> Stream m a
#

Like sampleInterval but samples at the beginning of the time window.

Example1 expression
sampleIntervalStart n = Stream.catMaybes . Stream.intervalsOf n Fold.one
valuesampleBurstEnd :: MonadAsync m => Double -> Stream m a -> Stream m a
#

Sample one event at the end of each burst of events. A burst is a group of events close together in time, it ends when an event is spaced by more than the specified time interval (in seconds) from the previous event.

This is known as debounce in some libraries.

The clock granularity is 10 ms.

Lifted Exceptions

3 declarations
valueafter
  1. :: (MonadIO m, MonadBaseControl IO m)
  2. => m b
  3. -> Stream m a
  4. -> Stream m a
#

Run the action m b whenever the stream Stream m a stops normally, or if it is garbage collected after a partial lazy evaluation.

The semantics of the action m b are similar to the semantics of cleanup action in bracket.

See also after_

valuebracket
  1. :: (MonadAsync m, MonadCatch m)
  2. => m b
  3. -> b -> m c
  4. -> b -> Stream m a
  5. -> Stream m a
#

Run the alloc action IO b with async exceptions disabled but keeping blocking operations interruptible (see mask). Use the output b of the IO action as input to the function b -> Stream m a to generate an output stream.

b is usually a resource under the IO monad, e.g. a file handle, that requires a cleanup after use. The cleanup action b -> m c, runs whenever (1) the stream ends normally, (2) due to a sync or async exception or, (3) if it gets garbage collected after a partial lazy evaluation. The exception is not caught, it is rethrown.

bracket only guarantees that the cleanup action runs, and it runs with async exceptions enabled. The action must ensure that it can successfully cleanup the resource in the face of sync or async exceptions.

When the stream ends normally or on a sync exception, cleanup action runs immediately in the current thread context, whereas in other cases it runs in the GC context, therefore, cleanup may be delayed until the GC gets to run.

See also: bracket_

Inhibits stream fusion

valuefinally :: (MonadAsync m, MonadCatch m) => m b -> Stream m a -> Stream m a
#

Run the action m b whenever the stream Stream m a stops normally, aborts due to an exception or if it is garbage collected after a partial lazy evaluation.

The semantics of running the action m b are similar to the cleanup action semantics described in bracket.

Example1 expression
finally action xs = Stream.bracket (return ()) (const action) (const xs)

See also finally_

Inhibits stream fusion

Deprecated