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.Fold.Prelude

All Fold related combinators including the streamly-core Streamly.Data.Fold module, concurrency, unordered container operations.

  • 3 types
  • 104 values
  • Packagestreamly-0.10.1
  • Exports107
  • LanguageHaskell2010
  • LicenceBSD-3-Clause
  • SourcePrelude.hs

Streamly.Data.Fold

101 declarations

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

datadata Fold (m :: Type -> Type) a b
#

The type Fold m a b represents a consumer of an input stream of values of type a and returning a final value of type b in Monad m. The constructor of a fold is Fold step initial extract final.

The fold uses an internal state of type s. The initial value of the state s is created by initial. This function is called once and only once before the fold starts consuming input. Any resource allocation can be done in this function.

The step function is called on each input, it consumes an input and returns the next intermediate state (see Step) or the final result b if the fold terminates.

If the fold is used as a scan, the extract function is used by the scan driver to map the current state s of the fold to the fold result. Thus extract can be called multiple times. In some folds, where scanning does not make sense, this function is left unimplemented; such folds cannot be used as scans.

Before a fold terminates, final is called once and only once (unless the fold terminated in initial itself). Any resources allocated by initial can be released in final. In folds that do not require any cleanup extract and final are typically the same.

When implementing fold combinators, care should be taken to cleanup any state of the argument folds held by the fold by calling the respective final at all exit points of the fold. Also, final should not be called more than once. Note that if a fold terminates by Done constructor, there is no state to cleanup.

NOTE: The constructor is not yet released, smart constructors are provided to create folds.

