# Implementing a fiber fed queue

**URL:** <https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090>\
**Category:** Help & Support\
**Created:** [October 15, 2023, 3:41pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090 "2023-10-15T15:41:13Z")\
**Posts on this page:** 7\
**Page:** 1

<div class="post-metadata">

**Author:** ![Xen](https://yyz2.discourse-cdn.com/flex036/user_avatar/forum.crystal-lang.org/xen/32/2288_2.png) [@Xen](https://forum.crystal-lang.org/u/Xen)\
**Post date:** [October 15, 2023, 3:41pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090/1 "2023-10-15T15:41:13Z")

</div>

So I want to make a queue that gets processed in my main fiber, which is fed by other fibers.

Obviously a loop that checks if there’s anything in the queue, processes it, or sleeps for a while would work. But I’d like to minimize the latency but the lesser the sleep period, the more CPU usage. And it’s going to spend most of its time waiting.

So something with a channel that gets written to when the other fibers enqueue messages, and the main fiber `receive`ing but as `send` blocks, we’re essentially using the channel as the queue and blocking the other fibers when it gets full.

Tried looking for a simple queue implementation, but most seems to be aiming towards using a service, and that’s overkill for my situation.

Any hints?

---

<div class="post-metadata">

**Author:** ![Blacksmoke16](https://yyz2.discourse-cdn.com/flex036/user_avatar/forum.crystal-lang.org/blacksmoke16/32/1241_2.png) [@Blacksmoke16](https://forum.crystal-lang.org/u/Blacksmoke16)\
**Post date:** [October 15, 2023, 4:26pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090/2 "2023-10-15T16:26:17Z")

</div>

Feel like [5 use cases for Crystal's select statement - lbarasti's blog](https://lbarasti.com/post/select_statement/) might push you in the right direction.

---

<div class="post-metadata">

**Author:** ![jgaskins](https://yyz2.discourse-cdn.com/flex036/user_avatar/forum.crystal-lang.org/jgaskins/32/2449_2.png) [@jgaskins](https://forum.crystal-lang.org/u/jgaskins)\
**Post date:** [October 15, 2023, 5:03pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090/3 "2023-10-15T17:03:02Z")

</div>

I wrote [a shard](https://github.com/jgaskins/mpsc) for this a couple years ago based on [the Rust `mpsc` library](https://doc.rust-lang.org/std/sync/mpsc/). It’s not a perfect recreation of it, but it’s an unbounded channel that was explicitly not designed to handle multiple consumers the way the stdlib `Channel` does. It will grow as large as it needs to in order to accommodate the data it’s being fed, allowing for surges in input.

If you need it, `MPSC::Channel#receive?` returns `nil` if it’s empty, so you also get atomic emptiness checks.

I originally wrote it to support [my `opentelemetry` shard](https://github.com/jgaskins/opentelemetry). During load tests, using a stdlib `Channel` throttled incoming HTTP requests to my app because the `Channel` buffer was full, so I needed something that would never block on `send`. Reallocating buffers was preferred over blocking for that use case.

Caveats:

- it’s strictly single-consumer, so you can’t consume it from multiple fibers — it will raise an exception if you try
- it doesn’t (currently) support `select`, so you can’t have it timeout

---

<div class="post-metadata">

**Author:** ![sablokgaurav](https://yyz2.discourse-cdn.com/flex036/user_avatar/forum.crystal-lang.org/sablokgaurav/32/2454_2.png) [@sablokgaurav](https://forum.crystal-lang.org/u/sablokgaurav)\
**Post date:** [October 15, 2023, 5:04pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090/4 "2023-10-15T17:04:39Z")

</div>

@Blacksmoke16 @Xen I checked this link which you have mentioned and it uses the select statement. i think this can be much shorter and easier if the select is merged with the .each. if multiple requests are there such .each.select { |each| puts “some text” if v=terminal.recieve } what do you think?.

---

<div class="post-metadata">

**Author:** ![Blacksmoke16](https://yyz2.discourse-cdn.com/flex036/user_avatar/forum.crystal-lang.org/blacksmoke16/32/1241_2.png) [@Blacksmoke16](https://forum.crystal-lang.org/u/Blacksmoke16)\
**Post date:** [October 15, 2023, 5:09pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090/5 "2023-10-15T17:09:39Z")

</div>

I’m not sure what `.each` you’re referring to, but the `select` talked about in the blog is _not_ the same thing as `Enumerable#select`. So I don’t think it’s as simple as you think it is.

---

<div class="post-metadata">

**Author:** ![sablokgaurav](https://yyz2.discourse-cdn.com/flex036/user_avatar/forum.crystal-lang.org/sablokgaurav/32/2454_2.png) [@sablokgaurav](https://forum.crystal-lang.org/u/sablokgaurav)\
**Post date:** [October 15, 2023, 5:17pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090/6 "2023-10-15T17:17:06Z")

</div>

@Blacksmoke16 do you think can it be put like this as an add-on such as if multiple request are coming then open an array such as array = Array.new() and instead of puts v, it can be added to the array as array.append(requests) and then if array.length == 0 break and terminate.

---

<div class="post-metadata">

**Author:** ![Xen](https://yyz2.discourse-cdn.com/flex036/user_avatar/forum.crystal-lang.org/xen/32/2288_2.png) [@Xen](https://forum.crystal-lang.org/u/Xen)\
**Post date:** [October 15, 2023, 6:48pm UTC](https://forum.crystal-lang.org/t/implementing-a-fiber-fed-queue/6090/7 "2023-10-15T18:48:22Z")

</div>

> [@Blacksmoke16](#):
>
> Feel like [5 use cases for Crystal’s select statement - lbarasti’s blog](https://lbarasti.com/post/select_statement/) might push you in the right direction.

Oh, excellent post. I had a feeling there was something with “select” but I couldn’t find it. Which is obvious now, it’s not really documented…

Pondering this while walking the dogs, I came to the conclusion that it’s just a matter of using non-blocking send and have the main fiber check the queue before receiving. Fibers takes a bit getting used to, my initial thought was “but what if another fiber queues another item in between the main fiber checking the queue and calling `receive`” and then realizing that’s not going to happen (until we bring in real multi threading).

> [@jgaskins](#):
>
> I wrote [a shard](https://github.com/jgaskins/mpsc) for this

Oh, nice.
