Backpressure in Elixir Pipelines with GenStage
This post explains how Elixir’s GenStage keeps streaming pipelines bounded with backpressure, walks through a file‑processing example, and lists trade‑offs and verification steps.
27 Apr 2026, 08:36 UTC

Problem: Unbounded Queues in Streaming Pipelines
When a fast producer feeds data into a slow consumer, the default message‑passing model in Elixir can cause the consumer’s mailbox to grow without bound. The system may eventually exhaust memory, leading to crashes or severe performance degradation. A common pattern for solving this is backpressure: the consumer signals how many items it can handle, and the producer only sends that many.
GenStage Overview: Demand‑Driven Flow Control
GenStage is built on top of Elixir processes and the :gen_event behaviour. It introduces three core behaviours:
GenStage.producer– emits events in response to demand.GenStage.consumer– requests a specific number of events.GenStage.subscriber– a hybrid that can both produce and consume.
The key idea is that each stage declares a demand window. When a consumer requests n events, the producer receives a {:demand, n} message and can safely emit up to n events. If the producer cannot keep up, it simply stops emitting until more demand arrives.
Building a Backpressure‑Aware Pipeline
Below is a minimal three‑stage pipeline that demonstrates GenStage’s backpressure in action:
- Producer – reads lines from a file lazily.
- Processor – transforms each line to uppercase.
- Consumer – writes the transformed line to
stdoutand simulates a slow consumer with a small delay.
1. Define the Producer
defmodule FileLineProducer do
use GenStage
def start_link(file_path) do
GenStage.start_link(__MODULE__, file_path, name: __MODULE__)
end
def init(file_path) do
{:ok, file} = File.open(file_path, [:read])
{:producer, %{file: file, demand: 0}}
end
def handle_demand(demand, state) do
lines = read_lines(state.file, demand)
{:noreply, lines, %{state | demand: 0}}
end
defp read_lines(_file, 0), do: []
defp read_lines(file, n) do
case IO.read(file, :line) do
:eof -> []
line -> [line | read_lines(file, n - 1)]
end
end
end
2. Define the Processor
defmodule UppercaseProcessor do
use GenStage
def start_link(opts) do
GenStage.start_link(__MODULE__, opts)
end
def init(_opts) do
{:processor, %{demand: 0}}
end
def handle_events(events, _from, state) do
processed = Enum.map(events, &String.upcase/1)
{:noreply, processed, state}
end
end
3. Define the Consumer
defmodule StdoutConsumer do
use GenStage
def start_link(opts) do
GenStage.start_link(__MODULE__, opts)
end
def init(_opts) do
{:consumer, %{demand: 0}}
end
def handle_events(events, _from, state) do
# Simulate a slow consumer
Enum.each(events, fn line ->
IO.puts(line)
:timer.sleep(100) # 100 ms per line
end)
{:noreply, [], state}
end
end
4. Wire the Pipeline
defmodule Pipeline do
def start(file_path) do
{:ok, prod} = FileLineProducer.start_link(file_path)
{:ok, proc} = UppercaseProcessor.start_link([])
{:ok, cons} = StdoutConsumer.start_link([])
# Subscribe in the correct order
GenStage.sync_subscribe(proc, to: prod)
GenStage.sync_subscribe(cons, to: proc, min_demand: 1, max_demand: 5)
end
end
Running the Example
Assuming the code is in a Mix project, run:
mix run -e "Pipeline.start(\"data.txt\")" --no-halt
Replace data.txt with a path to a text file. The consumer will print one line every 100 ms, and the producer will automatically slow down because it only emits new lines when the consumer’s demand queue is not full.
Trade‑offs & Limitations
- Latency – Each stage adds a context switch. For very low‑latency use‑cases, the overhead can be noticeable.
- Demand Configuration – Setting
min_demandormax_demandincorrectly can stall the pipeline. A demand of zero will freeze the producer. - Complex Back‑pressure Paths – In multi‑consumer scenarios, balancing demand across branches can become tricky.
- Elixir Version – GenStage is part of the core library since Elixir 1.11. Older releases require adding the external
gen_stagedependency and may have subtle behavioural differences.
How to Verify Backpressure Works
- Observe Output Rate – With the 100 ms consumer delay, the output should appear at roughly 10 lines per second. If you remove the delay, the consumer will keep up and the producer will run at file‑read speed.
- Introduce a Faster Consumer – Change
:timer.sleep(100)to0and note that the pipeline now processes lines as fast as the file can be read, but the mailbox size remains bounded. - Check Mailbox Size – In an IEx session, run
Process.info(pid, :message_queue_len)for each stage. The values should stay close to themax_demandsetting. - Use :observer – Start
:observer.startand inspect the process tree. Each stage should appear as a separate GenServer with a bounded message queue. - Stress Test – Feed a very large file and monitor memory usage. The system should not crash or run out of memory if backpressure is correctly wired.
Actionable Takeaway
When you need a streaming pipeline that can gracefully handle bursty input or slow consumers, GenStage’s demand‑driven model is a solid choice. Implement each stage as a GenStage module, set sensible min_demand and max_demand values, and test with varying consumer speeds. The built‑in backpressure will keep your system stable without manual queue size checks.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.