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-core-0.2.2Haskell2010

Streamly.Data.StreamK

Streams represented as chains of functions calls using Continuation Passing Style (CPS), suitable for dynamically composing potentially large number of streams.

Unlike the statically fused operations in Streamly.Data.Stream, StreamK operations are less efficient, involving a function call overhead for each element, but they exhibit linear O(n) time complexity wrt to the number of stream compositions. Therefore, they are suitable for dynamically composing streams e.g. appending potentially infinite streams in recursive loops. While fused streams can be used to efficiently process elements as small as a single byte, CPS streams are typically used on bigger chunks of data to avoid the larger overhead per element. For more details See the Stream vs StreamK section in the Streamly.Data.Stream module.

In addition to the combinators in this module, you can use operations from Streamly.Data.Stream for StreamK as well by converting StreamK to Stream (toStream), and vice-versa (fromStream). Please refer to Streamly.Internal.Data.StreamK for more functions that have not yet been released.

For documentation see the corresponding combinators in Streamly.Data.Stream. Documentation has been omitted in this module unless there is a difference worth mentioning or if the combinator does not exist in Streamly.Data.Stream.

  • 1 type
  • 32 values

Setup

0 declarations

To execute the code examples provided in this module in ghci, please run the following commands first.

Example4 expressions
:mimport Control.Concurrent (threadDelay)import Data.Function (fix, (&))import Data.Semigroup (cycle1)
Example1 expression
effect n = print n >> return n
Example6 expressions
import Streamly.Data.StreamK (StreamK)import qualified Streamly.Data.Fold as Foldimport qualified Streamly.Data.Parser as Parserimport qualified Streamly.Data.Stream as Streamimport qualified Streamly.Data.StreamK as StreamKimport qualified Streamly.FileSystem.Dir as Dir

For APIs that have not been released yet.

Example2 expressions
import qualified Streamly.Internal.Data.StreamK as StreamKimport qualified Streamly.Internal.FileSystem.Dir as Dir

Overview

0 declarations

Continuation passing style (CPS) stream implementation. The K in StreamK stands for Kontinuation.

StreamK can be constructed like lists, except that they use nil instead of '[]' and cons instead of :.

cons adds a pure value at the head of the stream:

Example2 expressions
import Streamly.Data.StreamK (StreamK, cons, consM, nil)stream = 1 `cons` 2 `cons` nil :: StreamK IO Int

You can use operations from Streamly.Data.Stream for StreamK as well by converting StreamK to Stream (toStream), and vice-versa (fromStream).

Example1 expression
Stream.fold Fold.toList $ StreamK.toStream stream -- IO [Int][1,2]

consM adds an effect at the head of the stream:

Example2 expressions
stream = effect 1 `consM` effect 2 `consM` nilStream.fold Fold.toList $ StreamK.toStream stream12[1,2]

Type

1 declaration
newtypenewtype StreamK (m :: Type -> Type) a
#
Instances10Functor, Foldable, Traversable, IsList, Read, Show, …

Construction

0 declarations

Primitives

Primitives to construct a stream from pure values or monadic actions. All other stream construction and generation combinators described later can be expressed in terms of these primitives. However, the special versions provided in this module can be much more efficient in some cases. Users can create custom combinators using these primitives.

valuenil :: StreamK m a
#

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

Example1 expression
Stream.fold Fold.toList (StreamK.toStream StreamK.nil)[]
valuenilM :: Applicative m => m b -> StreamK m a
#

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

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

Pre-release

valuecons :: a -> StreamK m a -> StreamK m a
#

A right associative prepend operation to add a pure value at the head of an existing stream::

Example2 expressions
s = 1 `StreamK.cons` 2 `StreamK.cons` 3 `StreamK.cons` StreamK.nilStream.fold Fold.toList (StreamK.toStream s)[1,2,3]

It can be used efficiently with foldr:

Example1 expression
fromFoldable = Prelude.foldr StreamK.cons StreamK.nil

Same as the following but more efficient:

Example1 expression
cons x xs = return x `StreamK.consM` xs
valueconsM :: Monad m => m a -> StreamK m a -> StreamK m a
#

A right associative prepend operation to add an effectful value at the head of an existing stream::

Example2 expressions
s = putStrLn "hello" `StreamK.consM` putStrLn "world" `StreamK.consM` StreamK.nilStream.fold Fold.drain (StreamK.toStream s)helloworld

It can be used efficiently with foldr:

Example1 expression
fromFoldableM = Prelude.foldr StreamK.consM StreamK.nil

