mmmmillar
Problem:
As a first go at using GenStage, I’m building a simple task queue. Initially it used consumers that I would add explicitly in my Application supervisor. I’ve since updated it to use a ConsumerSupervisor to add/remove Consumers as required.
The issue I’m having is that every task added to the producer gets picked up by the ProducerConsumer and passed on, but about half of them seem to disappear into thin air (ie I only see the output of IO.inspect("#{job_id} is starting") for about half of them)!
Does anyone know why this is happening or how I can debug this further?
Thanks
Code
producer.ex
defmodule Q.Producer do
use GenStage
require Logger
def start_link(_init_args) do
GenStage.start_link(__MODULE__, {:queue.new(), 0}, name: __MODULE__)
end
@impl true
def init(initial) do
{:producer, initial}
end
@impl true
def handle_demand(demand, {backlog, existing_demand}) do
# IO.inspect("Received demand: #{demand}, existing_demand: #{existing_demand}")
case :queue.len(backlog) do
0 ->
{:noreply, [], {backlog, existing_demand + demand}}
n ->
n = min(demand, n)
{items, backlog} = :queue.split(n, backlog)
:queue.len(backlog) |> Q.Stats.set_waiting()
{:noreply, :queue.to_list(items), {backlog, existing_demand + demand - n}}
end
end
@impl true
def handle_cast({:enqueue, item}, {backlog, 0}) do
{:noreply, [], {:queue.in(item, backlog), 0}}
end
@impl true
def handle_cast({:enqueue, item}, {backlog, existing_demand}) do
IO.inspect("Received item: existing_demand: #{existing_demand}")
backlog = :queue.in(item, backlog)
{{:value, item}, backlog} = :queue.out(backlog)
{:noreply, [item], {backlog, existing_demand - 1}}
end
def enqueue(item), do: GenStage.cast(__MODULE__, {:enqueue, item})
end
producer_consumer.ex
defmodule Q.ProducerConsumer do
use GenStage
def start_link(_init_args) do
GenStage.start_link(__MODULE__, :ok, name: __MODULE__)
end
def init(initial) do
{:producer_consumer, initial, subscribe_to: [{Q.Producer, max_demand: 1}]}
end
def handle_events(events, _from, state) do
# producer "middleware" - do things like filter before passing on to consumer
IO.inspect("passing events to consumer: #{inspect(events)}")
{:noreply, events, state}
end
end
consumer_supervisor.ex
defmodule Q.ConsumerSupervisor do
use ConsumerSupervisor
def start_link(_args) do
{:ok, pid} = ConsumerSupervisor.start_link(__MODULE__, :ok, name: __MODULE__)
{:ok, pid}
end
def init(:ok) do
children = [
%{
id: Q.Consumer,
start: {Q.Consumer, :start_link, []},
restart: :transient
}
]
ConsumerSupervisor.init(children,
strategy: :one_for_one,
subscribe_to: [
{Q.ProducerConsumer, max_demand: 5}
]
)
end
end
consumer.ex
defmodule Q.Consumer do
alias Q.JobRecord
use GenStage
import Q.Constants
@max_job_duration max_job_duration()
def start_link(_init_args) do
GenStage.start_link(__MODULE__, :ok)
end
def init(initial) do
Process.flag(:trap_exit, true)
Q.Stats.increment_consumer_count()
{:consumer, initial, subscribe_to: [{Q.ProducerConsumer, max_demand: 1}]}
end
def handle_events(events, _from, state) do
Enum.each(events, fn job_id ->
task =
Task.async(fn ->
IO.inspect("#{job_id} is starting")
JobRecord.set_started(job_id)
run_job()
end)
case Task.yield(task, @max_job_duration) || Task.shutdown(task) do
{:ok, _result} ->
JobRecord.set_completed(job_id)
nil ->
JobRecord.retry_job(job_id)
end
end)
# As a consumer we never emit events
{:noreply, [], state}
end
def handle_info({:EXIT, _pid, reason}, state) do
Q.Stats.decrement_consumer_count()
{:stop, reason, state}
end
defp run_job do
Process.sleep(100)
end
end
Trending in Questions
I having some trouble figuring out if I have set myself too strict of standards for my production server. Currently I can handle 75% of r...
New
Hello,
I’m trying to build a basic Phoenix web-app, and I’d like to use Tailwind.
However, when I launch mix phx.server, I get an error...
New
Hi everyone.
My team and I have been working on a fairly modest app based around video streaming and chat, but we’ve landed a customer t...
New
I really like the adapter patterns that ecto, nebulex, waffle, etc. use and would love find something similar for a key management servic...
New
Hello folks!
So at work, we are seeing some situations where we have to define some “fixed” strings that are used across the codebase in...
New
I’m working on a small exercise involving update_in/3, and I came up with this solution:
data = %{
name: "Periodic Table",
category:...
New
How Can I Optimise Compile Time Dependencies
I have been building an elixir application for about 2 years now. Many modules and files ha...
New
Other Trending Topics
Edit: 2026 May 15 - This post is archived.
Mob is alive!!
Main docs: mob v0.7.11 — Documentation
A bit of explanation for the slightly c...
New
Hobbes is a low-level distributed database for the Elixir programming language.
Hobbes provides a simple, safe, and scalable storage lay...
New
A little off-topic, but I feel like people here have a good head on their shoulders.
I used to be quite good at making software. Was luc...
New
Hey. Is there anyone here who creates agents in their apps? Not talking about using agents, but creating them. I’m finding it pretty diff...
New
I fully migrated to my own harness from Anthropic/Gemini and I think it’s time to share it. Welcome DSH, the DeepSeek Harness, fully writ...
New
ExRatatui lets you cook up rich terminal UIs in Elixir, powered by Rust’s ratatui via Rustler NIFs. Build interactive terminal applicatio...
New
Categories:
Sub Categories:
Forums
Popular Tags
- #ecto
- #liveview
- #troubleshooting
- #learning-elixir
- #library
- #deployment
- #erlang
- #testing
- #genserver
- #mix
- #absinthe
- #remote-other
- #otp
- #plug
- #how-to-question
- #macros
- #postgres
- #elixirconf
- #channels
- #exunit
- #discussion
- #code-sync
- #podcasts
- #javascript
- #onsite
- #dialyzer
- #docker
- #authentication
- #umbrella
- #ai
- #full-time-contract
- #podcasts-by-brainlid
- #ecto-query
- #blog-post
- #elixirconf-us
- #elixir-ls
- #phoenix_html
- #iex
- #graphql
- #genstage
- #websockets
- #supervisor
- #advent-of-code
- #distillery
- #processes
- #elixirconf-eu
- #api
- #forms
- #security
- #metaprogramming











Showing Posts 1 to 1- Show Best Posts
- Show All (oldest first)
- Show All (newest first)
mmmmillar
My problem was that I was still using GenStage for the “consumer” when in fact I just needed to start a task (as the consumer supervisor was now picking up the events)