# Parallel mutable reduction / mapreduce

**URL:** <https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207>\
**Category:** General Usage\
**Tags:** question, package, multithreading\
**Created:** [February 14, 2024, 2:35pm UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207 "2024-02-14T14:35:34Z")\
**Posts on this page:** 8\
**Page:** 1

<div class="post-metadata">

**Author:** ![foobar\_lv2](https://avatars.discourse-cdn.com/v4/letter/f/ee59a6/32.png) [@foobar\_lv2](https://discourse.julialang.org/u/foobar_lv2)\
**Post date:** [February 14, 2024, 2:35pm UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/1 "2024-02-14T14:35:34Z")

</div>

Hi,  
I wanted to ask about recommendations for packages that offer multithreaded mapreduce-style computation. That is unfortunately not part of stdlib, so I was wondering which packages are good for that.

The way I like to think about my problems is in terms like the (imo very good!) [java stream collect](https://docs.oracle.com/javase/8/docs/api/java/util/stream/Stream.html#collect-java.util.function.Supplier-java.util.function.BiConsumer-java.util.function.BiConsumer-) API:

I have an iterator of `T` items, I have a `makeCollector()::Collector`, I have an `accept(collector::Collector, item::T)::Collector`, and a function `merge(left::Collector, right::Collector)::Collector`.

On a list e.g. `[item1, item2, item3]` this should compute

```julia
foldcollectthing(makeCollector, [item1, item2, item3]) = accept(accept(accept(makeCollector(), item1), item2), item3)

```

somewhat `foldl` style. Parallelism is enabled by associativity:

```julia
accept(left, item) == merge(left, accept(makeCollector(), item))
merge(merge(A, B), C) == merge(A, merge(B, C))

```

This allows the implementation of the desired `foldcollectthing` to split the input into chunks, run on each chunk on a different thread, and then use `merge` to combine the result.

Typically `accept` and `merge` are mutating, i.e. `accept(collector, item) === collector` and `merge(left, right) === left`, and `merge` can destroy the right collector. If the amount of time taken for each item is known to be homogeneous, ideally `makeCollector` should be called once per thread only.

A typical example how that kind view works is the following baby example:

```julia
julia> foo(n)=collect(1:n)
julia> mapreduce(foo, vcat, 1:3)
6-element Vector{Int64}:
 1
 1
 2
 1
 2

```

Using this paradigm, this would be spelled:

```julia
julia> mutable struct Collector items::Vector{Int} end

julia> makeCollector() = Collector(Int[])
makeCollector (generic function with 1 method)

julia> function accept(collector::Collector, i)
       for res in foo(i)
       push!(collector, res)
       end
       collector
       end
accept (generic function with 1 method)
julia> function merge(left::Collector, right::Collector)
       if length(left.items) < length(right.items)
       prepend!(right.items, left.items)
       left.items = right.items
       else
       append!(left.items, right.items)
       end
       left
       end

```

As a side-note, the ` if length(left.items) < length(right.items)` case is essential: It makes the difference between O(N log N) and O(N^2) runtime for collecting N items in the worst-case schedule; this kind of thing really needs a double-ended queue and cannot be done with e.g. C++ Vector.

=================

Any recommendations?

Am I too blind to see how to do this with Transducers.jl @tk3369 ? I.e. given some `makeCollector`, `merge` and `accept` functions, how would I use `foldxt` to do this job on e.g. a plain `Vector` of items?

PS. Another way of writing the desired result would be

```julia
mapreduce(item -> accept(makeCollector(), item), merge, collection; init = makeCollector())

```

The essential part, however, is that `accept` is possibly non-allocating and cheaper than `merge`. I don’t want brittle compiler optimizations to possibly maybe elide such allocations, I want it guaranteed.

---

<div class="post-metadata">

**Author:** ![tk3369](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tk3369/32/2824_2.png) [@tk3369](https://discourse.julialang.org/u/tk3369)\
**Post date:** [February 15, 2024, 2:37am UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/2 "2024-02-15T02:37:10Z")

</div>

You probably want @tkf not me? 🙂

---

<div class="post-metadata">

**Author:** ![jar1](https://avatars.discourse-cdn.com/v4/letter/j/c0e974/32.png) [@jar1](https://discourse.julialang.org/u/jar1)\
**Post date:** [February 15, 2024, 3:48am UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/3 "2024-02-15T03:48:57Z")

</div>

Is there something wrong with just substituting `ThreadsX.mapreduce` for `mapreduce` and using `push!(collector.items, res)`? Sorry if I’m misunderstanding the requirements here.

```julia
using ThreadsX

mutable struct Collector
    items::Vector{Int}
end

makeCollector() = Collector(Int[])

function accept(collector::Collector, i)
    for res in foo(i)
        push!(collector.items, res)
    end
    collector
end

function merge(left::Collector, right::Collector)
    if length(left.items) < length(right.items)
        prepend!(right.items, left.items)
        left.items = right.items
    else
        append!(left.items, right.items)
    end
    left
end

foo(n)=collect(1:n)

ThreadsX.mapreduce(item -> accept(makeCollector(), item), merge, 1:5; init = makeCollector())

```

---

<div class="post-metadata">

**Author:** ![foobar\_lv2](https://avatars.discourse-cdn.com/v4/letter/f/ee59a6/32.png) [@foobar\_lv2](https://discourse.julialang.org/u/foobar_lv2)\
**Post date:** [February 15, 2024, 10:21am UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/4 "2024-02-15T10:21:11Z")

</div>

> [@tk3369](#):
>
> You probably want @tkf not me? 🙂

Oops, sorry and thank you! The question was for you, @tkf

> [@jar1](#):
>
> Is there something wrong with just substituting `ThreadsX.mapreduce` for `mapreduce` and using `push!(collector.items, res)`?

Yes: The wrong thing is that any `mapreduce` style API will allocate `length(input)` many collectors – it will, in a loop call something like

```julia
current = merge(current, accept(makeCollector(), idx))

```

The point of the java stream mutable reduction interface is that this can be written as

```julia
current = accept(current, idx)

```

and no new `Collector` needs to be created and almost immediately discarded. Indeed, the number of `makeCollector()` calls scales only with the number of threads / tasks.

In this sense, this API is conceptually superior to mapreduce for settings where it is possible to optimize the case where a single item is merged from the right into a collector.

The implicit claims in my questions were:

1. This is common – many settings permit a fastpath for handling a single element
2. This is not inherently conceptually confusing – once you get used to it, it is no more brain-warping to express a mapreduce problem in terms of `merge / accept` as opposed to `merge / map`.

Both claims are backed by the astounding success and precedent of this part of the java stream API.

---

<div class="post-metadata">

**Author:** ![jar1](https://avatars.discourse-cdn.com/v4/letter/j/c0e974/32.png) [@jar1](https://discourse.julialang.org/u/jar1)\
**Post date:** [February 16, 2024, 11:50pm UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/5 "2024-02-16T23:50:31Z")

</div>

I took another shot. This defines nonallocating `NilCollector` and `SingletonCollector`, and adds a counter property `i` to `Collector` so we can check that it is only instantiated some number of times smaller than `length(xs)` . In this case the number of data items `length(xs) == 100` and `Collector` is instantiated 16 times on my computer with `nthreads() == 12`.

```julia
using Folds

mutable struct Collector
    items::Vector{Int}
    i::Int # count how many times this type is instantiated.
end
makeCollector() = Collector(Int[], 1)

accept!(col::Collector, item) = begin
    push!(col.items, item)
    col
end

# `!!` methods use the mutate-or-widen pattern: mutate if possible, otherwise create a new object that can hold the result.

accept!!(c::Collector, item) = accept!(c, item)

struct SingletonCollector
    item::Int
end

accept!!(c::SingletonCollector, item) = merge!!(c, SingletonCollector(item))

struct NilCollector end

accept!!(::NilCollector, item) = SingletonCollector(item)

function merge!(left::Collector, right::Collector)
    if length(left.items) < length(right.items)
        prepend!(right.items, left.items)
        left.items = right.items
    else
        append!(left.items, right.items)
    end
    left.i += right.i
    left
end

merge!!(::NilCollector, ::NilCollector) = NilCollector()
merge!!(left, right::NilCollector) = left
merge!!(left::NilCollector, right) = right
merge!!(left::SingletonCollector, right::SingletonCollector) = accept!!(accept!!(makeCollector(), left.item), right.item)
merge!!(left::Collector, right::SingletonCollector) = begin
    push!(left.items, right.item)
    left
end
merge!!(left::SingletonCollector, right::Collector) = begin
    insert!(right.items, firstindex(right.items), left.item)
    right
end
merge!!(left::Collector, right::Collector) = merge!(left, right)

Folds.mapreduce(SingletonCollector, merge!!, 1:100, ThreadedEx(); init=NilCollector())

```

---

<div class="post-metadata">

**Author:** ![bertschi](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/bertschi/32/33462_2.png) [@bertschi](https://discourse.julialang.org/u/bertschi)\
**Post date:** [February 17, 2024, 1:00pm UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/6 "2024-02-17T13:00:22Z")

</div>

Imho, no Java API has ever been good as they tend to over-enigineer things in terms of types, i.e., it’s hard to see the trees for the forest.  
Thus, here is my take using `Transducers`:

```julia
using Transducers
# sequential version
newCollector() = @show Int[] # show to see how often it got called
accept(c::Vector, v) = push!(c, v) # No need for a new type
foldl(accept, Map(x -> x^2), 1:100; init = OnInit(newCollector))

# possibly parallel version
singleton(x) = (x,) # cheap non-allocating wrapper for singleton
function merge(cl::Vector, cr::Vector)
    if length(cl) < length(cr)
        prepend!(cr, cl)
    else
       append!(cl, cr)
    end
end
merge(c::Vector, x::Tuple) = push!(c, x[1])
# Could use a new type wrapping singletons, but a tuple does just fine
foldxt(merge, Map(x -> singleton(x^2)), 1:100; init = OnInit(newCollector))

```

---

<div class="post-metadata">

**Author:** ![foobar\_lv2](https://avatars.discourse-cdn.com/v4/letter/f/ee59a6/32.png) [@foobar\_lv2](https://discourse.julialang.org/u/foobar_lv2)\
**Post date:** [February 18, 2024, 10:49pm UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/7 "2024-02-18T22:49:50Z")

</div>

> [@bertschi](#):
>
> Imho, no Java API has ever been good as they tend to over-enigineer things in terms of types, i.e., it’s hard to see the trees for the forest.

This is nonsense. They do tend to over-engineer and many java APIs suck, but sometimes they hit an exceptionally good spot. The stream collector API is really really good.

> [@bertschi](#):
>
> Thus, here is my take using `Transducers`:

Thanks! I like that.

A shortened form would be:

```julia
julia> using Transducers

julia> newCollector() = Int[]
newCollector (generic function with 1 method)

julia> merge(left::Vector, right::Vector) = length(left) < length(right) ? prepend!(right, left) : append!(left, right)
merge (generic function with 1 method)

julia> merge(left::Vector, item) = push!(left, item)
merge (generic function with 2 methods)

```

This uses some additional properties of `foldxt` in order to be type-stable. I do wonder whether that is guaranteed to work, or only works incidentially?

Apriori it would be valid for foldxt to evaluate this via e.g.

```julia
merge(input[1], merge(oninit.f(), input[2]))

```

This would crash and burn with the provided definitions.

In other words: This entire thing works and is type-stable with `OnInit.f()::CollectorT`, and `merge(::CollectorT, ::CollectorT)::CollectorT` and `merge(::CollectorT, ::InputT)::CollectorT`, which is exactly the API I asked for.

A naive reading would have suggested that `merge(::InputT, ::InputT)::CollectorT` and `merge(::InputT, ::CollectorT)::CollectorT` are also required, and that this is not type-stable / fully inferred unless `InputT == CollectorT`.

After reading a little more Transducers.jl code, it appears that it internally does use the sane API (java accept is called `Transducers.next`, java combine is called `Transducers.combine`).

---

<div class="post-metadata">

**Author:** ![bertschi](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/bertschi/32/33462_2.png) [@bertschi](https://discourse.julialang.org/u/bertschi)\
**Post date:** [February 19, 2024, 6:27pm UTC](https://discourse.julialang.org/t/parallel-mutable-reduction-mapreduce/110207/8 "2024-02-19T18:27:42Z")

</div>

> [@foobar\_lv2](#):
>
> > [@bertschi](#):
> >
> > Imho, no Java API has ever been good as they tend to over-enigineer things in terms of types, i.e., it’s hard to see the trees for the forest.
> 
> This is nonsense. They do tend to over-engineer and many java APIs suck, but sometimes they hit an exceptionally good spot. The stream collector API is really really good.

Ok, fair enough. Had used some Java APIs previously, e.g., the technically very good Fork-Join pool. Nevertheless, the API always became nicer when used from Clojure (or Scala) as you could just pass a couple functions instead of implementing yet another interface.

> [@foobar\_lv2](#):
>
> Apriori it would be valid for foldxt to evaluate this via e.g.
> 
> ```julia
> merge(input[1], merge(oninit.f(), input[2]))
> 
> ```
> 
> This would crash and burn with the provided definitions.

Well, as it’s supposed to generalize `foldl` I would assume something like:

```julia
merge(merge(merge(oninit.f(), input[1]), input[2]),
      merge(merge(oninit.f(), input[3]), input[4]))

```

You are right though, that this does not seem to be explicitly documented.
