Concurrent Data Processing in Elixir - Chapter 2
Published 2026-10-03
2. Long-Running Processes Using GenServer
So what we have learned so far is a quick and easy way to build a concurrent set of instruction and use functions. These small sets of processes deal with trash clean-up and don’t require much on the server side.
What happens when we need a more robust and longer lasting structure? Well that is where GenServer comes in. We will start to work with GenServers in this chapter.
Starting with a Basic GenServer
We will continue to work with the sender project for this next set of code. First create a new file in the lib folder called send_server.ex we can then define the module like this.
defmodule SendServer do
use GenServer
end
We can start the program with the normal start we used before
iex -S mix
With the output
Erlang/OTP 27 [erts-15.1.2] [source] [64-bit] [smp:24:24] [ds:24:24:10] [async-threads:1] [jit:ns]
Compiling 1 file (.ex)
warning: function init/1 required by behaviour GenServer is not implemented (in module SendServer).
We will inject a default implementation for now:
def init(init_arg) do
{:ok, init_arg}
end
You can copy the implementation above or define your own that converts the arguments given to GenServer.start_link/3 to the server state.
│
1 │ defmodule SendServer do
│ ~~~~~~~~~~~~~~~~~~~~~~~
│
└─ lib/send_server.ex:1: SendServer (module)
Generated sender app
Interactive Elixir (1.18.0) - press Ctrl+C to exit (type h() ENTER for help)
iex(1)>
You did it you have the first GenServer
GenServer Callbacks In Depth
So we will now need to talk about callbacks that will be used to do all the work within the GenServer. there are defaults that we will be overriding to get the functionality that we want. In order to override something we need to know:
- What arguments the callback takes
- What the return values that are supported
We will cover the following within this chapter:
- handle_call/3
- handle_cast/2
- handle_continue/2
- handle_info/2
- init/1
- terminate/2
Learning just these will get you very far into understanding how the GenServer functions.
Initializing the Process
Let’s start with init/1 this is callback that is started when you first start the process. We will now start to make it so the SendServer can see faux emails.
@impl true
def init(args) do
IO.puts("Received arguments: #{inspect(args)}")
max_retries = Keyword.get(args, :max_retries, 5)
state = %{emails: [], max_retries: max_retries}
{:ok, state}
end
We set up some standard setting within the init and we can now use it as we need. Head back to your iex session and test it out.
iex(3)> {:ok, pid} = GenServer.start(SendServer, [max_retries: 1])
Received arguments: [max_retries: 1]
{:ok, #PID<0.187.0>}
iex(4)> GenServer.stop(pid)
:ok
There are some standard returns for the init some common ones are the following:
}
:ignore
We covered the but the is great for doing post-initialization work.
Lets now use the handle_continue/2 to deal with some post init work to get more information. As its not good to do much besides set the initial state of the process within the init/1.
@impl true
def handle_continue(:fetch_from_database, state) do
# called after init/1
end
Breaking Down Work in Multiple Steps
Okay so some of the return values for handle_continue/2 are:
}
We can use one of those now with this implementation. It will assume that we have a DB that has a users
@impl true
def handle_continue(:fetch_from_database, state) do
users = []
# get `users` from the database
{:noreply, Map.put(state, :users, users)}
end
Sending Process Messages
One of the nice things about this GenServer is that you can send and receive messages from the process while its running that is where GenServer.call/3 and GenServer.cast/2 will come in. You can call and expect a reply and you can cast and send info without a need for a reply. Let’s add in some more code.
@impl true
def handle_call(:get_state, _from, state) do
{:reply, state, state}
end
We have said that we want to get a reply and set the state as the reply, notice that you also need to send the state in the return meaning that you can do work on the state within a call. Here are some of the normal returns:
}
Now we can recompile the iex session and test it out.
iex(7)> {:ok, pid} = GenServer.start(SendServer, [max_retries: 1])
Received arguments: [max_retries: 1]
{:ok, #PID<0.214.0>}
iex(8)> GenServer.call(pid, :get_state)
%{emails: [], max_retries: 1}
Okay that is great now we can implement sending emails with the handle_cast/2
def handle_cast({:send, email}, state) do
Sender.send_email(email)
emails = [%{email: email, status: "sent", retries: 0}] ++ state.emails
{:noreply, %{state | emails: emails}}
end
This assumes that you remember that we implemented the Sender.send_email earlier in the book. We can now try that out.
iex(10)> {:ok, pid} = GenServer.start(SendServer, [max_retries: 1])
Received arguments: [max_retries: 1]
{:ok, #PID<0.232.0>}
iex(11)> GenServer.cast(pid, {:send, "hello@email.com"})
:ok
iex(12)> GenServer.call(pid, :get_state)
Email to hello@email.com sent
%{
emails: [%{status: "sent", email: "hello@email.com", retries: 0}],
max_retries: 1
}
Notifying the Process of Events
Now we can talk about an other callback that we can implement the handle_info/2 this is meant to deal with times that you get new information. We will use the Process.send/2 to send info to the Process, it will then trigger the handle_info/2 callback. Let’s work with it now. We can even make it have to work through its retries. Augment the send_email/1 logic to the following and we can augment the handle_cast/2 as well.
def send_email("konnichiwa@world.com" = email),
do: :error
def send_email(email) do
Process.sleep(3000)
IO.puts("Email to #{email} sent")
{:ok, "email_sent"}
end
@impl true
def handle_cast({:send, email}, state) do
status =
case Sender.send_email(email) do
{:ok, "email_sent"} -> "sent"
:error -> "failed"
end
emails = [%{email: email, status: status, retries: 0}] ++ state.emails
{:noreply, %{state | emails: emails}}
end
Okay so we have a way to deal with errors and a way to retry. Let’s add a line to the init/1 so that we try to send out an email asap.
Process.send_after(self(), :retry, 5000)
Now let’s setup the handle_info/2
@impl true
def handle_info(:retry, state) do
{failed, done} =
Enum.split_with(state.emails, fn item ->
item.status == "failed" && item.retries < state.max_retries
end)
retried =
Enum.map(failed, fn item ->
IO.puts("Retrying email #{item.email}...")
new_status =
case Sender.send_email(item.email) do
{:ok, "email_sent"} -> "sent"
:error -> "failed"
end
%{email: item.email, status: new_status, retries: item.retries + 1}
end)
Process.send_after(self(), :retry, 5000)
{:noreply, %{state | emails: retried ++ done}}
end
We now have a way of retrying the emails after we start. It will continue to do so over and over as in the beginning it will send a retry, then end of the handle_info/2 it will send out an other :retry. Let’s test it out in the iex session.
iex(15)> {:ok, pid} = GenServer.start(SendServer, max_retries: 2)
Received arguments: [max_retries: 2]
{:ok, #PID<0.257.0>}
iex(16)> GenServer.cast(pid, {:send, "hello@world.com"})
:ok
iex(17)> GenServer.cast(pid, {:send, "aloha@world.com"})
:ok
iex(18)> GenServer.cast(pid, {:send, "konnichiwa@world.com"})
:ok
Email to aloha@world.com sent
iex(20)> GenServer.call(pid, :get_state)
%{
emails: [
%{status: "failed", email: "konnichiwa@world.com", retries: 0},
%{status: "sent", email: "aloha@world.com", retries: 0},
%{status: "sent", email: "hello@world.com", retries: 0}
],
max_retries: 2
}
Process Teardown
Okay there is an other callback that we should talk about and its the terminate it is what is run when you kill a process. This is a small bit of code but add it to the Module.
@impl true
def terminate(reason, _state) do
IO.puts("Terminating with reason #{reason}")
end
Let’s test it out.
iex(22)> {:ok, pid} = GenServer.start(SendServer, [])
Received arguments: []
{:ok, #PID<0.276.0>}
iex(23)> GenServer.stop(pid)
Terminating with reason normal
:ok
Building a Job-Processing System
Okay so we want to build a process server that will scale to make it run a process for every email that we want to send, with all the functionality that we had already the retries and so on. We will start a new project for this so let’s get started.
mix new jobber --sup
We can set up a new file job.ex in lib/jobber with the following:
defmodule Jobber.Job do
use GenServer
require Logger
end
Initializing the Job Process
Lets set-up the init/1 for the new job.ex
defstruct [:work, :id, :max_retries, retries: 0, status: "new"]
def init(args) do
work = Keyword.fetch!(args, :work)
id = Keyword.get(args, :id, random_job_id())
max_retries = Keyword.get(args, :max_retries, 3)
state = %Jobber.Job{id: id, work: work, max_retries: max_retries}
{:ok, state, {:continue, :run}}
end
defp random_job_id() do
:crypto.strong_rand_bytes(5) |> Base.url_encode64(padding: false)
end
We want to have a unique :id for each job so we will use the :crypto to do so. Because we added the need for :crypto we need to add it into the mix.exs
def application do
[
extra_applications: [:logger, :crypto],
mod: {Jobber.Application, []}
]
end
Performing Work
We set the return to be } so we will need to handle that continue right away in-order to finish the init/1
@impl true
def handle_continue(:run, state) do
new_state = state.work.() |> handle_job_result(state)
if new_state.status == "errored" do
Process.send_after(self(), :retry, 5000)
{:noreply, new_state}
else
Logger.info("Job exiting #{state.id}")
{:stop, :normal, new_state}
end
end
Okay so now we need a way to deal with the different states of the work. So we will need to define some private functions with pattern matching. Based off these states:
- Success, when the job completes and returns {:ok, data}
- Initial error, when it fails the first time with :error
- Retry error, when we attempt to rerun the job and also receive :error
defp handle_job_result({:ok, _data}, state) do
Logger.info("Job completed #{state.id}")
%Jobber.Job{state | status: "done"}
end
defp handle_job_result(:error, %{status: "new"} = state) do
Logger.warn("Job errored #{state.id}")
%Jobber.Job{state | status: "errored"}
end
defp handle_job_result(:error, %{status: "errored"} = state) do
Logger.warn("Job retry failed #{state.id}")
new_state = %Jobber.Job{state | retries: state.retries + 1}
if new_state.retries == state.max_retries do
%Jobber.Job{new_state | status: "failed"}
else
new_state
end
end
Now we can handle the handle_info/2 for the :retry message that will be sent.
@impl true
def handle_info(:retry, state) do
# Delegate work to the `handle_continue/2` callback.
{:noreply, state, {:continue, :run}}
end
You should be able to see that we will go back to the continue branch of the logic and that will start it over again. Let’s test it out.
iex(1)> GenServer.start(Jobber.Job, work: fn -> Process.sleep(5000); {:ok, []} end)
{:ok, #PID<0.188.0>}
iex(2)>
nil
iex(3)> good_job = fn ->
...(3)> Process.sleep(5000)
...(3)> {:ok, []}
...(3)> end
#Function<43.39164016/0 in :erl_eval.expr/6>
iex(4)>
nil
iex(5)> GenServer.start(Jobber.Job, work: good_job)
{:ok, #PID<0.189.0>}
iex(6)>
nil
iex(7)> bad_job = fn ->
...(7)> Process.sleep(5000)
...(7)> :error
...(7)> end
#Function<43.39164016/0 in :erl_eval.expr/6>
iex(8)>
nil
iex(9)> GenServer.start(Jobber.Job, work: bad_job)
{:ok, #PID<0.190.0>}
iex(10)>
nil
19:35:56.155 [info] Job completed uOdxGQg
19:35:56.155 [warning] Job errored cQ0pp-s
19:35:56.158 [info] Job exiting uOdxGQg
19:35:56.155 [info] Job completed -xPCT84
19:35:56.158 [info] Job exiting -xPCT84
Introducing DynamicSupervisor
Okay so before we had everything did start again but it needed to be started with a GenServer.start, that starts the process without a link-process that will die with the child and the process will be lost. There is a better way and that is the DynamicSupervisor. We will need to head into the application.ex in-order to create the supervisor tree for this.
@impl true
def start(_type, _args) do
children = [
{DynamicSupervisor, strategy: :one_for_one, name: Jobber.JobRunner}
]
# See https://hexdocs.pm/elixir/Supervisor.html
# for other strategies and supported options
opts = [strategy: :one_for_one, name: Jobber.Supervisor]
Supervisor.start_link(children, opts)
end
end
Let’s now create the .iex.exs so that we don’t have to add in the good and bad jobs.
good_job = fn ->
Process.sleep(5000)
{:ok, []}
end
bad_job = fn ->
Process.sleep(5000)
:error
end
Okay so we have the functions now let’s start to work with the jobber.ex this will be what we use to start jobs within the DynamicSupervisor.
#jobber/lib/jobber.ex
defmodule Jobber do
alias Jobber.{JobRunner, Job}
def start_job(args) do
DynamicSupervisor.start_child(JobRunner, {Job, args})
end
end
Let’s test it.
iex(1)> Jobber.start_job(work: good_job)
{:error,
{:undef,
[
{Jobber.Job, :start_link,
[[work: #Function<43.39164016/0 in :erl_eval.expr/6>]], []},
{DynamicSupervisor, :start_child, 3,
[file: ~c"lib/dynamic_supervisor.ex", line: 800]},
{DynamicSupervisor, :handle_start_child, 2,
[file: ~c"lib/dynamic_supervisor.ex", line: 786]},
{:gen_server, :try_handle_call, 4, [file: ~c"gen_server.erl", line: 2381]},
{:gen_server, :handle_msg, 6, [file: ~c"gen_server.erl", line: 2410]},
{:proc_lib, :init_p_do_apply, 3, [file: ~c"proc_lib.erl", line: 329]}
]}}
Looks like we don’t have a start_link/2 let’s get that defined.
#jobber/lib/jobber/job.ex
def start_link(args) do
GenServer.start_link(__MODULE__, args)
end
Now it works.
iex(4)> Jobber.start_job(work: good_job)
{:ok, #PID<0.182.0>}
Process Restart Values Revisited
There are ways to make it so it will not restart if the process exits normally. Otherwise it will restart forever. Replace the macro for the GenServer with
use GenServer, restart: :transient
Now lets test the 2 different jobs.
iex(1)> Jobber.start_job(work: good_job)
{:ok, #PID<0.164.0>}
19:51:13.954 [info] Job completed og7u9bA
19:51:13.956 [info] Job exiting og7u9bA
iex(2)> Jobber.start_job(work: bad_job)
{:ok, #PID<0.165.0>}
19:52:15.013 [warning] Job errored b16AJNU
19:52:25.017 [warning] Job retry failed b16AJNU
19:52:32.155 [warning] Job retry failed b16AJNU
19:52:42.167 [warning] Job retry failed b16AJNU
19:52:42.167 [info] Job exiting b16AJNU
That all worked as expected but what happens when we have a doomed job? Add the following to the .iex.exs
doomed_job = fn ->
Process.sleep(5000)
raise "Boom!"
end
iex(1)> Jobber.start_job(work: doomed_job)
{:ok, #PID<0.150.0>}
19:54:29.974 [error] GenServer #PID<0.150.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "FjeMdjg", max_retries: 3, retries: 0, status: "new"}
19:54:35.014 [error] GenServer #PID<0.151.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "60TqQzk", max_retries: 3, retries: 0, status: "new"}
Adjusting Restart Frequency
Right now we have the doomed job because the timeout is longer than the frequency can handle. With the max_restarts set to 3 and the max_seconds set to 5. We can change those. Go to the application.ex and add this extra setting.
@impl true
def start(_type, _args) do
job_runner_config = [
strategy: :one_for_one,
max_seconds: 30,
name: Jobber.JobRunner
]
children = [
{DynamicSupervisor, job_runner_config}
]
opts = [strategy: :one_for_one, name: Jobber.Supervisor]
Supervisor.start_link(children, opts)
end
Lets test it out and keep track of the pids.
iex(1)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
iex(2)> Jobber.start_job(work: doomed_job)
{:ok, #PID<0.163.0>}
iex(3)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
19:59:06.705 [error] GenServer #PID<0.163.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "6U5INHM", max_retries: 3, retries: 0, status: "new"}
19:59:11.731 [error] GenServer #PID<0.164.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "K_c4Xmg", max_retries: 3, retries: 0, status: "new"}
iex(4)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
19:59:16.740 [error] GenServer #PID<0.165.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "jUeDQis", max_retries: 3, retries: 0, status: "new"}
iex(5)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
19:59:21.747 [error] GenServer #PID<0.166.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "-FI_W5E", max_retries: 3, retries: 0, status: "new"}
iex(6)> Process.whereis(Jobber.JobRunner)
#PID<0.167.0>
Looking at the output we see that when a process fails the attached process is killed and we are left with a new Jobber. We can make sure that the same process is kept with a Supervisor.
Implementing a Supervisor
Both of the 2 Supervisors that we have reached for are good for simple trees but we want to have more control over the Supervisors and the children. So now we can reach for The Supervisor module. We can start by creating a new file job_supervisor.ex.
defmodule Jobber.JobSupervisor do
use Supervisor, restart: :temporary
def start_link(args) do
Supervisor.start_link(__MODULE__, args)
end
def init(args) do
children = [
{Jobber.Job, args}
]
options = [
strategy: :one_for_one,
max_seconds: 30
]
Supervisor.init(children, options)
end
end
Now that we have the new Supervisor we can change the Jobber to use the new Supervisor
defmodule Jobber do
alias Jobber.{JobRunner, Job, JobSupervisor}
def start_job(args) do
DynamicSupervisor.start_child(JobRunner, {JobSupervisor, args})
end
end
Let’s test it out.
iex(1)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
iex(2)> Jobber.start_job(work: doomed_job)
{:ok, #PID<0.163.0>}
20:10:12.559 [error] GenServer #PID<0.163.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "mvdyZfs", max_retries: 3, retries: 0, status: "new"}
iex(3)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
20:10:17.601 [error] GenServer #PID<0.164.0> terminating
** (RuntimeError) Boom!
(elixir 1.18.0) src/elixir.erl:386: :elixir.eval_external_handler/3
(jobber 0.1.0) lib/jobber/job.ex:22: Jobber.Job.handle_continue/2
(stdlib 6.1.2) gen_server.erl:2335: :gen_server.try_handle_continue/3
(stdlib 6.1.2) gen_server.erl:2244: :gen_server.loop/7
Last message: {:continue, :run}
State: %Jobber.Job{work: #Function<43.39164016/0 in :erl_eval.expr/6>, id: "jJD80-Y", max_retries: 3, retries: 0, status: "new"}
iex(4)> Process.whereis(Jobber.JobRunner)
#PID<0.161.0>
See the same PID?
Okay so understand what we have here there are 3 types of restart structure that you can use:
- :one_for_one
- :one_for_all
- :rest_for_one
Depending on which one you use, everything will be restarted, nothing will be restarted or pieces will be restarted.
Naming Processes Using the Registry
Okay so one of the most powerful options is the ability to name a process to call it by name, we have been using the pid identifier so far but we can change that.
children = [
{DynamicSupervisor, strategy: :one_for_one, name: Jobber.JobRunner},
]
Typically, the :name is set when we link a process using start_link/1. There are three types of accepted values when naming a process:
- An atom, like :job_runner. This includes module names, since they’re atoms under the hood, for example, Jobber.JobRunner.
- A {:global, term} tuple, like {:global, :job_runner}, which registers the process globally. Useful for distributed applications
- A {:via, module, term} tuple, where module is an Elixir module that would take care of the registration process, using the value term
Starting a Registry Process
This is actually a very simple process head back to the application.ex and make the following changes.
@impl true
def start(_type, _args) do
job_runner_config = [
strategy: :one_for_one,
max_seconds: 30,
name: Jobber.JobRunner
]
children = [
# This next line
{Registry, keys: :unique, name: Jobber.JobRegistry},
{DynamicSupervisor, job_runner_config}
]
opts = [strategy: :one_for_one, name: Jobber.Supervisor]
Supervisor.start_link(children, opts)
end
Registering New Processes
Okay so we can now use the Registry to name a process, lets head to the job.ex and make a few changes starting by adding in a via helper.
defp via(key, value) do
{:via, Registry, {Jobber.JobRegistry, key, value}}
end
def start_link(args) do
args =
if Keyword.has_key?(args, :id) do
args
else
Keyword.put(args, :id, random_job_id())
end
id = Keyword.get(args, :id)
type = Keyword.get(args, :type)
GenServer.start_link(__MODULE__, args, name: via(id, type))
end
@impl true
def init(args) do
work = Keyword.fetch!(args, :work)
# This line next
id = Keyword.get(args, :id)
max_retries = Keyword.get(args, :max_retries, 3)
state = %Jobber.Job{id: id, work: work, max_retries: max_retries}
{:ok, state, {:continue, :run}}
end
Querying the Registry
Registry allows you to use the select/2 it might look cryptic but once you get to understand it you should be fine.
def running_imports() do
match_all = {:"$1", :"$2", :"$3"}
guards = [{:==, :"$3", "import"}]
map_result = [%{id: :"$1", pid: :"$2", type: :"$3"}]
Registry.select(Jobber.JobRegistry, [{match_all, guards, map_result}])
end
The match specification we have is a bit lengthy, so we have broken it down to three variables:
- match_all is a wildcard that matches all entries in the registry.
- guards is a filter that filters results by the third element in the tuple, which has to be equal to “import”.
- map_result is transforming the result by creating a list of maps, assigning each element of the tuple to a key, which makes the result a bit more readable.
Now go to the .iex.exs and increase the sleep for the good job.
good_job = fn ->
Process.sleep(60_000)
{:ok, []}
end
So now we have a way of querying the Registry. Let’s test it out.
iex(1)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.178.0>}
iex(2)> Jobber.start_job(work: good_job, type: "send_email")
{:ok, #PID<0.179.0>}
iex(3)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.180.0>}
iex(4)> Jobber.running_imports()
[
%{id: "xjzkd1k", pid: #PID<0.178.0>, type: "import"},
%{id: "MaTRVTw", pid: #PID<0.180.0>, type: "import"}
]
iex(5)>
Limiting Concurrency of Important Jobs
Now we can really use the running_imports/0 to good use. Add this check to Jobber.start_job/1
alias Jobber.{JobRunner, Job, JobSupervisor}
# ...
def start_job(args) do
if Enum.count(running_imports()) >= 5 do
{:error, :import_quota_reached}
else
DynamicSupervisor.start_child(JobRunner, {JobSupervisor, args})
end
end
Now test it out by trying to add more than 5 jobs.
iex(1)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.166.0>}
iex(2)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.168.0>}
iex(3)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.170.0>}
iex(4)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.172.0>}
iex(5)> Jobber.start_job(work: good_job, type: "import")
{:ok, #PID<0.174.0>}
iex(6)> Jobber.start_job(work: good_job, type: "import")
{:error, :import_quota_reached}
Inspecting Supervisors at Runtime
You can query a running supervision tree with two functions. Both are available in Supervisor and DynamicSupervisor.
count_children/1
Takes a supervisor’s PID or name and returns a map of child counts:
- :supervisors counts child supervisors, both active and inactive.
- :workers counts child workers, both active and inactive.
- :specs counts all children, both active and inactive.
- :active counts the children that are running right now.
iex(1)> DynamicSupervisor.count_children(Jobber.JobRunner)
%{active: 0, specs: 0, supervisors: 0, workers: 0}
which_children/1
Returns a list of {id, pid, type, [module]} tuples:
- id is :undefined for dynamically started children.
- pid is the child’s PID, or :restarting while it is being restarted.
- type is :worker or :supervisor.
- [module] is the implementing module, inside a list.
The Idle JobSupervisor
After a job finishes, JobRunner still shows active: 1, supervisors: 1.
iex(2)> DynamicSupervisor.which_children(Jobber.JobRunner)
[{:undefined, #PID<0.152.0>, :supervisor, [Jobber.JobSupervisor]}]
The Job process has exited, but its parent JobSupervisor keeps running with nothing to do. To confirm this, check the JobSupervisor directly:
iex(3)> {_, pid, _, _} = List.first(DynamicSupervisor.which_children(Jobber.JobRunner))
iex(4)> Supervisor.count_children(pid)
%{active: 0, specs: 1, supervisors: 0, workers: 1}
Idle processes are cheap and usually harmless. If you want to clean one up, stop it:
iex(5)> Supervisor.stop(pid)
:ok
After that, DynamicSupervisor.which_children(Jobber.JobRunner) returns [].
Use count_children/1 to get the numbers and which_children/1 to see what each child is. Together they let you find leftover processes in the tree and stop them.
Wrapping Up
We now have the supervisor trees set and you can work with them.
I would take some time now and get used to doing more and more with the tree. You can work on persistency and even adding more into queues.