Concurrent Data Processing in Elixir - Chapter 3
Published 2026-10-04
3. Data-Processing Pipelines with GenStage
Okay we started to get into it in the last chapter but that was a very simple implementation where we just stopped giving out new processes if the amount of processes was greater than a specific number. Here we will get into a new issue: the amount of data you can process vs. the amount coming in.
Understanding Back-Pressure
You are signing books and you can only sign them so fast. If everyone came in at the same time you would get overwhelmed and get stressed and want to leave early. If however there was an orderly line you would feel much better about signing and maybe want to stay later.
How this relates is that we will be using the GenStage library to build a pipeline to ensure that we only process what we want when we want it.
Introducing GenStage
“GenStage is a new Elixir behaviour for exchanging events with back-pressure between Elixir processes.”
We will not be building processes as much as we will be building Stages. They are meant to receive events and use them to do useful work. Stages are meant to run together, across multiple processes and cores, which means they can branch and be replicated.
You might think that the data flows from left to right, passing along the events, but in a stage the requests (demand) flow from right to left. That means the last stage requests data from the stage before it and will only accept the amount it asks for.
The Producer
The producer is the stage that is producing the events that need to be processed. You might say that the producer publishes events.
The Consumer
Events created by the producer need to be processed by the consumer. The consumer will subscribe to the producer.
The Producer-Consumer
There is sometimes a need for a middleman that can produce-consume. The book talks about how a restaurant is a producer-consumer: it produces food and consumes the produce that comes in.
Building Your Data-Processing Pipeline
Okay let’s start a new project.
mix new scraper --sup
Head to mix.exs and add the dependency for :gen_stage
defp deps do
[
{:gen_stage, "~> 1.0"}
]
end
Next we can build in the dummy code in scraper.ex so that we can simulate scraping a webpage.
defmodule Scraper do
def work() do
# For simplicity, this function is
# just a placeholder and does not contain
# real scraping logic.
1..5
|> Enum.random()
|> :timer.seconds()
|> Process.sleep()
end
end
Creating a Producer
Remember, a producer creates events that we will need to work on. Its job in this case is to produce URLs. Let’s create a new file called page_producer.ex; it will sit in the lib folder.
defmodule PageProducer do
use GenStage
require Logger
def start_link(_args) do
initial_state = []
GenStage.start_link(__MODULE__, initial_state, name: __MODULE__)
end
def init(initial_state) do
Logger.info("PageProducer init")
{:producer, initial_state}
end
def handle_demand(demand, state) do
Logger.info("PageProducer received demand for #{demand} pages")
events = []
{:noreply, events, state}
end
end
This should look familiar, as we have start_link/1, init/1 and some way of dealing with events. Since this is a producer, it uses the handle_demand/2 callback, which receives the number of events being requested and the internal state of the producer.
Creating a Consumer
Now let us create the page_consumer.ex that will take the events that it asks for and do the work on them. It will sit in the lib folder.
defmodule PageConsumer do
use GenStage
require Logger
def start_link(_args) do
initial_state = []
GenStage.start_link(__MODULE__, initial_state)
end
def init(initial_state) do
Logger.info("PageConsumer init")
{:consumer, initial_state, subscribe_to: [PageProducer]}
end
def handle_events(events, _from, state) do
Logger.info("PageConsumer received #{inspect(events)}")
# Pretending that we're scraping web pages.
Enum.each(events, fn _page ->
Scraper.work()
end)
{:noreply, [], state}
end
end
The real big difference here is the init/1 return as it says which stage it will subscribe to. We also have the handle_events/3 that will help us deal with the events that we receive from the producer. Within that function will be the work that we would do but for now we are just going to use the dummy function.
Now let’s head to the application and set up the supervision tree.
@impl true
def start(_type, _args) do
children = [
PageProducer,
PageConsumer
]
opts = [strategy: :one_for_one, name: Scraper.Supervisor]
Supervisor.start_link(children, opts)
end
We can now run the program and see the output.
Generated scraper app
10:17:25.638 [info] PageProducer init
10:17:25.640 [info] PageConsumer init
10:17:25.641 [info] PageProducer received demand for 1000 pages
Wow, a whopping 1000 page demand.
Understanding Consumer Demand
You might be wondering where the 1000 came from: by default a :consumer subscribes with a min_demand of 500 and a max_demand of 1000. You can set these within the init/1 of the consumer. Here is the line that you would add.
sub_opts = [{PageProducer, min_demand: 500, max_demand: 1000}]
{:consumer, initial_state, subscribe_to: sub_opts}
The idea is that once the consumer has fewer than min_demand events left to process, it will ask for more. It can be sent anywhere from 1 event up to max_demand, but never more than that. You can set these values to be anything that you want and what will work for the data load that you are processing. Here are some questions that you might ask while setting the numbers.
Tuning min_demand and max_demand
-
How long does it take to process a single event vs. 100 vs. 1000?
- This tells you whether it’s faster to demand events in batches rather than one by one, and what batch size to use.
-
How often do you expect new events?
-
If the producer emits events often and they’re always available, raising
min_demandmay help. This matters most when working in batches is ideal.
-
If the producer emits events often and they’re always available, raising
-
Can processing be delayed?
-
Consumers wait until the producer has satisfied
min_demand. If processing can’t be delayed, lowermin_demandor set it to0, so a consumer demands each event as soon as the producer has one.
-
Consumers wait until the producer has satisfied
-
What system resources are available to me?
-
You can spread the load by lowering
min_demandandmax_demand, buffering events, or running several consumers (covered later).
-
You can spread the load by lowering
Okay now that we have that all within our minds let’s head back to our project. We should lower the throughput to something that will work for us.
def init(initial_state) do
Logger.info("PageConsumer init")
sub_opts = [{PageProducer, min_demand: 1, max_demand: 3}]
{:consumer, initial_state, subscribe_to: sub_opts}
end
Revisiting the Producer
So the handle_demand/2 works well when we want to handle demand ASAP and keep the consumer busy, but we want to peel back the curtain a bit, so we will use some other callbacks for this project. We could use one of the following:
- handle_call/3
- handle_cast/2
- handle_info/2
Those all will be called in some way when we send info to the page_producer but for the GenStage to work these are the return tuples that are allowed:
-
{:reply, reply, [event], new_state}
-
{:noreply, [event], new_state}
Let’s make this work for our page_producer
def scrape_pages(pages) when is_list(pages) do
GenStage.cast(__MODULE__, {:pages, pages})
end
def handle_cast({:pages, pages}, state) do
{:noreply, pages, state}
end
We have now exposed a public API function for clients to call. You can send a list of pages that you want scraped and it will take care of the backend while only exposing a small fraction of what is needed to run the GenStage. Let’s test it out.
iex(1)> pages = [
...(1)> "google.com",
...(1)> "facebook.com",
...(1)> "apple.com",
...(1)> "netflix.com",
...(1)> "amazon.com"
...(1)> ]
["google.com", "facebook.com", "apple.com", "netflix.com", "amazon.com"]
iex(2)> PageProducer.scrape_pages(pages)
10:35:06.845 [info] PageConsumer received ["google.com", "facebook.com"]
:ok
iex(3)>
nil
10:35:13.849 [info] PageConsumer received ["apple.com"]
10:35:18.850 [info] PageConsumer received ["netflix.com", "amazon.com"]
10:35:24.852 [info] PageProducer received demand for 3 pages
See how the consumer only took a few at a time. Now we can work on adding in more consumers
Adding More Consumers
Okay so now let’s change this up and make the max it can deal with 1 and then add another consumer and see what happens.
sub_opts = [{PageProducer, min_demand: 0, max_demand: 1}]
...
children = [
PageProducer,
Supervisor.child_spec(PageConsumer, id: :consumer_a),
Supervisor.child_spec(PageConsumer, id: :consumer_b)
]
Let’s try that again. Make a .iex.exs file so you don’t have to keep typing in all the pages.
pages = [
"google.com",
"facebook.com",
"apple.com",
"netflix.com",
"amazon.com"
]
PageProducer.scrape_pages(pages)
# Here is the output
10:40:20.298 [info] PageConsumer init
10:40:20.298 [info] PageProducer received demand for 1 pages
10:40:20.298 [info] PageConsumer init
10:40:20.298 [info] PageProducer received demand for 1 pages
Interactive Elixir (1.18.0) - press Ctrl+C to exit (type h() ENTER for help)
10:40:20.312 [info] PageConsumer received ["google.com"]
10:40:20.312 [info] PageConsumer received ["facebook.com"]
10:40:22.314 [info] PageConsumer received ["apple.com"]
10:40:23.315 [info] PageConsumer received ["netflix.com"]
10:40:24.316 [info] PageConsumer received ["amazon.com"]
10:40:25.314 [info] PageProducer received demand for 1 pages
10:40:26.317 [info] PageProducer received demand for 1 pages
iex(1)>
Buffering Events
Producers have a buffer for events that no consumer has asked for yet. The default size is 10_000 events, and if more than that pile up, events get dropped. Let’s set the buffer to 1 and see what happens.
{:producer, initial_state, buffer_size: 1}
# Here is the output.
10:58:44.289 [warning] GenStage producer PageProducer has discarded 2 events from buffer
10:58:44.299 [info] PageConsumer received ["google.com"]
10:58:44.299 [info] PageConsumer received ["facebook.com"]
10:58:46.301 [info] PageConsumer received ["amazon.com"]
10:58:49.301 [info] PageProducer received demand for 1 pages
10:58:51.302 [info] PageProducer received demand for 1 pages
iex(1)>
Adding Concurrency with ConsumerSupervisor
We added another consumer to process the events, but we had to name them ourselves. That isn’t the best way to automate this, which is where ConsumerSupervisor comes in. It lets us deal with more and more consumers automatically.
Creating a ConsumerSupervisor
Let’s create a new file called page_consumer_supervisor.ex that will sit in the lib folder.
defmodule PageConsumerSupervisor do
use ConsumerSupervisor
require Logger
def start_link(_args) do
ConsumerSupervisor.start_link(__MODULE__, :ok)
end
def init(:ok) do
Logger.info("PageConsumerSupervisor init")
children = [
%{
id: PageConsumer,
start: {PageConsumer, :start_link, []},
restart: :transient
}
]
opts = [
strategy: :one_for_one,
subscribe_to: [
{PageProducer, max_demand: 2}
]
]
ConsumerSupervisor.init(children, opts)
end
end
This looks like a lot of new code, but in the end it’s not that much more. We are creating a ConsumerSupervisor that will start a PageConsumer child for each event, with a max demand set in init/1.
The Simplified Consumer
Now that we have the ConsumerSupervisor for the PageConsumer we can simplify the PageConsumer. We need a way to link the PageConsumer to the parent process so the supervisor knows what happens to it. The easiest way is with Task.start_link/1.
defmodule PageConsumer do
require Logger
def start_link(event) do
Logger.info("PageConsumer received #{event}")
Task.start_link(fn ->
Scraper.work()
end)
end
end
If you really think about the structure we just created it all makes sense for what we are trying to accomplish. We could have just used the old functionality but there would be no way to know what happened within the child process. Also make sure to remove the small buffer.
Putting It All Together
Now we need to set up the new application.ex to reflect the current state of the structure.
children = [
PageProducer,
PageConsumerSupervisor
]
# Here is the output.
Interactive Elixir (1.18.0) - press Ctrl+C to exit (type h() ENTER for help)
11:08:36.826 [info] PageConsumer received google.com
11:08:36.826 [info] PageConsumer received facebook.com
11:08:38.829 [info] PageConsumer received apple.com
11:08:41.829 [info] PageConsumer received netflix.com
11:08:41.830 [info] PageConsumer received amazon.com
11:08:42.830 [info] PageProducer received demand for 1 pages
11:08:42.831 [info] PageProducer received demand for 1 pages
Creating Multi-Stage Data Pipelines
Now we can start to talk about pipelines, a way of taking a process and breaking it up into smaller processes and then trying to make them as concurrent as possible.
If your domain has to process the data in multiple steps, you should write that
logic in separate modules and not directly in a GenStage. You only add stages
according to the runtime needs, typically when you need to provide back-pressure
or leverage concurrency.
Start with functions and when you feel the need to deal with back-pressure, create a two-stage data pipeline, then keep going if you need it.
We will build from small to large. The first thing we need is to check if the site is even up. Let’s change scraper.ex.
def online?(_url) do
# Pretend we are checking if the
# service is online or not.
work()
# Select result randomly.
Enum.random([false, true, true])
end
Adding a Producer-Consumer
Let’s add another file called online_page_producer_consumer.ex that will hold the producer-consumer for our new set of data.
defmodule OnlinePageProducerConsumer do
use GenStage
require Logger
def start_link(_args) do
initial_state = []
GenStage.start_link(__MODULE__, initial_state, name: __MODULE__)
end
@impl true
def init(initial_state) do
Logger.info("OnlinePageProducerConsumer init")
subscription = [
{PageProducer, min_demand: 0, max_demand: 1}
]
{:producer_consumer, initial_state, subscribe_to: subscription}
end
@impl true
def handle_events(events, _from, state) do
Logger.info("OnlinePageProducerConsumer received #{inspect(events)}")
events = Enum.filter(events, &Scraper.online?/1)
{:noreply, events, state}
end
end
Okay, so we have the producer-consumer set up with start_link/1, init/1 and handle_events/3. These will help us to take care of the pipeline that we are setting up.
Rewiring Our Pipeline
Let’s head to page_consumer_supervisor.ex and application.ex.
subscribe_to: [
{OnlinePageProducerConsumer, max_demand: 2}
]
children = [
PageProducer,
OnlinePageProducerConsumer,
PageConsumerSupervisor
]
Now let’s try it out.
Generated scraper app
18:10:16.014 [info] PageProducer init
18:10:16.023 [info] OnlinePageProducerConsumer init
18:10:16.023 [info] PageProducer received demand for 1 pages
18:10:16.023 [info] PageConsumerSupervisor init
Interactive Elixir (1.18.0) - press Ctrl+C to exit (type h() ENTER for help)
18:10:16.039 [info] OnlinePageProducerConsumer received ["google.com"]
18:10:21.043 [info] OnlinePageProducerConsumer received ["facebook.com"]
18:10:21.044 [info] PageConsumer received google.com
18:10:25.044 [info] PageConsumer received facebook.com
18:10:25.046 [info] OnlinePageProducerConsumer received ["apple.com"]
18:10:27.046 [info] PageConsumer received apple.com
18:10:27.047 [info] OnlinePageProducerConsumer received ["netflix.com"]
18:10:32.050 [info] PageConsumer received netflix.com
18:10:32.050 [info] OnlinePageProducerConsumer received ["amazon.com"]
18:10:36.051 [info] PageConsumer received amazon.com
18:10:36.051 [info] PageProducer received demand for 1 pages
Scaling Up a Stage with Extra Processes
Now we can add a registry for this with just one more child in application.ex.
children = [
{Registry, keys: :unique, name: ProducerConsumerRegistry},
PageProducer,
OnlinePageProducerConsumer,
PageConsumerSupervisor
]
Next we can add in a new OnlinePageProducerConsumer and it will need a new id for it. So we can add a helper function that will take care of naming the new children.
def producer_consumer_spec(id: id) do
id = "online_page_producer_consumer_#{id}"
Supervisor.child_spec({OnlinePageProducerConsumer, id}, id: id)
end
# Then we can use this to add in more children
children = [
{Registry, keys: :unique, name: ProducerConsumerRegistry},
PageProducer,
producer_consumer_spec(id: 1),
producer_consumer_spec(id: 2),
PageConsumerSupervisor
]
Now that we have the children named we can now use the names within the online_page_producer_consumer.ex
def start_link(id) do
initial_state = []
GenStage.start_link(__MODULE__, initial_state, name: via(id))
end
def via(id) do
{:via, Registry, {ProducerConsumerRegistry, id}}
end
Okay so when we pass in the id of the child to the start_link we will then be able to make sure we are calling them by the right name.
We now need to head to the page_consumer_supervisor.ex and make sure that we are subscribing to the right process by name.
opts = [
strategy: :one_for_one,
subscribe_to: [
{OnlinePageProducerConsumer.via("online_page_producer_consumer_1"), []},
{OnlinePageProducerConsumer.via("online_page_producer_consumer_2"), []}
]
]
We can now test it out.
18:38:04.649 [info] PageProducer init
18:38:04.654 [info] OnlinePageProducerConsumer init
18:38:04.654 [info] PageProducer received demand for 1 pages
18:38:04.654 [info] OnlinePageProducerConsumer init
18:38:04.654 [info] PageProducer received demand for 1 pages
18:38:04.657 [info] PageConsumerSupervisor init
Interactive Elixir (1.18.0) - press Ctrl+C to exit (type h() ENTER for help)
18:38:04.672 [info] OnlinePageProducerConsumer received ["google.com"]
18:38:04.672 [info] OnlinePageProducerConsumer received ["facebook.com"]
18:38:07.677 [info] OnlinePageProducerConsumer received ["apple.com"]
18:38:07.678 [info] PageConsumer received google.com
18:38:09.691 [info] PageConsumer received facebook.com
18:38:09.691 [info] OnlinePageProducerConsumer received ["netflix.com"]
18:38:10.681 [info] PageConsumer received apple.com
18:38:10.681 [info] OnlinePageProducerConsumer received ["amazon.com"]
18:38:10.692 [info] PageConsumer received netflix.com
18:38:10.692 [info] PageProducer received demand for 1 pages
18:38:14.686 [info] PageProducer received demand for 1 pages
Choosing the Right Dispatcher
The last bit of this chapter covers the dispatcher, which decides how events are sent to consumers. When we use :producer and :producer_consumer we get the default DemandDispatcher, but there are other types.
What we currently use when we init/1 is equivalent to:
def init(initial_state) do
Logger.info("PageProducer init")
{:producer, initial_state}
end
# Is equal to
def init(state) do
{:producer, state, dispatcher: GenStage.DemandDispatcher}
end
GenStage ships with two other dispatchers.
Using BroadcastDispatcher
{:producer, state, dispatcher: GenStage.BroadcastDispatcher}
This is used when we want all the consumers subscribed to the producer to get every event, even if they each use it differently. With that being said, a subscriber can pass a selector to opt in or out of individual events.
def init(state) do
selector =
fn incoming_event ->
nil
# you can use the event to decide whether
# to return `true` and accept it, or `false` to reject it.
end
sub_opts = [
{SomeProducer, selector: selector}
]
{:consumer, state, subscribe_to: sub_opts}
end
Using PartitionDispatcher
This is the way in which the producer will partition the events into different buckets, and it’s the consumer that will be responsible for grabbing the events from those buckets.
def init(state) do
hash =
fn event ->
# you can use the event to decide which partition
# to assign it to, or use `:none` to ignore it.
{event, :c}
end
opts = [
partitions: [:a, :b, :c],
hash: hash
]
{:producer, state, dispatcher: {GenStage.PartitionDispatcher, opts}}
end
sub_opts = [
{SomeProducer, partition: :b}
]
{:consumer, state, subscribe_to: sub_opts}
Wrapping Up
We now have GenStage, which allows us to better spread out the work. As of right now we are creating the children as needed and nothing is there to completely automate the process. We can work on that later but let’s talk about what we learned.
The rate at which events come in might not be consistent, so we need to deal with them at the rate we choose.
First we took the normal GenServer and then changed it into a GenStage. That allowed us to talk about the Producer/Consumer relationship. Once we had that, we worked on the Producer-Consumer that can do both production and consuming.
After that we talked about having a supervisor for the consumers, and then we added another stage that checks whether a site is online. Once that was done, we scaled that stage up with extra processes. You can see how powerful this can become to take a process and break it up into the parts that we want and make sure that each stage handles events at the rate it can.