Generate an infinite stream by repeating a pure value.
repeat x = Stream.repeatM (pure x):: a typeCtrl KGHC 9.10.3 · lts/ghc-9.10.x · 248f8f0 · 2026-10-05
Modulestreamly-0.10.1Haskell2010
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.
All Streamly.Data.Stream combinators are re-exported via this module. For more pre-release combinators also see Streamly.Internal.Data.Stream module.
Generate an infinite stream by repeating a pure value.
repeat x = Stream.repeatM (pure x)A stream consists of a step function that generates the next step given a current state, and the current state.
Monad m => Functor (Stream m)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Type(Foldable m, Monad m) => Foldable (Stream m)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.TypeIsList (Stream Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.TypeEq a => Eq (Stream Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.TypeOrd a => Ord (Stream Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.TypeRead a => Read (Stream Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.TypeShow a => Show (Stream Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Typea ~ Char => IsString (Stream Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Typetype Item (Stream Identity a) = aDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.TypeRun 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
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
Construct a stream from a list of pure values.
Definitions:
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.
A stream that terminates without producing any output or side effect.
Stream.toList Stream.nil[]
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::
s = 1 `Stream.cons` Stream.fromList [2,3]Stream.toList s[1,2,3]
Definition:
cons x xs = return x `Stream.consM` xsLike cons but fuses an effect instead of a pure value.
Convert an Unfold into a stream by supplying it an input seed.
s = Stream.unfold Unfold.replicateM (3, putStrLn "hello")Stream.fold Fold.drain shellohellohello
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,
:{let f b = if b > 2 then Nothing else Just (b, b + 1)in Stream.toList $ Stream.unfoldr f 0:}[0,1,2]
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,
:{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]
Create a singleton stream from a pure value.
fromPure a = a `Stream.cons` Stream.nilfromPure = purefromPure = Stream.fromEffect . pure
Create a singleton stream from a monadic action.
fromEffect m = m `Stream.consM` Stream.nilfromEffect = Stream.sequence . Stream.fromPure
Stream.fold Fold.drain $ Stream.fromEffect (putStrLn "hello")hello
repeatM = Stream.sequence . Stream.repeatGenerate a stream by repeatedly executing a monadic action forever.
:{repeatAction = Stream.repeatM (threadDelay 1000000 >> print 1) & Stream.take 10 & Stream.fold Fold.drain:}
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.
replicateM n = Stream.sequence . Stream.replicate nGenerate a stream by performing a monadic action n times.
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.
enumerateFrom :: Monad m => a -> Stream m aenumerateFrom 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.
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.
Stream.toList $ Stream.take 4 $ Stream.enumerateFrom 1.1[1.1,2.1,3.1,4.1]
enumerateFromTo :: Monad m => a -> a -> Stream m aGenerate 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.
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.
Stream.toList $ Stream.enumerateFromTo 1.1 4[1.1,2.1,3.1,4.1]
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 aenumerateFromThen 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.
Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 2[0,2,4,6]
Stream.toList $ Stream.take 4 $ Stream.enumerateFromThen 0 (-2)[0,-2,-4,-6]
enumerateFromThenTo :: Monad m => a -> a -> a -> Stream m aenumerateFromThenTo 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.
Stream.toList $ Stream.enumerateFromThenTo 0 2 6[0,2,4,6]
Stream.toList $ Stream.enumerateFromThenTo 0 (-2) (-6)[0,-2,-4,-6]
Enumerable IntegerDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable NaturalDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Int16Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Int32Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Int64Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Int8Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Word16Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Word32Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Word64Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable Word8Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable BoolDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable CharDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable DoubleDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable FloatDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable IntDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable OrderingDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable WordDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable ()Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateIntegral a => Enumerable (Ratio a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateEnumerable a => Enumerable (Identity a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateHasResolution a => Enumerable (Fixed a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.GenerateGenerate an infinite stream with x as the first element and each
successive element derived by applying the function f on the previous
element.
Stream.toList $ Stream.take 5 $ Stream.iterate (+1) 1[1,2,3,4,5]
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.
:{Stream.iterateM (\x -> print x >> return (x + 1)) (return 0) & Stream.take 3 & Stream.toList:}01[0,1,2]
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:
fold f = fmap fst . Stream.foldBreak ffold f = Stream.parse (Parser.fromFold f)
Example:
Stream.fold Fold.sum (Stream.enumerateFromTo 1 100)5050
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:
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.
:{uncons xs = do r <- Stream.foldBreak Fold.one xs return $ case r of (Nothing, _) -> Nothing (Just h, t) -> Just (h, t):}
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:
s = Stream.fromList (2:4:5:undefined)step x xs = if odd x then return True else xsStream.foldrM step (return False) sTrue
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.
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.
Compare two streams for equality
Compare two streams lexicographically.
Returns True if the first stream is the same as or a prefix of the second. A stream is a prefix of itself.
Stream.isPrefixOf (Stream.fromList "hello") (Stream.fromList "hello" :: Stream IO Char)True
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.
Stream.isSubsequenceOf (Stream.fromList "hlo") (Stream.fromList "hello" :: Stream IO Char)True
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)
sequence = Stream.mapM idReplace the elements of a stream of monadic actions with the outputs of those actions.
s = Stream.fromList [putStr "a", putStr "b", putStrLn "c"]Stream.fold Fold.drain $ Stream.sequence sabc
mapM f = Stream.sequence . fmap fApply a monadic function to each element of the stream and replace it with the output of the resulting action.
s = Stream.fromList ["a", "b", "c"]Stream.fold Fold.drain $ Stream.mapM putStr sabc
Apply a monadic function to each element flowing through the stream and discard the results.
s = Stream.enumerateFromTo 1 2Stream.fold Fold.drain $ Stream.trace print s12
Compare with tap.
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-----
s = Stream.enumerateFromTo 1 2Stream.fold Fold.drain $ Stream.tap (Fold.drainMapM print) s12
Compare with trace.
Introduce a delay of specified seconds between elements of the stream.
Definition:
sleep n = liftIO $ threadDelay $ round $ n * 1000000delay = Stream.intersperseM_ . sleep
Example:
input = Stream.enumerateFromTo 1 3Stream.fold (Fold.drainMapM print) $ Stream.delay 1 input123
Strict left scan. Scan a stream using the given monadic fold.
s = Stream.fromList [1..10]Stream.fold Fold.toList $ Stream.takeWhile (< 10) $ Stream.scan Fold.sum s[0,1,3,6]
See also: usingStateT
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:
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]
Include only those elements that pass a predicate.
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)
Same as filter but with a monadic predicate.
f p x = p x >>= \r -> return $ if r then Just x else NothingfilterM p = Stream.mapMaybeM (f p)
Take first n elements from the stream and discard the rest.
End the stream as soon as the predicate fails on an element.
Same as takeWhile but with a monadic predicate.
Discard first n elements from the stream and take the rest.
Drop elements in the stream as long as the predicate succeeds and then take the rest of the stream.
Same as dropWhile but with a monadic predicate.
insertBy cmp elem stream inserts elem before the first element in
stream that is less than elem when compared using cmp.
insertBy cmp x = Stream.mergeBy cmp (Stream.fromPure x)input = Stream.fromList [1,3,5]Stream.fold Fold.toList $ Stream.insertBy compare 2 input[1,2,3,5]
Insert an effect and its output before consuming an element of a stream except the first one.
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".
Stream.fold Fold.toList $ Stream.intersperseM (putChar '.' >> return ',') $ Stream.trace putChar inputhe.l.l.o."h,e,l,l,o"
Insert a pure value between successive elements of a stream.
input = Stream.fromList "hello"Stream.fold Fold.toList $ Stream.intersperse ',' input"h,e,l,l,o"
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:
reverse m = Stream.concatEffect $ Stream.fold Fold.toListRev m >>= return . Stream.fromListf = 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.
Stream.fold Fold.toList $ Stream.indexed $ Stream.fromList "hello"[(0,'h'),(1,'e'),(2,'l'),(3,'l'),(4,'o')]
Like mapMaybe but maps a monadic function.
Equivalent to:
mapMaybeM f = Stream.catMaybes . Stream.mapM fmapM f = Stream.mapMaybeM (\x -> Just <$> f x)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.
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]
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:
:{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]
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.
s1 = Stream.fromList [1,2,3]s2 = Stream.fromList [4,5,6]Stream.fold Fold.toList $ Stream.zipWith (+) s1 s2[5,7,9]
Like zipWith but using a monadic zipping function.
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.
intersperse followed by unfold and concat.
intercalate u a = Stream.unfoldMany u . Stream.intersperse aintersperse = Stream.intercalate Unfold.identityunwords = Stream.intercalate Unfold.fromList " "
input = Stream.fromList ["abc", "def", "ghi"]Stream.fold Fold.toList $ Stream.intercalate Unfold.fromList " " input"abc def ghi"
intersperseMSuffix followed by unfold and concat.
intercalateSuffix u a = Stream.unfoldMany u . Stream.intersperseMSuffix aintersperseMSuffix = Stream.intercalateSuffix Unfold.identityunlines = Stream.intercalateSuffix Unfold.fromList "\n"
input = Stream.fromList ["abc", "def", "ghi"]Stream.fold Fold.toList $ Stream.intercalateSuffix Unfold.fromList "\n" input"abc\ndef\nghi\n"
Map a stream producing function on each element of the stream and then flatten the results into a single stream.
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.
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.
Apply a Fold repeatedly on a stream and emit the results in the output stream.
Definition:
foldMany f = Stream.parseMany (Parser.fromFold f)Example, empty stream:
f = Fold.take 2 Fold.sumfmany = Stream.fold Fold.toList . Stream.foldMany ffmany $ Stream.fromList [][]
Example, last fold empty:
fmany $ Stream.fromList [1..4][3,7]
Example, last fold non-empty:
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.
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:
chunksOf n = Stream.foldMany (Array.writeN n)Pre-release
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:
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:
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:
splitOn' (== '.') "."["",""]
splitOn' (== '.') ".a"["","a"]
splitOn' (== '.') "a."["a",""]
splitOn' (== '.') "a..b"["a","","b"]
splitOn is an inverse of intercalating single element:
Stream.intercalate (Stream.fromPure '.') Unfold.fromList . Stream.splitOn (== '.') Fold.toList === idAssuming the input stream does not contain the separator:
Stream.splitOn (== '.') Fold.toList . Stream.intercalate (Stream.fromPure '.') Unfold.fromList === idSplit 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.
Run the action m b before the stream yields its first element.
Same as the following but more efficient due to fusion:
before action xs = Stream.nilM action <> xsbefore action xs = Stream.concatMap (const xs) (Stream.fromEffect action)
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
Lift the inner monad m of Stream m a to t m where t is a monad
transformer.
Evaluate the inner monad of a stream as ReaderT.
Evaluate the inner monad of a stream as StateT and emit the resulting state and value pair after each step.
A stream that terminates without producing any output, but produces a side effect.
Stream.fold Fold.toList (Stream.nilM (print "nil"))"nil"[]
Pre-release
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:
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.
Insert a side effect before consuming an element of a stream except the first one.
input = Stream.fromList "hello"Stream.fold Fold.drain $ Stream.trace putChar $ Stream.intersperseM_ (putChar '.') inputh.e.l.l.o
Pre-release
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.
s1 = Stream.fromList [1,2]s2 = Stream.fromList [3,4]Stream.fold Fold.toList $ s1 `Stream.append` s2[1,2,3,4]
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.
Apply a Parser repeatedly on a stream and emit the parsed values in the
output stream.
Example:
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.
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)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.
Transform the inner monad of a stream using a natural transformation.
Example, generalize the inner monad from Identity to any other:
generalizeInner = Stream.morphInner (return . runIdentity)Also known as hoist.
Remove the either wrapper and flatten both lefts and as well as rights in the output stream.
catEithers = fmap (either id id)Pre-release
Use a filtering fold on a stream.
scanMaybe f = Stream.catMaybes . Stream.postscan fDefinition:
crossWith f m1 m2 = fmap f m1 `Stream.crossApply` m2Note that the second stream is evaluated multiple times.
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
Like bracketIO but can use 3 separate cleanup actions depending on the mode of termination:
When the stream stops normally
When the stream is garbage collected
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
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.
finallyIO release = Stream.bracketIO (return ()) (const release)See also finallyUnsafe
Inhibits stream fusion
Like fold but also returns the remaining stream. The resulting stream
would be Stream.nil if the stream finished before the fold.
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.
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.
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.
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.
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.
A monad that can perform concurrent or parallel IO operations. Streams that can be composed concurrently require the underlying monad to be MonadAsync.
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.
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.
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.
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.
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.
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.
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.
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.
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.
Specify when the Channel should stop.
FirstStopsStop when the first stream ends.
AllStopStop when all the streams end.
AnyStopsStop when any one stream ends.
Specify when the Channel should stop.
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.
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.
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.
Print debug information about the Channel when the stream ends.
Stream combinators using a concurrent channel.
Evaluate a stream as a whole concurrently with respect to the consumer of the stream.
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:
parUnfoldrM step = Stream.parEval id . Stream.unfoldrM stepparIterateM step = Stream.parEval id . Stream.iterateM step
Generate a stream by evaluating multiple actions concurrently.
Definition:
parRepeatM cfg = Stream.parSequence cfg . Stream.repeatGenerate a stream by repeatedly executing a monadic action forever.
Generate a stream by concurrently performing a monadic action n times.
Definition:
parReplicateM cfg n = Stream.parSequence cfg . Stream.replicate nExample, parReplicateM in the following example executes all the replicated actions concurrently, thus taking only 1 second:
Stream.fold Fold.drain $ Stream.parReplicateM id 10 $ delay 1...
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 actions on a stream such that the mapped actions are evaluated concurrently with each other.
Definition:
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:
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]
Definition:
parSequence modifier = Stream.parMapM modifier idUseful idioms:
parFromListM = Stream.parSequence id . Stream.fromListparFromFoldableM = Stream.parSequence id . StreamK.toStream . StreamK.fromFoldable
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.
Evaluates the streams being zipped in separate threads than the consumer. The zip function is evaluated in the consumer thread.
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.
parZipWith cfg f = Stream.parZipWithM cfg (\a b -> return $ f a b)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)]
Like mergeByM but evaluates both the streams concurrently.
Definition:
parMergeByM cfg f m1 m2 = Stream.mergeByM f (Stream.parEval cfg m1) (Stream.parEval cfg m2)Like mergeBy but evaluates both the streams concurrently.
Definition:
parMergeBy cfg f = Stream.parMergeByM cfg (\a b -> return $ f a b)Shares a single channel across many streams.
Like parConcat but works on a list of streams.
parList modifier = Stream.parConcat modifier . Stream.fromListApply an argument stream to a function stream concurrently. Uses a shared channel for all individual applications within a stream application.
Shares a single channel across many streams.
Evaluate the streams in the input stream concurrently and combine them.
parConcat modifier = Stream.parConcatMap modifier idMap 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:
parConcatMap modifier f stream = Stream.parConcat modifier $ fmap f streamExamples:
f cfg xs = Stream.fold Fold.toList $ Stream.parConcatMap cfg id $ Stream.fromList xsThe following streams finish in 4 seconds:
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:
f (Stream.maxThreads 2) [stream1, stream2, stream3]...[2,1,4]
When used with a Single thread it behaves like serial concatMap:
f (Stream.maxThreads 1) [stream1, stream2, stream3]...[4,2,1]
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:
f cfg xs = Stream.fold Fold.toList $ Stream.parConcatMap (Stream.interleaved True . cfg) id $ Stream.fromList xsstream1 = Stream.fromList [1,2,3]stream2 = Stream.fromList [4,5,6]f (Stream.maxThreads 1) [stream1, stream2][1,4,2,5,3,6]
Same as concatIterate but concurrent.
Pre-release
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:
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
Intersperse a monadic action into the input stream after every n
seconds.
Definition:
interject n f xs = Stream.parListEagerFst [xs, Stream.periodic f n]Example:
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
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.
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.
Group the input stream into windows of n second each and then fold each
group using the provided fold function.
twoPerSec = Stream.parEval (Stream.constRate 2) $ Stream.enumerateFrom 1intervals = Stream.intervalsOf 1 Fold.toList twoPerSecStream.fold Fold.toList $ Stream.take 2 intervals[...,...]
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.
sampleIntervalEnd n = Stream.catMaybes . Stream.intervalsOf n Fold.latestLike sampleInterval but samples at the beginning of the time window.
sampleIntervalStart n = Stream.catMaybes . Stream.intervalsOf n Fold.oneSample 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.
Like sampleBurstEnd but samples the event at the beginning of the burst instead of at the end of it.
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_
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
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.
finally action xs = Stream.bracket (return ()) (const action) (const xs)See also finally_
Inhibits stream fusion
Deprecated. Please use parTapCount instead.
Same as parTapCount. Deprecated.