~/adam.log

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_demand may help. This matters most when working in batches is ideal.
  • Can processing be delayed?

    • Consumers wait until the producer has satisfied min_demand. If processing can’t be delayed, lower min_demand or set it to 0, so a consumer demands each event as soon as the producer has one.
  • What system resources are available to me?

    • You can spread the load by lowering min_demand and max_demand, buffering events, or running several consumers (covered later).

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.