Streams represented as state machines, that fuse together when composed
statically, eliminating function calls or intermediate constructor
allocations - generating tight, efficient loops. Suitable for high
performance looping operations.
If you need to call these operations recursively in a loop (i.e. composed
dynamically) then it is recommended to use the continuation passing style
(CPS) stream operations from the Streamly.Data.StreamK module. Stream
and StreamK types are interconvertible. See more details in the
documentation below regarding Stream vs StreamK.
>>> hSetBuffering stdout LineBuffering>>> effect n = print n >> return n
Example8 expressions
>>> import Streamly.Data.Stream (Stream)>>> import qualified Streamly.Data.Array as Array>>> import qualified Streamly.Data.Fold as Fold>>> import qualified Streamly.Data.Stream as Stream>>> import qualified Streamly.Data.StreamK as StreamK>>> import qualified Streamly.Data.Unfold as Unfold>>> import qualified Streamly.Data.Parser as Parser>>> import qualified Streamly.FileSystem.Dir as Dir
For APIs that have not been released yet.
Example5 expressions
>>> import qualified Streamly.Internal.Data.Fold as Fold>>> import qualified Streamly.Internal.Data.Parser as Parser>>> import qualified Streamly.Internal.Data.Stream as Stream>>> import qualified Streamly.Internal.Data.Unfold as Unfold>>> import qualified Streamly.Internal.FileSystem.Dir as Dir
Overview
0 declarations
Streamly is a framework for modular data flow based programming and
declarative concurrency. Powerful stream fusion framework in streamly
allows high performance combinatorial programming even when using byte level
streams. Streamly API is similar to Haskell lists.
Console Echo Example
In the following example, repeatM generates an infinite stream of String
by repeatedly performing the getLine IO action. mapM then applies
putStrLn on each element in the stream converting it to stream of ().
Finally, drain folds the stream to IO discarding the () values, thus
producing only effects.
This is a console echo program. It is an example of a declarative loop
written using streaming combinators. Compare it with an imperative while
loop.
Hopefully, this gives you an idea how we can program declaratively by
representing loops using streams. In this module, you can find all
Data.List like functions and many more powerful combinators to perform
common programming tasks.
Stream Fusion
The fused Stream type in this module employs stream fusion for C-like
performance when looping over data. It represents the stream as a state
machine using an explicit state, and a step function working on the state. A
typical stream operation consumes elements from the previous state machine
in a stream pipeline, transforms the elements and yields new values for the
next stage to consume. The stream operations are modular and represent a
single task, they have no knowledge of previous or next operation on the
elements.
A typical stream pipeline consists of a stream producer, several stream
transformation operations and a stream consumer. All these operations taken
together form a closed loop processing the stream elements. Elements are
transferred between stages using a boxed data constructor. However, all the
stages of the pipeline are fused together by GHC, eliminating the boxing of
intermediate constructors, and thus forming a tight C like loop without any
boxed data being used in the loop.
Stream fusion works effectively when:
the stream pipeline is composed statically (known at compile time)
all the operations forming the loop are inlined
the loop is not recursively defined, recursion breaks inlining
If these conditions cannot be met, the CPS style stream type StreamK may
turn out to be a better choice than the fused stream type Stream.
Stream vs StreamK
The fused stream model avoids constructor allocations and function call
overheads. However, the stream is represented as a state machine, and to
generate stream elements it has to navigate the decision tree of the state
machine. Moreover, the state machine is cranked for each element in the
stream. This performs extremely well when the number of states are limited.
The state machine starts getting expensive as the number of states increase.
For example, generating a stream from a list requires a single state and is
very efficient, even if it has millions of elements. However, using cons
to construct a million element stream would be a disaster.
A typical worst case scenario for fused stream model is a large number of
cons or append operations. A few static cons or append operations
are very fast and much faster than a CPS style stream because CPS involves a
function call for each element whereas fused stream involves a few
conditional branches in the state machine. However, constructing a large
stream using cons introduces as many states in the state machine as the
number of elements. If we compose cons as a balanced binary tree it will
take n * log n time to navigate the tree, and n * n if it is a right
associative composition.
Operations like cons or append; are typically recursively called to
construct a lazy infinite stream. For such use cases the CPS style StreamK
should be used. CPS streams do not have a state machine that needs to be
cranked for each element, past state has no effect on the future element
processing. However, CPS incurs a function call overhead for each element
processed, the overhead could be large compared to a fused state machine
even if it has many states. However, because of its linear performance
characterstics, after a certain threshold of stream compositions the CPS
stream would perform much better than the quadratic fused stream operations.
As a general guideline, you need to use StreamK when you have to use
cons, append or other operations having quadratic complexity at a large
scale. Typically, in such cases you need to compose the stream recursively,
by calling an operation in a loop. The decision to compose the stream is
taken at run time rather than statically at compile time.
Typically you would compose a StreamK of chunks of data so that the
StreamK overhead is not high, and then process the chunks using Stream by
using statically fused stream pipeline operations on the chunks.
typeItem (StreamIdentitya) = aDefined in streamly-core-0.2.2 · Streamly.Internal.Data.Stream.Type
Construction
0 declarations
Functions ending in the general shape b -> Stream m a.
Useful Idioms:
Example5 expressions
>>> fromIndices f = fmap f $ Stream.enumerateFrom 0>>> fromIndicesM f = Stream.mapM f $ Stream.enumerateFrom 0>>> fromListM = Stream.sequence . Stream.fromList>>> fromFoldable = StreamK.toStream . StreamK.fromFoldable>>> fromFoldableM = Stream.sequence . fromFoldable
Primitives
A fused Stream is never constructed using these primitives, they are
typically generated by converting containers like list into streams, or
generated using custom functions provided in this module. The cons
primitive in this module has a rare use in fusing a small number of
elements. On the other hand, it is common to construct StreamK stream
using the StreamK.cons primitive.
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]
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]
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]
From Values
Generate a monadic stream from a seed value or values.
However, this is not particularly efficient.
The Enumerable type class provides corresponding functions that
generate a stream instead of a list, efficiently.
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 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.
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.
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.
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.
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.
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.nil>>> Just ("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):}
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 f>>> fold f = Stream.parse (Parser.fromFold f)
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:
Note: parse p is not the same as head . parseMany p on an empty stream.
Lazy Right Folds
Consuming a stream to build a right associated expression, suitable
for lazy evaluation. Evaluation of the input happens when the output of
the fold is evaluated, the fold output is a lazy thunk.
This is suitable for stream transformation operations, for example,
operations like mapping a function over the stream.
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 xs>>> Stream.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.
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.
Specific Folds
Usually you can use the folds in Streamly.Data.Fold. However, some
folds that may be commonly used or may have an edge in performance in
some cases are provided here.
Useful idioms:
Example7 expressions
>>> foldlM' f a = Stream.fold (Fold.foldlM' f a)>>> foldl1 f = Stream.fold (Fold.foldl1' f)>>> foldl' f a = Stream.fold (Fold.foldl' f a)>>> drain = Stream.fold Fold.drain>>> mapM_ f = Stream.fold (Fold.drainMapM f)>>> length = Stream.fold Fold.length>>> head = Stream.fold Fold.one
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.
Mapping
5 declarations
Stateless one-to-one transformations. Use fmap for mapping a pure
function on a stream.
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-----
>>> scanl' f z = Stream.scan (Fold.foldl' f z)>>> scanlM' f z = Stream.scan (Fold.foldlM' f z)>>> postscanl' f z = Stream.postscan (Fold.foldl' f z)>>> postscanlM' f z = Stream.postscan (Fold.foldlM' f z)>>> scanl1' f = Stream.catMaybes . Stream.scan (Fold.foldl1' f)>>> scanl1M' f = Stream.catMaybes . Stream.scan (Fold.foldlM1' f)
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".
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)
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
Stateful Filters
scanMaybe is the most general stateful filtering operation. The
filtering folds (folds returning a Maybe type) in
Streamly.Internal.Data.Fold can be used along with scanMaybe to
perform stateful filtering operations in general.
Useful idioms:
Example4 expressions
>>> deleteBy cmp x = Stream.scanMaybe (Fold.deleteBy cmp x)>>> findIndices p = Stream.scanMaybe (Fold.findIndices p)>>> elemIndices a = findIndices (== a)>>> uniq = Stream.scanMaybe (Fold.uniqBy (==))
Note that these operations are suitable for statically fusing a few
streams, they have a quadratic O(n^2) time complexity wrt to the number
of streams. If you want to compose many streams dynamically using binary
combining operations see the corresponding operations in
Streamly.Data.StreamK.
When fusing more than two streams it is more efficient if the binary
operations are composed as a balanced tree rather than a right
associative or left associative one e.g.:
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.
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.
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.
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]
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.
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.
Stream operations like map and filter represent loop processing in
imperative programming terms. Similarly, the imperative concept of
nested loops are represented by streams of streams. The concatMap
operation represents nested looping.
A concatMap operation loops over the input stream and then for each
element of the input stream generates another stream and then loops over
that inner stream as well producing effects and generating a single
output stream.
One dimension loops are just a special case of nested loops. For
example, concatMap can degenerate to a simple map operation:
map f m = S.concatMap (\x -> S.fromPure (f x)) m
Similarly, concatMap can perform filtering by mapping an element to a
nil stream:
filter p m = S.concatMap (\x -> if p x then S.fromPure x else S.nil) m
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.
>>> splitWithSuffix p f = Stream.foldMany (Fold.takeEndBy p f)>>> splitOnSuffix p f = Stream.foldMany (Fold.takeEndBy_ p f)>>> groupsBy eq f = Stream.parseMany (Parser.groupBy eq f)>>> groupsByRolling eq f = Stream.parseMany (Parser.groupByRolling eq f)>>> groupsOf n f = Stream.foldMany (Fold.take n f)
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:
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.
Buffered Operations
1 declaration
Operations that require buffering of the stream.
Reverse is essentially a left fold followed by an unfold.
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.
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)
Exceptions
2 declarations
Note that the stream exception handling routines catch and handle
exceptions only in the stream generation steps and not in the consumer
of the stream. For example, if we are folding or parsing a stream - any
exceptions in the fold or parse steps won't be observed by the stream
exception handlers. Exceptions in the fold or parse steps can be handled
using the fold or parse exception handling routines. You can wrap the
stream elimination function in the monad exception handler to observe
exceptions in the stream as well as the consumer.
Most of these combinators inhibit stream fusion, therefore, when
possible, they should be called in an outer loop to mitigate the cost.
For example, instead of calling them on a stream of chars call them on a
stream of arrays before flattening it to a stream of chars.
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
Resource Management
5 declarations
bracket is the most general resource management operation, all other
operations can be expressed using it. These functions have IO suffix
because the allocation and cleanup functions are IO actions. For
generalized allocation and cleanup functions, see the functions without
the IO suffix in the "streamly" package.
Note that these operations bracket the stream generation only, they do
not cover the stream consumer. This means if an exception occurs in
the consumer of the stream (e.g. in a fold or parse step) then the
exception won't be observed by the stream resource handlers, in that
case the resource cleanup handler runs when the stream is garbage
collected.
Monad level resource management can always be used around the stream
elimination functions, such a function can observe exceptions in both
the stream and its consumer.
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.
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.
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.