Same as the following but more efficient:

Example1 expression
consM x xs = StreamK.fromEffect x `StreamK.append` xs

From Values

From Stream

Please note that Stream type does not observe any exceptions from the consumer of the stream whereas StreamK does.

valuefromStream :: Monad m => Stream m a -> StreamK m a
#

Convert a fused Stream to StreamK.

For example:

Example3 expressions
s1 = StreamK.fromStream $ Stream.fromList [1,2]s2 = StreamK.fromStream $ Stream.fromList [3,4]Stream.fold Fold.toList $ StreamK.toStream $ s1 `StreamK.append` s2[1,2,3,4]

From Containers

valuefromFoldable :: Foldable f => f a -> StreamK m a
#
Example1 expression
fromFoldable = Prelude.foldr StreamK.cons StreamK.nil

Construct a stream from a Foldable containing pure values:

Elimination

0 declarations

Primitives

Parsing

Transformation

3 declarations

Combining Two Streams

0 declarations

Unlike the operations in Streamly.Data.Stream, these operations can be used to dynamically compose large number of streams e.g. using the concatMapWith and mergeMapWith operations. They have a linear O(n) time complexity wrt to the number of streams being composed.

Appending

Interleaving

valueinterleave :: StreamK m a -> StreamK m a -> StreamK m a
#

Note: When joining many streams in a left associative manner earlier streams will get exponential priority than the ones joining later. Because of exponentially high weighting of left streams it can be used with concatMapWith even on a large number of streams.

Merging

Zipping

Cross Product

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

Definition:

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

Note that the second stream is evaluated multiple times.

Stream of streams

3 declarations

Some useful idioms:

Example3 expressions
concatFoldableWith f = Prelude.foldr f StreamK.nilconcatMapFoldableWith f g = Prelude.foldr (f . g) StreamK.nilconcatForFoldableWith f xs g = Prelude.foldr (f . g) StreamK.nil xs
valuemergeMapWith
  1. :: StreamK m b -> StreamK m b -> StreamK m b
  2. -> a -> StreamK m b
  3. -> StreamK m a
  4. -> StreamK m b
#

Combine streams in pairs using a binary combinator, the resulting streams are then combined again in pairs recursively until we get to a single combined stream. The composition would thus form a binary tree.

For example, you can sort a stream using merge sort like this:

Example4 expressions
s = StreamK.fromStream $ Stream.fromList [5,1,7,9,2]generate = StreamK.fromPurecombine = StreamK.mergeBy compareStream.fold Fold.toList $ StreamK.toStream $ StreamK.mergeMapWith combine generate s[1,2,5,7,9]

Note that if the stream length is not a power of 2, the binary tree composed by mergeMapWith would not be balanced, which may or may not be important depending on what you are trying to achieve.

Caution: the stream of streams must be finite

Pre-release

Buffered Operations

2 declarations
valuesortBy :: Monad m => (a -> a -> Ordering) -> StreamK m a -> StreamK m a
#

Sort the input stream using a supplied comparison function.

Sorting can be achieved by simply:

Example1 expression
sortBy cmp = StreamK.mergeMapWith (StreamK.mergeBy cmp) StreamK.fromPure

However, this combinator uses a parser to first split the input stream into down and up sorted segments and then merges them to optimize sorting when pre-sorted sequences exist in the input stream.

O(n) space

Exceptions

1 declaration

Please note that Stream type does not observe any exceptions from the consumer of the stream whereas StreamK does.

valuehandle
  1. :: (MonadCatch m, Exception e)
  2. => e -> m (StreamK m a)
  3. -> StreamK m a
  4. -> StreamK m a
#

Like Streamly.Data.Stream.Streamly.Data.Stream.handle but with one significant difference, this function observes exceptions from the consumer of the stream as well.

You can also convert StreamK to Stream and use exception handling from Stream module:

Example1 expression
handle f s = StreamK.fromStream $ Stream.handle (\e -> StreamK.toStream (f e)) (StreamK.toStream s)

Resource Management

1 declaration

Please note that Stream type does not observe any exceptions from the consumer of the stream whereas StreamK does.

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

Like Streamly.Data.Stream.Streamly.Data.Stream.bracketIO but with one significant difference, this function observes exceptions from the consumer of the stream as well. Therefore, it cleans up the resource promptly when the consumer encounters an exception.

You can also convert StreamK to Stream and use resource handling from Stream module:

Example1 expression
bracketIO bef aft bet = StreamK.fromStream $ Stream.bracketIO bef aft (StreamK.toStream . bet)