# Thread safe buffer correctness

**URL:** <https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183>\
**Category:** General Usage\
**Tags:** multithreading\
**Created:** [April 28, 2021, 5:24pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183 "2021-04-28T17:24:02Z")\
**Posts on this page:** 14\
**Page:** 1

<div class="post-metadata">

**Author:** ![halleysfifthinc](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/halleysfifthinc/32/206280_2.png) [@halleysfifthinc](https://discourse.julialang.org/u/halleysfifthinc)\
**Post date:** [April 28, 2021, 5:24pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/1 "2021-04-28T17:24:02Z")

</div>

For a recent multi-threaded simulation, I needed a thread-safe way to use pre-allocated arrays, and the simple “`nthreads()` length vector of arrays” wasn’t working due to the presence of `yield` points in the spawn’ed tasks. (The problem noted by @tkf [here](https://juliafolds.github.io/FLoops.jl/dev/explanation/faq/#faq-state-threadid).)

I came up with a fairly simple type to add some locks around the pre-allocated arrays, and it solved my problem and the simulation results look correct. However, I’m not very familiar with multi-threading, so I’m curious if there are any obvious drawbacks or potential issues with this approach:

```julia
struct SafeBuffer{T,N}
    held::Semaphore
    locks::NTuple{N,ReentrantLock}
    buffers::Vector{T}
end

function acquirebuffer(buffer)
    acquire(buffer.held)
    @inbounds for ci in eachindex(buffer.locks)
        if trylock(buffer.locks[ci])
            return buffer.locks[ci], buffer.buffers[ci]
        end
    end

    return nothing
end

function releasebuffer(lock, buffer)
    unlock(lock)
    release(buffer.held)
end

function withbuffer(f, buffer)
    l, buf = acquirebuffer(buffer)
    try
        f(buf)
    finally
        releasebuffer(l, buffer)
    end
end

```

---

<div class="post-metadata">

**Author:** ![djholiver](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/djholiver/32/50470_2.png) [@djholiver](https://discourse.julialang.org/u/djholiver)\
**Post date:** [April 28, 2021, 6:10pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/2 "2021-04-28T18:10:52Z")

</div>

Hi,

Not to reply to a question with a question, but aren’t you describing a need to do a thread safe multi producer - multi consumer pattern? Would that not be better suited to Channels of fixed size?

Regards

---

<div class="post-metadata">

**Author:** ![halleysfifthinc](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/halleysfifthinc/32/206280_2.png) [@halleysfifthinc](https://discourse.julialang.org/u/halleysfifthinc)\
**Post date:** [April 28, 2021, 6:35pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/3 "2021-04-28T18:35:36Z")

</div>

I’ve never been able to understand concurrent producer/consumer patterns with Channels, so it isn’t clear to me what that would look like. In particular, I’m not sure how you would reuse arrays with that approach.

---

<div class="post-metadata">

**Author:** ![djholiver](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/djholiver/32/50470_2.png) [@djholiver](https://discourse.julialang.org/u/djholiver)\
**Post date:** [April 28, 2021, 7:02pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/4 "2021-04-28T19:02:11Z")

</div>

The [async docs](https://docs.julialang.org/en/v1/manual/asynchronous-programming/) provide a really good example halfway down.

As you will see - you can create a channel of fixed type and size (which would be your array replacement) and write / read to them asynchronously.

If you need to lock access to an array / vector of fixed size and then read / write by index, the above wont suit since the Channel approach is analogous to pipes and you only get access to top / bottom element. It may also require a bit of a rewrite. My assumption here is that you need to fill something up with multiple threads and then potentially empty or consume as its filled; with the fixed size for performance.

Obviously a bit of a guess (tricky without more info) and still not answering your original question…

Regards

---

<div class="post-metadata">

**Author:** ![tkf](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tkf/32/17635_2.png) [@tkf](https://discourse.julialang.org/u/tkf)\
**Post date:** [April 28, 2021, 7:58pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/5 "2021-04-28T19:58:24Z")

</div>

For `Channel`-based producer-consumer approach with consumer-local buffer reuse, see the first example in [Pattern for managing thread local storage? - #2 by tkf](https://discourse.julialang.org/t/pattern-for-managing-thread-local-storage/54023/2)

This pattern is very handy and performs reasonably well in many situations. But, since the whole purpose of `Channel` is to serialize the access, it’s not really the best approach. If your input collection is [supported by SplittablesBase.jl](https://github.com/JuliaFolds/SplittablesBase.jl#supported-collections), FLoops.jl’s [`@init`](https://juliafolds.github.io/FLoops.jl/dev/howto/parallel/#private-variables) can be useful.

But, if you need a buffer pool that works outside producer-consumer pattern or parallel loop, I see that something like `SafeBuffer` in the OP could be handy. One way to simplify the synchronization is:

```julia
struct BufferPool{T}
    buffers::Vector{T}
    condition::Threads.Condition
end

function bufferpool(factory, n::Integer)
    n > 0 || error("need at least one buffer; got n = $n")
    return BufferPool([factory() for _ in 1:n], Threads.Condition())
end

function withbuffer(f, pool::BufferPool)
    b = lock(pool.condition) do
        while isempty(pool.buffers)
            wait(pool.condition)
        end
        pop!(pool.buffers)
    end
    try
        f(b)
    finally
        lock(pool.condition) do
            push!(pool.buffers, b)
            notify(pool.condition)
        end
    end
end

```

---

<div class="post-metadata">

**Author:** ![halleysfifthinc](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/halleysfifthinc/32/206280_2.png) [@halleysfifthinc](https://discourse.julialang.org/u/halleysfifthinc)\
**Post date:** [April 28, 2021, 8:08pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/6 "2021-04-28T20:08:18Z")

</div>

Do you have something along the lines of my `SafeBuffer` or `BufferPool` implemented in one of your packages? (Or are you aware of such a thing in a published package?)

Can you comment on any benefits (performance or otherwise) to `BufferPool` over `SafeBuffer`?

---

<div class="post-metadata">

**Author:** ![tkf](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tkf/32/17635_2.png) [@tkf](https://discourse.julialang.org/u/tkf)\
**Post date:** [April 28, 2021, 8:28pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/7 "2021-04-28T20:28:30Z")

</div>

No, I don’t think I’ve used something like this in released packages so far.

`SafeBuffer` vs `BufferPool` is just my reaction to “hmm… there are many locks… can I use less?” Since `Base.Semaphore` uses a `Threads.Condtion` anyway, I think the performance can be better for `BufferPool` since it uses strictly less locks. But I don’t know if that’s noticeable. Also, ideally, it’d be better if I we had a lock-free (-ish) stack for something like buffer pool (which I’ve been wanting to play with).

---

<div class="post-metadata">

**Author:** ![Ralph\_Smith](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/ralph_smith/32/10344_2.png) [@Ralph\_Smith](https://discourse.julialang.org/u/Ralph_Smith)\
**Post date:** [April 29, 2021, 2:25pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/8 "2021-04-29T14:25:05Z")

</div>

> [@tkf](#):
>
> ```julia
> lock(pool.condition) do
> push!(pool.buffers, b)
> notify(pool.condition)
> end
> 
> ```

Why are you waking all the tasks when you are only adding one buffer?

---

<div class="post-metadata">

**Author:** ![pbayer](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pbayer/32/11675_2.png) [@pbayer](https://discourse.julialang.org/u/pbayer)\
**Post date:** [April 29, 2021, 5:14pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/9 "2021-04-29T17:14:20Z")

</div>

You can implement a lock-free and thread-safe buffer in [Actors](https://github.com/JuliaActors/Actors.jl). This is based on channels and buffer access is serialized through the actor’s message channel (link). A simple implementation would be:

```julia
module SafeBuffer
using Actors

export newBuffer

struct Buffer{T}
    lk::Link
end
(b::Buffer)() = call(b.lk)

# actor behavior
buf(x) = copy(x)
buf(x, f, args...) = f(x, args...)

# interface function
newBuffer(x::Vector{T}) where T = Buffer{Vector{T}}(Actors.spawn(buf, x))

Actors.call(b::Buffer, args...) = call(b.lk, args...) 
end

```

with that you can do:

```julia
julia> include("safebuffer.jl")
Main.SafeBuffer

julia> using Actors, .SafeBuffer, .Threads

julia> sf = newBuffer([1,2,3,4,5])
Main.SafeBuffer.Buffer{Vector{Int64}}(Link{Channel{Any}}(Channel{Any}(32), 1, :default))

julia> sf()
5-element Vector{Int64}:
 1
 2
 3
 4
 5

julia> @threads for _ in 1:495
           call(sf, push!, threadid())
       end

julia> length(sf())
500

julia> map(x->length(filter(==(x), sf())), 1:nthreads())
16-element Vector{Int64}:
 32
 32
 32
 32
  ⋮
 31
 31
 31
 30

julia> @threads for _ in 1:499
           call(sf, pop!)
       end

julia> sf()
1-element Vector{Int64}:
 1

```

Note: the [Guards](https://github.com/JuliaActors/Guards.jl) actors library provides a more comfortable implementation.

---

<div class="post-metadata">

**Author:** ![halleysfifthinc](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/halleysfifthinc/32/206280_2.png) [@halleysfifthinc](https://discourse.julialang.org/u/halleysfifthinc)\
**Post date:** [April 29, 2021, 5:55pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/10 "2021-04-29T17:55:04Z")

</div>

How does that enable the use and reuse of separate arrays? The use case being described here is multiple calls along the lines of:

```julia
@spawn withbuffer(pool) do buf
    fill_array_somehow!(buf)
    reduce_array_func(buf)
end

```

where there are a number of pre-allocated arrays, one of which is “claimed” and used by mutating-bang functions with a guarantee that the contents of the array `buf` won’t be changed by another spawned task. It looks like there is only one array in your solution?

---

<div class="post-metadata">

**Author:** ![pbayer](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pbayer/32/11675_2.png) [@pbayer](https://discourse.julialang.org/u/pbayer)\
**Post date:** [April 29, 2021, 7:32pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/11 "2021-04-29T19:32:48Z")

</div>

You would have an actor for each array and access them concurrently as you want. There will be no races. If multiple tasks have the actor’s link, they can send it messages and therefore change the guarded buffer.

> [@halleysfifthinc](#):
>
> with a guarantee that the contents of the array `buf` won’t be changed by another spawned task

For that you would have to ensure that only one task at a time has access to a buffer. With `Actors` you can write a pool-actor managing the access to a buffer pool. If a task calls `acquire`, the pool-actor would give it access to a buffer, if it calls `release`, the actor takes it back. In your example the `withbuffer` task finishes after calling release. The reference to the buffer then is destroyed and it would be safe for the pool-actor to hand over variables or references.

---

<div class="post-metadata">

**Author:** ![tkf](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tkf/32/17635_2.png) [@tkf](https://discourse.julialang.org/u/tkf)\
**Post date:** [April 30, 2021, 8:28pm UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/12 "2021-04-30T20:28:02Z")

</div>

> [@Ralph\_Smith](#):
>
> Why are you waking all the tasks when you are only adding one buffer?

Since Julia does not support migration of tasks across OS threads (yet), waking up one task may not maximize the thread usage. But yeah, you can manage waiter tasks more cleverly like per-thread queues.

> [@pbayer](#):
>
> You can implement a lock-free and thread-safe buffer in [Actors](https://github.com/JuliaActors/Actors.jl)

IIUC, Actors.jl looks like an abstraction based on `Channel`. Since `Channel` does use locks, it is not lock-free. Usually, lock-free data structure means, very roughly speaking, “use atomics instead of locks and alike for non-blocking communication” and in loose conversation often include wait-free and obstruction-free. Of course, since you can trivially write a spin lock using atomics, the key point is how the atomics are used. See, e.g., [Non-blocking algorithm - Wikipedia](https://en.wikipedia.org/wiki/Non-blocking_algorithm) for proper definitions of them. I saw people using lock-free to just simply mean “not using lock” but that’s not what the term means for experts.

---

<div class="post-metadata">

**Author:** ![pbayer](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/pbayer/32/11675_2.png) [@pbayer](https://discourse.julialang.org/u/pbayer)\
**Post date:** [May 1, 2021, 9:26am UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/13 "2021-05-01T09:26:59Z")

</div>

> [@tkf](#):
>
> I saw people using lock-free to just simply mean “not using lock” but that’s not what the term means for experts.

Thank you for the clarification.

Still it makes a difference if

- a language construct or a standard library function is implemented using a lock - where [lock ordering can be carefully controlled](https://docs.julialang.org/en/v1.6/devdocs/locks/) – or if a user or 2nd party package uses them,
- in solving a concurrency problem you share memory using locks or if you share data by communicating (as [Go people put it](https://golang.org/doc/effective_go#sharing)).

We probably agree that locks should be avoided or are over-used. One problem with them is, that they are so easy to use and to botch everything in a concurrent program.

A `Channel` is a higher-level language construct to serialize concurrent access to a shared resource or service. But it is not an easy alternative to locks since it always needs some function (task, goroutine …) to serve the resource. [Actors.jl](https://github.com/JuliaActors/Actors.jl) abstracts away the channel and basically focusses on the serving.

BTW if in a future Julia version – maybe with atomics – there will be a lock-free implementation of `Channel`, I will be happy to use that. But a stateful receiver (a `take!`) on a channel will always block until there is data to receive.

---

<div class="post-metadata">

**Author:** ![tkf](https://sea2.discourse-cdn.com/julialang/user_avatar/discourse.julialang.org/tkf/32/17635_2.png) [@tkf](https://discourse.julialang.org/u/tkf)\
**Post date:** [May 2, 2021, 1:17am UTC](https://discourse.julialang.org/t/thread-safe-buffer-correctness/60183/14 "2021-05-02T01:17:21Z")

</div>

> [@pbayer](#):
>
> We probably agree that locks should be avoided or are over-used.

Absolutely! This is why my first suggestion was to use `Channel`. I agree “share memory by communicating” is the best starting point for concurrent programming.