Instances2Functor, Applicative
  • Functor m => Functor (Fold m a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Type

    Maps a function on the output of the fold (the type b).

  • Monad m => Applicative (Fold m a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Type

    Applicative form of splitWith. Split the input serially over two folds. Note that this fuses but performance degrades quadratically with respect to the number of compositions. It should be good to use for less than 8 compositions.

valuetoList :: Monad m => Fold m a [a]
#

Folds the input stream to a list.

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

Example1 expression
toList = Fold.foldr' (:) []
valuelength :: Monad m => Fold m a Int
#

Determine the length of the input stream.

Definition:

Example2 expressions
length = Fold.lengthGenericlength = fmap getSum $ Fold.foldMap (Sum . const  1)
newtypenewtype Tee (m :: Type -> Type) a b
#

Tee is a newtype wrapper over the Fold type providing distributing Applicative, Semigroup, Monoid, Num, Floating and Fractional instances.

The input received by the composed Tee is replicated and distributed to the constituent folds of the Tee.

For example, to compute the average of numbers in a stream without going through the stream twice:

Example2 expressions
avg = (/) <$> (Tee Fold.sum) <*> (Tee $ fmap fromIntegral Fold.length)Stream.fold (unTee avg) $ Stream.fromList [1.0..100.0]50.5

Similarly, the Semigroup and Monoid instances of Tee distribute the input to both the folds and combine the outputs using Monoid or Semigroup instances of the output types:

Example3 expressions
import Data.Monoid (Sum(..))t = Tee Fold.one <> Tee Fold.latestStream.fold (unTee t) (fmap Sum $ Stream.enumerateFromTo 1.0 100.0)Just (Sum {getSum = 101.0})

The Num, Floating, and Fractional instances work in the same way.

Constructors

Instances7Functor, Applicative, Floating, Fractional, Num, Semigroup, …
  • Functor m => Functor (Tee m a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Tee
  • Monad m => Applicative (Tee m a)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Tee

    <*> distributes the input to both the argument Tees and combines their outputs using function application.

  • (Monad m, Floating b) => Floating (Tee m a b)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Tee

    Binary Floating operations distribute the input to both the argument Tees and combine their outputs using the Floating instance of the output type.

  • (Monad m, Fractional b) => Fractional (Tee m a b)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Tee

    Binary Fractional operations distribute the input to both the argument Tees and combine their outputs using the Fractional instance of the output type.

  • (Monad m, Num b) => Num (Tee m a b)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Tee

    Binary Num operations distribute the input to both the argument Tees and combine their outputs using the Num instance of the output type.

  • (Semigroup b, Monad m) => Semigroup (Tee m a b)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Tee

    <> distributes the input to both the argument Tees and combines their outputs using the Semigroup instance of the output type.

  • (Monoid b, Monad m) => Monoid (Tee m a b)Defined in streamly-core-0.2.2 · Streamly.Internal.Data.Fold.Tee

    <> distributes the input to both the argument Tees and combines their outputs using the Monoid instance of the output type.

valuefoldr :: Monad m => (a -> b -> b) -> b -> Fold m a b
#

Deprecated. Please use foldr' instead.

valuefoldl' :: Monad m => (b -> a -> b) -> b -> Fold m a b
#

Make a fold from a left fold style pure step function and initial value of the accumulator.

If your Fold returns only Partial (i.e. never returns a Done) then you can use foldl'* constructors.

A fold with an extract function can be expressed using fmap:

mkfoldlx :: Monad m => (s -> a -> s) -> s -> (s -> b) -> Fold m a b
mkfoldlx step initial extract = fmap extract (foldl' step initial)
valuefoldl1' :: Monad m => (a -> a -> a) -> Fold m a (Maybe a)
#

Make a strict left fold, for non-empty streams, using first element as the starting value. Returns Nothing if the stream is empty.

Pre-release

valuefoldlM' :: Monad m => (b -> a -> m b) -> m b -> Fold m a b
#

Make a fold from a left fold style monadic step function and initial value of the accumulator.

A fold with an extract function can be expressed using rmapM:

mkFoldlxM :: Functor m => (s -> a -> m s) -> m s -> (s -> m b) -> Fold m a b
mkFoldlxM step initial extract = rmapM extract (foldlM' step initial)
valuedrain :: Monad m => Fold m a ()
#

A fold that drains all its input, running the effects and discarding the results.

Example2 expressions
drain = Fold.drainMapM (const (return ()))drain = Fold.foldl' (\_ _ -> ()) ()
valuesum :: (Monad m, Num a) => Fold m a a
#

Determine the sum of all elements of a stream of numbers. Returns additive identity (0) when the stream is empty. Note that this is not numerically stable for floating point numbers.

Example1 expression
sum = Fold.cumulative Fold.windowSum

Same as following but numerically stable:

Example2 expressions
sum = Fold.foldl' (+) 0sum = fmap Data.Monoid.getSum $ Fold.foldMap Data.Monoid.Sum
valueproduct :: (Monad m, Num a, Eq a) => Fold m a a
#

Determine the product of all elements of a stream of numbers. Returns multiplicative identity (1) when the stream is empty. The fold terminates when it encounters (0) in its input.

Same as the following but terminates on multiplication by 0:

Example1 expression
product = fmap Data.Monoid.getProduct $ Fold.foldMap Data.Monoid.Product
valuemaximum :: (Monad m, Ord a) => Fold m a (Maybe a)
#

Determine the maximum element in a stream.

Definitions:

Example2 expressions
maximum = Fold.maximumBy comparemaximum = Fold.foldl1' max

Same as the following but without a default maximum. The Max Monoid uses the minBound as the default maximum:

Example1 expression
maximum = fmap Data.Semigroup.getMax $ Fold.foldMap Data.Semigroup.Max
valueminimum :: (Monad m, Ord a) => Fold m a (Maybe a)
#

Determine the minimum element in a stream using the supplied comparison function.

Definitions:

Example2 expressions
minimum = Fold.minimumBy compareminimum = Fold.foldl1' min

Same as the following but without a default minimum. The Min Monoid uses the maxBound as the default maximum:

Example1 expression
maximum = fmap Data.Semigroup.getMin $ Fold.foldMap Data.Semigroup.Min
valuethe :: (Monad m, Eq a) => Fold m a (Maybe a)
#

Terminates with Nothing as soon as it finds an element different than the previous one, returns the element if the entire input consists of the same element.

valuehead :: Monad m => Fold m a (Maybe a)
#

Deprecated. Please use "one" instead

Extract the first element of the stream, if any.

Example1 expression
head = Fold.one
valuefindM :: Monad m => (a -> m Bool) -> Fold m a (Maybe a)
#

Returns the first element that satisfies the given predicate.

Pre-release

valuefind :: Monad m => (a -> Bool) -> Fold m a (Maybe a)
#

Returns the first element that satisfies the given predicate.

valuelookup :: (Eq a, Monad m) => a -> Fold m (a, b) (Maybe b)
#

In a stream of (key-value) pairs (a, b), return the value b of the first pair where the key equals the given value a.

Definition:

Example1 expression
lookup x = fmap snd <$> Fold.find ((== x) . fst)
valueelemIndex :: (Eq a, Monad m) => a -> Fold m a (Maybe Int)
#

Returns the first index where a given value is found in the stream.

Definition:

Example1 expression
elemIndex a = Fold.findIndex (== a)
valuenull :: Monad m => Fold m a Bool
#

Consume one element, return True if successful else return False. In other words, test if the input is empty or not.

WARNING! It consumes one element if the stream is not empty. If that is not what you want please use the eof parser instead.

Definition:

Example1 expression
null = fmap isJust Fold.one
valueelem :: (Eq a, Monad m) => a -> Fold m a Bool
#

Return True if the given element is present in the stream.

Definition:

Example1 expression
elem a = Fold.any (== a)
valuenotElem :: (Eq a, Monad m) => a -> Fold m a Bool
#

Returns True if the given element is not present in the stream.

Definition:

Example1 expression
notElem a = Fold.all (/= a)
valueall :: Monad m => (a -> Bool) -> Fold m a Bool
#

Returns True if all elements of the input satisfy the predicate.

Definition:

Example1 expression
all p = Fold.lmap p Fold.and

Example:

Example1 expression
Stream.fold (Fold.all (== 0)) $ Stream.fromList [1,0,1]False
valueany :: Monad m => (a -> Bool) -> Fold m a Bool
#

Returns True if any element of the input satisfies the predicate.

Definition:

Example1 expression
any p = Fold.lmap p Fold.or

Example:

Example1 expression
Stream.fold (Fold.any (== 0)) $ Stream.fromList [1,0,1]True
valuesequence :: Monad m => Fold m a (m b) -> Fold m a b
#

Deprecated. Use "rmapM id" instead

Flatten the monadic output of a fold to pure output.

valuemapM :: Monad m => (b -> m c) -> Fold m a b -> Fold m a c
#

Deprecated. Use rmapM instead

Map a monadic function on the output of a fold.

valuescan :: Monad m => Fold m a b -> Fold m b c -> Fold m a c
#

Scan the input of a Fold to change it in a stateful manner using another Fold. The scan stops as soon as the fold terminates.

Pre-release

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

Postscan the input of a Fold to change it in a stateful manner using another Fold.

postscan scanner collector

Pre-release

valuedeleteBy :: Monad m => (a -> a -> Bool) -> a -> Fold m a (Maybe a)
#

Returns the latest element omitting the first occurrence that satisfies the given equality predicate.

Example:

Example2 expressions
input = Stream.fromList [1,3,3,5]Stream.fold Fold.toList $ Stream.scanMaybe (Fold.deleteBy (==) 3) input[1,3,5]
valuefilter :: Monad m => (a -> Bool) -> Fold m a r -> Fold m a r
#

Include only those elements that pass a predicate.

Example1 expression
Stream.fold (Fold.filter (> 5) Fold.sum) $ Stream.fromList [1..10]40
Example3 expressions
filter p = Fold.scanMaybe (Fold.filtering p)filter p = Fold.filterM (return . p)filter p = Fold.mapMaybe (\x -> if p x then Just x else Nothing)
valuefilterM :: Monad m => (a -> m Bool) -> Fold m a r -> Fold m a r
#

Like filter but with a monadic predicate.

Example2 expressions
f p x = p x >>= \r -> return $ if r then Just x else NothingfilterM p = Fold.mapMaybeM (f p)
valuetake :: Monad m => Int -> Fold m a b -> Fold m a b
#

Take at most n input elements and fold them using the supplied fold. A negative count is treated as 0.

Example1 expression
Stream.fold (Fold.take 2 Fold.toList) $ Stream.fromList [1..10][1,2]
valueelemIndices :: (Monad m, Eq a) => a -> Fold m a (Maybe Int)
#

Returns the index of the latest element if the element matches the given value.

Definition:

Example1 expression
elemIndices a = Fold.findIndices (== a)
valuemapMaybe :: Monad m => (a -> Maybe b) -> Fold m b r -> Fold m a r
#

mapMaybe f fold maps a Maybe returning function f on the input of the fold, filters out Nothing elements, and return the values extracted from Just.

Example2 expressions
mapMaybe f = Fold.lmap f . Fold.catMaybesmapMaybe f = Fold.mapMaybeM (return . f)
Example3 expressions
f x = if even x then Just x else Nothingfld = Fold.mapMaybe f Fold.toListStream.fold fld (Stream.enumerateFromTo 1 10)[2,4,6,8,10]
valueconcatMap :: Monad m => (b -> Fold m a c) -> Fold m a b -> Fold m a c
#

Map a Fold returning function on the result of a Fold and run the returned fold. This operation can be used to express data dependencies between fold operations.

Let's say the first element in the stream is a count of the following elements that we have to add, then:

Example4 expressions
import Data.Maybe (fromJust)count = fmap fromJust Fold.onetotal n = Fold.take n Fold.sumStream.fold (Fold.concatMap total count) $ Stream.fromList [10,9..1]45

This does not fuse completely, see refold for a fusible alternative.

Time: O(n^2) where n is the number of compositions.

See also: foldIterateM, refold

valuemany :: Monad m => Fold m a b -> Fold m b c -> Fold m a c
#

Collect zero or more applications of a fold. many first second applies the first fold repeatedly on the input stream and accumulates it's results using the second fold.

Example3 expressions
two = Fold.take 2 Fold.toListtwos = Fold.many two Fold.toListStream.fold twos $ Stream.fromList [1..10][[1,2],[3,4],[5,6],[7,8],[9,10]]

Stops when second fold stops.

See also: Data.Stream.concatMap, Data.Stream.foldMany

valuemconcat :: (Monad m, Monoid a) => Fold m a a
#

Monoid concat. Fold an input stream consisting of monoidal elements using mappend and mempty.

Definition:

Example1 expression
mconcat = Fold.sconcat mempty
Example2 expressions
monoids = fmap Data.Monoid.Sum $ Stream.enumerateFromTo 1 10Stream.fold Fold.mconcat monoidsSum {getSum = 55}
valuetoListRev :: Monad m => Fold m a [a]
#

Buffers the input stream to a list in the reverse order of the input.

Definition:

Example1 expression
toListRev = Fold.foldl' (flip (:)) []

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

valueuniqBy :: Monad m => (a -> a -> Bool) -> Fold m a (Maybe a)
#

Return the latest unique element using the supplied comparison function. Returns Nothing if the current element is same as the last element otherwise returns Just.

Example, strip duplicate path separators:

Example3 expressions
input = Stream.fromList "//a//b"f x y = x == '/' && y == '/'Stream.fold Fold.toList $ Stream.scanMaybe (Fold.uniqBy f) input"/a/b"

Space: O(1)

Pre-release

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

Modify a fold to receive a Maybe input, the Just values are unwrapped and sent to the original fold, Nothing values are discarded.

Example2 expressions
catMaybes = Fold.mapMaybe idcatMaybes = Fold.filter isJust . Fold.lmap fromJust
valuetakeEndBy_ :: Monad m => (a -> Bool) -> Fold m a b -> Fold m a b
#

Like takeEndBy but drops the element on which the predicate succeeds.

Example:

Example3 expressions
input = Stream.fromList "hello\nthere\n"line = Fold.takeEndBy_ (== '\n') Fold.toListStream.fold line input"hello"
Example1 expression
Stream.fold Fold.toList $ Stream.foldMany line input["hello","there"]
valuetakeEndBy :: Monad m => (a -> Bool) -> Fold m a b -> Fold m a b
#

Take the input, stop when the predicate succeeds taking the succeeding element as well.

Example:

Example3 expressions
input = Stream.fromList "hello\nthere\n"line = Fold.takeEndBy (== '\n') Fold.toListStream.fold line input"hello\n"
Example1 expression
Stream.fold Fold.toList $ Stream.foldMany line input["hello\n","there\n"]
valuegroupsOf :: Monad m => Int -> Fold m a b -> Fold m b c -> Fold m a c
#

groupsOf n split collect repeatedly applies the split fold to chunks of n items in the input stream and supplies the result to the collect fold.

Definition:

Example1 expression
groupsOf n split = Fold.many (Fold.take n split)

Example:

Example2 expressions
twos = Fold.groupsOf 2 Fold.toList Fold.toListStream.fold twos $ Stream.fromList [1..10][[1,2],[3,4],[5,6],[7,8],[9,10]]

Stops when collect stops.

valuemorphInner :: (forall x. m x -> n x) -> Fold m a b -> Fold n a b
#

Change the underlying monad of a fold. Also known as hoist.

Pre-release

valueunzip :: Monad m => Fold m a x -> Fold m b y -> Fold m (a, b) (x, y)
#

Send the elements of tuples in a stream of tuples through two different folds.


                          |-------Fold m a x--------|
---------stream of (a,b)--|                         |----m (x,y)
                          |-------Fold m b y--------|

Definition:

Example1 expression
unzip = Fold.unzipWith id

This is the consumer side dual of the producer side zip operation.

valuefoldMap :: (Monad m, Monoid b) => (a -> b) -> Fold m a b
#

Definition:

Example1 expression
foldMap f = Fold.lmap f Fold.mconcat

Make a fold from a pure function that folds the output of the function using mappend and mempty.

Example2 expressions
sum = Fold.foldMap Data.Monoid.SumStream.fold sum $ Stream.enumerateFromTo 1 10Sum {getSum = 55}
valueaddStream :: Monad m => Stream m a -> Fold m a b -> m (Fold m a b)
#

Append a stream to a fold to build the fold accumulator incrementally. We can repeatedly call addStream on the same fold to continue building the fold and finally use drive to finish the fold and extract the result. Also see the Streamly.Data.Fold.addOne operation which is a singleton version of addStream.

Definitions:

Example1 expression
addStream stream = Fold.drive stream . Fold.duplicate

Example, build a list incrementally:

Example1 expression
:{pure (Fold.toList :: Fold IO Int [Int])    >>= Fold.addOne 1    >>= Fold.addStream (Stream.enumerateFromTo 2 4)    >>= Fold.drive Stream.nil    >>= print:}[1,2,3,4]

This can be used as an O(n) list append compared to the O(n^2) ++ when used for incrementally building a list.

Example, build a stream incrementally:

Example1 expression
:{pure (Fold.toStream :: Fold IO Int (Stream Identity Int))    >>= Fold.addOne 1    >>= Fold.addStream (Stream.enumerateFromTo 2 4)    >>= Fold.drive Stream.nil    >>= print:}fromList [1,2,3,4]

This can be used as an O(n) stream append compared to the O(n^2) <> when used for incrementally building a stream.

Example, build an array incrementally:

Example1 expression
:{pure (Array.write :: Fold IO Int (Array Int))    >>= Fold.addOne 1    >>= Fold.addStream (Stream.enumerateFromTo 2 4)    >>= Fold.drive Stream.nil    >>= print:}fromList [1,2,3,4]

Example, build an array stream incrementally:

Example1 expression
:{let f :: Fold IO Int (Stream Identity (Array Int))    f = Fold.groupsOf 2 (Array.writeN 3) Fold.toStreamin pure f    >>= Fold.addOne 1    >>= Fold.addStream (Stream.enumerateFromTo 2 4)    >>= Fold.drive Stream.nil    >>= print:}fromList [fromList [1,2],fromList [3,4]]
valuedistribute :: Monad m => [Fold m a b] -> Fold m a [b]
#

Distribute one copy of the stream to each fold and collect the results in a container.


                |-------Fold m a b--------|
---stream m a---|                         |---m [b]
                |-------Fold m a b--------|
                |                         |
                           ...
Example1 expression
Stream.fold (Fold.distribute [Fold.sum, Fold.length]) (Stream.enumerateFromTo 1 5)[15,5]
Example1 expression
distribute = Prelude.foldr (Fold.teeWith (:)) (Fold.fromPure [])

This is the consumer side dual of the producer side sequence operation.

Stops when all the folds stop.

valuedrainMapM :: Monad m => (a -> m b) -> Fold m a ()
#

Definitions:

Example2 expressions
drainMapM f = Fold.lmapM f Fold.draindrainMapM f = Fold.foldMapM (void . f)

Drain all input after passing it through a monadic function. This is the dual of mapM_ on stream producers.

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

Drive a fold using the supplied Stream, reducing the resulting expression strictly at each step.

Definition:

Example1 expression
drive = flip Stream.fold

Example:

Example1 expression
Fold.drive (Stream.enumerateFromTo 1 100) Fold.sum5050
valuefoldMapM :: (Monad m, Monoid b) => (a -> m b) -> Fold m a b
#

Definition:

Example1 expression
foldMapM f = Fold.lmapM f Fold.mconcat

Make a fold from a monadic function that folds the output of the function using mappend and mempty.

Example2 expressions
sum = Fold.foldMapM (return . Data.Monoid.Sum)Stream.fold sum $ Stream.enumerateFromTo 1 10Sum {getSum = 55}
valueindex :: Monad m => Int -> Fold m a (Maybe a)
#

Return the element at the given index.

Definition:

Example1 expression
index = Fold.indexGeneric
valuelatest :: Monad m => Fold m a (Maybe a)
#

Returns the latest element of the input stream, if any.

Example2 expressions
latest = Fold.foldl1' (\_ x -> x)latest = fmap getLast $ Fold.foldMap (Last . Just)
valuemean :: (Monad m, Fractional a) => Fold m a a
#

Compute a numerically stable arithmetic mean of all elements in the input stream.

valueone :: Monad m => Fold m a (Maybe a)
#

Take one element from the stream and stop.

Definition:

Example1 expression
one = Fold.maybe Just

This is similar to the stream Stream.uncons operation.

valuepartition
  1. :: Monad m
  2. => Fold m b x
  3. -> Fold m c y
  4. -> Fold m (Either b c) (x, y)
#

Compose two folds such that the combined fold accepts a stream of Either and routes the Left values to the first fold and Right values to the second fold.

Definition:

Example1 expression
partition = Fold.partitionBy id
valuerollingHash :: (Monad m, Enum a) => Fold m a Int64
#

Compute an Int sized polynomial rolling hash of a stream.

Example1 expression
rollingHash = Fold.rollingHashWithSalt Fold.defaultSalt
valuesconcat :: (Monad m, Semigroup a) => a -> Fold m a a
#

Semigroup concat. Append the elements of an input stream to a provided starting value.

Definition:

Example1 expression
sconcat = Fold.foldl' (<>)
Example2 expressions
semigroups = fmap Data.Monoid.Sum $ Stream.enumerateFromTo 1 10Stream.fold (Fold.sconcat 10) semigroupsSum {getSum = 65}
valuestdDev :: (Monad m, Floating a) => Fold m a a
#

Deprecated. Use the streamly-statistics package instead

Compute a numerically stable (population) standard deviation over all elements in the input stream.

valuetee :: Monad m => Fold m a b -> Fold m a c -> Fold m a (b, c)
#

Distribute one copy of the stream to each fold and zip the results.

                |-------Fold m a b--------|
---stream m a---|                         |---m (b,c)
                |-------Fold m a c--------|

Definition:

Example1 expression
tee = Fold.teeWith (,)

Example:

Example2 expressions
t = Fold.tee Fold.sum Fold.lengthStream.fold t (Stream.enumerateFromTo 1.0 100.0)(5050.0,100)
valuetopBy
  1. :: (MonadIO m, Unbox a)
  2. => a -> a -> Ordering
  3. -> Int
  4. -> Fold m a (MutArray a)
#

Get the top n elements using the supplied comparison function.

To get bottom n elements instead:

Example1 expression
bottomBy cmp = Fold.topBy (flip cmp)

Example:

Example2 expressions
stream = Stream.fromList [2::Int,7,9,3,1,5,6,11,17]Stream.fold (Fold.topBy compare 3) stream >>= MutArray.toList[17,11,9]

Pre-release

valuevariance :: (Monad m, Fractional a) => Fold m a a
#

Deprecated. Use the streamly-statistics package instead

Compute a numerically stable (population) variance over all elements in the input stream.

valueclassify
  1. :: (Monad m, Ord k)
  2. => a -> k
  3. -> Fold m a b
  4. -> Fold m a (m (Map k b), Maybe (k, b))
#

Folds the values for each key using the supplied fold. When scanning, as soon as the fold is complete, its result is available in the second component of the tuple. The first component of the tuple is a snapshot of the in-progress folds.

Once the fold for a key is done, any future values of the key are ignored.

Definition:

Example1 expression
classify f fld = Fold.demux f (const fld)
valueclassifyIO
  1. :: (MonadIO m, Ord k)
  2. => a -> k
  3. -> Fold m a b
  4. -> Fold m a (m (Map k b), Maybe (k, b))
#

Same as classify except that it uses mutable IORef cells in the Map providing better performance. Be aware that if this is used as a scan, the values in the intermediate Maps would be mutable.

Definitions:

Example1 expression
classifyIO f fld = Fold.demuxIO f (const fld)
valuecountDistinct :: (Monad m, Ord a) => Fold m a Int
#

Count non-duplicate elements in the stream.

Definition:

Example2 expressions
countDistinct = fmap Set.size Fold.toSetcountDistinct = Fold.postscan Fold.nub $ Fold.catMaybes $ Fold.length

The memory used is proportional to the number of distinct elements in the stream, to guard against using too much memory use it as a scan and terminate if the count reaches more than a threshold.

Space: \mathcal{O}(n)

Pre-release

valuecountDistinctInt :: Monad m => Fold m Int Int
#

Like countDistinct but specialized to a stream of Int, for better performance.

Definition:

Example2 expressions
countDistinctInt = fmap IntSet.size Fold.toIntSetcountDistinctInt = Fold.postscan Fold.nubInt $ Fold.catMaybes $ Fold.length

Pre-release

valuedemux
  1. :: (Monad m, Ord k)
  2. => a -> k
  3. -> a -> m (Fold m a b)
  4. -> Fold m a (m (Map k b), Maybe (k, b))
#

demux getKey getFold: In a key value stream, fold values corresponding to each key using a key specific fold. getFold is invoked to generate a key specific fold when a key is encountered for the first time in the stream.

The first component of the output tuple is a key-value Map of in-progress folds. The fold returns the fold result as the second component of the output tuple whenever a fold terminates.

If a fold terminates, another instance of the fold is started upon receiving an input with that key, getFold is invoked again whenever the key is encountered again.

This can be used to scan a stream and collect the results from the scan output.

Since the fold generator function is monadic we can add folds dynamically. For example, we can maintain a Map of keys to folds in an IORef and lookup the fold from that corresponding to a key. This Map can be changed dynamically, folds for new keys can be added or folds for old keys can be deleted or modified.

Compare with classify, the fold in classify is a static fold.

Pre-release

valuedemuxIO
  1. :: (MonadIO m, Ord k)
  2. => a -> k
  3. -> a -> m (Fold m a b)
  4. -> Fold m a (m (Map k b), Maybe (k, b))
#

This is specialized version of demux that uses mutable IO cells as fold accumulators for better performance.

Keep in mind that the values in the returned Map may be changed by the ongoing fold if you are using those concurrently in another thread.

valuefrequency :: (Monad m, Ord a) => Fold m a (Map a Int)
#

Determine the frequency of each element in the stream.

You can just collect the keys of the resulting map to get the unique elements in the stream.

Definition:

Example1 expression
frequency = Fold.toMap id Fold.length
valuenub :: (Monad m, Ord a) => Fold m a (Maybe a)
#

Used as a scan. Returns Just for the first occurrence of an element, returns Nothing for any other occurrences.

Example:

Example2 expressions
stream = Stream.fromList [1::Int,1,2,3,4,4,5,1,5,7]Stream.fold Fold.toList $ Stream.scanMaybe Fold.nub stream[1,2,3,4,5,7]

Pre-release

valuetoIntSet :: Monad m => Fold m Int IntSet
#

Fold the input to an int set. For integer inputs this performs better than toSet.

Definition:

Example1 expression
toIntSet = Fold.foldl' (flip IntSet.insert) IntSet.empty
valuetoMap :: (Monad m, Ord k) => (a -> k) -> Fold m a b -> Fold m a (Map k b)
#

Split the input stream based on a key field and fold each split using the given fold. Useful for map/reduce, bucketizing the input in different bins or for generating histograms.

Example:

Example2 expressions
import Data.Map.Strict (Map):{ let input = Stream.fromList [("ONE",1),("ONE",1.1),("TWO",2), ("TWO",2.2)]     classify = Fold.toMap fst (Fold.lmap snd Fold.toList)  in Stream.fold classify input :: IO (Map String [Double]):}fromList [("ONE",[1.0,1.1]),("TWO",[2.0,2.2])]

Once the classifier fold terminates for a particular key any further inputs in that bucket are ignored.

Space used is proportional to the number of keys seen till now and monotonically increases because it stores whether a key has been seen or not.

See demuxToMap for a more powerful version where you can use a different fold for each key. A simpler version of toMap retaining only the last value for a key can be written as:

Example1 expression
toMap = Fold.foldl' (\kv (k, v) -> Map.insert k v kv) Map.empty

Stops: never

Pre-release

valuetoMapIO
  1. :: (MonadIO m, Ord k)
  2. => a -> k
  3. -> Fold m a b
  4. -> Fold m a (Map k b)
#

Same as toMap but maybe faster because it uses mutable cells as fold accumulators in the Map.

valuetoSet :: (Monad m, Ord a) => Fold m a (Set a)
#

Fold the input to a set.

Definition:

Example1 expression
toSet = Fold.foldl' (flip Set.insert) Set.empty
valueaddOne :: Monad m => a -> Fold m a b -> m (Fold m a b)
#

Append a singleton value to the fold.

See examples under addStream.

Pre-release

valuecatEithers :: Fold m a b -> Fold m (Either a a) b
#

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

Definition:

Example1 expression
catEithers = Fold.lmap (either id id)

Pre-release

valueduplicate :: Monad m => Fold m a b -> Fold m a (Fold m a b)
#

duplicate provides the ability to run a fold in parts. The duplicated fold consumes the input and returns the same fold as output instead of returning the final result, the returned fold can be run later to consume more input.

duplicate essentially appends a stream to the fold without finishing the fold. Compare with snoc which appends a singleton value to the fold.

Pre-release

valuefoldlM1' :: Monad m => (a -> a -> m a) -> Fold m a (Maybe a)
#

Like 'foldl1'' but with a monadic step function.

Pre-release

valuefoldr' :: Monad m => (a -> b -> b) -> b -> Fold m a b
#

Make a fold using a right fold style step function and a terminal value. It performs a strict right fold via a left fold using function composition. Note that a strict right fold can only be useful for constructing strict structures in memory. For reductions this will be very inefficient.

Definitions:

Example2 expressions
foldr' f z = fmap (flip appEndo z) $ Fold.foldMap (Endo . f)foldr' f z = fmap ($ z) $ Fold.foldl' (\g x -> g . f x) id

Example:

Example1 expression
Stream.fold (Fold.foldr' (:) []) $ Stream.enumerateFromTo 1 5[1,2,3,4,5]
valuelmap :: (a -> b) -> Fold m b r -> Fold m a r
#

lmap f fold maps the function f on the input of the fold.

Definition:

Example1 expression
lmap = Fold.lmapM return

Example:

Example2 expressions
sumSquared = Fold.lmap (\x -> x * x) Fold.sumStream.fold sumSquared (Stream.enumerateFromTo 1 100)338350
valuelmapM :: Monad m => (a -> m b) -> Fold m b r -> Fold m a r
#

lmapM f fold maps the monadic function f on the input of the fold.

valuermapM :: Monad m => (b -> m c) -> Fold m a b -> Fold m a c
#

Map a monadic function on the output of a fold.

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

Use a Maybe returning fold as a filtering scan.

Example1 expression
scanMaybe p f = Fold.postscan p (Fold.catMaybes f)

Pre-release

valuesplitWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> Fold m x a
  4. -> Fold m x b
  5. -> Fold m x c
#

Sequential fold application. Apply two folds sequentially to an input stream. The input is provided to the first fold, when it is done - the remaining input is provided to the second fold. When the second fold is done or if the input stream is over, the outputs of the two folds are combined using the supplied function.

Example:

Example4 expressions
header = Fold.take 8 Fold.toListline = Fold.takeEndBy (== '\n') Fold.toListf = Fold.splitWith (,) header lineStream.fold f $ Stream.fromList "header: hello\n"("header: ","hello\n")

Note: This is dual to appending streams using Data.Stream.append.

Note: this implementation allows for stream fusion but has quadratic time complexity, because each composition adds a new branch that each subsequent fold's input element has to traverse, therefore, it cannot scale to a large number of compositions. After around 100 compositions the performance starts dipping rapidly compared to a CPS style implementation.

For larger number of compositions you can convert the fold to a parser and use ParserK.

Time: O(n^2) where n is the number of compositions.

valueteeWith
  1. :: Monad m
  2. => a -> b -> c
  3. -> Fold m x a
  4. -> Fold m x b
  5. -> Fold m x c
#

teeWith k f1 f2 distributes its input to both f1 and f2 until both of them terminate and combines their output using k.

Definition:

Example1 expression
teeWith k f1 f2 = fmap (uncurry k) (Fold.tee f1 f2)

Example:

Example2 expressions
avg = Fold.teeWith (/) Fold.sum (fmap fromIntegral Fold.length)Stream.fold avg $ Stream.fromList [1.0..100.0]50.5

For applicative composition using this combinator see Streamly.Data.Fold.Tee.

See also: Streamly.Data.Fold.Tee

Note that nested applications of teeWith do not fuse.

Concurrent Operations

0 declarations

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.

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.

valueboundThreads :: Bool -> Config -> Config
#

Spawn bound threads (i.e., spawn threads using forkOS instead of forkIO). The default value is False.

Currently, this only takes effect only for concurrent folds.

Combinators

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

Evaluate a fold asynchronously using a concurrent channel. The driver just queues the input stream values to the fold channel buffer and returns. The fold evaluates the queued values asynchronously. On finalization, parEval waits for the asynchronous fold to complete before it returns.

Container Related

1 declaration
valuetoHashMapIO
  1. :: (MonadIO m, Hashable k, Ord k)
  2. => a -> k
  3. -> Fold m a b
  4. -> Fold m a (HashMap k b)
#

Split the input stream based on a hashable component of the key field and fold each split using the given fold. Useful for map/reduce, bucketizing the input in different bins or for generating histograms.

Example3 expressions
import Data.HashMap.Strict (HashMap, fromList)import qualified Streamly.Data.Fold.Prelude as Foldimport qualified Streamly.Data.Stream as Stream

Consider a stream of key value pairs:

Example1 expression
input = Stream.fromList [("k1",1),("k1",1.1),("k2",2), ("k2",2.2)]

Classify each key to a different hash bin and fold the bins:

Example2 expressions
classify = Fold.toHashMapIO fst (Fold.lmap snd Fold.toList)Stream.fold classify input :: IO (HashMap String [Double])fromList [("k2",[2.0,2.2]),("k1",[1.0,1.1])]

Pre-release