Elixir Study Projects
Project Inventory
| Project | Concepts combined | Start from |
|---|---|---|
| Supervised KV store | GenServer, ETS backup, supervision tree | GenServer demo |
| Job queue | Oban, Ecto, telemetry, retry semantics | Oban demo |
| Metrics pipeline | Broadway, back-pressure, batching, telemetry | Broadway demo |
| Distributed chat | PubSub, Registry, LiveView, clustering | Registry demo |
Project 1: Supervised KV Store
A key/value store that survives crashes: a GenServer owns the map, an ETS table (owned by a
separate "guardian" process) holds the last snapshot, and on restart the server rebuilds
from ETS. Reading order: Guardian (ETS owner) → Store (GenServer)
→ StoreSupervisor (tree).
defmodule Guardian do
# A process whose ONLY job is owning the ETS table. If the store
# crashes, this process still holds the data.
use GenServer
@table :kv_backup
def start_link(_), do: GenServer.start_link(__MODULE__, nil, name: __MODULE__)
@impl true
def init(_) do
# :public + named_table lets the Store process read/write directly.
:ets.new(@table, [:set, :named_table, :public, read_concurrency: true])
{:ok, nil}
end
def snapshot(map), do: :ets.insert(@table, {:data, map})
def restore do
case :ets.lookup(@table, :data) do
[{:data, map}] -> map
[] -> %{}
end
end
end
defmodule Store do
use GenServer
def start_link(_), do: GenServer.start_link(__MODULE__, nil, name: __MODULE__)
@impl true
def init(_) do
# REBUILD state from the guardian after every restart — the crash
# loses nothing but the writes since the last snapshot.
{:ok, Guardian.restore(), {:continue, :snapshot_loop}}
end
@impl true
def handle_continue(:snapshot_loop, state) do
Process.send_after(self(), :snapshot, 5_000)
{:noreply, state}
end
@impl true
def handle_info(:snapshot, state) do
Guardian.snapshot(state)
Process.send_after(self(), :snapshot, 5_000)
{:noreply, state}
end
@impl true
def handle_call({:put, k, v}, _from, state), do: {:reply, :ok, Map.put(state, k, v)}
def handle_call({:get, k}, _from, state), do: {:reply, Map.get(state, k), state}
end
# The tree — rest_for_one guarantees the Guardian boots FIRST, so
# Store.restore/0 always finds the table it depends on.
defmodule StoreSupervisor do
use Supervisor
def start_link(_), do: Supervisor.start_link(__MODULE__, nil, name: __MODULE__)
@impl true
def init(_) do
Supervisor.init([Guardian, Store], strategy: :rest_for_one)
end
end
# Boot and try it:
{:ok, _} = Supervisor.start_link(StoreSupervisor, [])
:ok = GenServer.call(Store, {:put, :lang, "elixir"})
GenServer.call(Store, {:get, :lang}) #=> "elixir"
# Kill the store — it restarts and RESTORES from ETS:
[{_, store, _, _}] = Supervisor.which_children(StoreSupervisor)
Process.exit(store, :kill)
Process.sleep(100)
GenServer.call(Store, {:get, :lang}) #=> "elixir" (rebuilt!)
Experiments: move the snapshot interval to 100ms and measure the cost; add del/1; make Guardian snapshot to disk on shutdown.
Project 2: Job Queue
Background jobs with retries and observability, using Oban on PostgreSQL. Reading order:
Worker → config → insert a job. The point to study is how much of the OTP
machinery (polling, retries, supervision) a production library gives you — and that its pieces are
the same behaviours you wrote by hand.
# application.ex — mount Oban in the tree (needs a Postgres repo).
children = [
MyApp.Repo,
{Oban, queues: [default: 10], repo: MyApp.Repo}
]
defmodule MyApp.Jobs.EmailWorker do
use Oban.Worker,
queue: :default,
max_attempts: 5,
unique: [period: 60] # dedupe identical jobs for 60s
@impl Oban.Worker
def perform(%Oban.Job{args: %{"email" => to}}) do
case Mailer.deliver(to) do
:ok -> :ok
# Returning {:error, _} schedules a retry with backoff —
# transient failures recover without human help.
{:error, reason} -> {:error, reason}
end
end
end
# Enqueue from anywhere — a plain struct, persisted in Postgres:
%{email: "ada@example.com"}
|> MyApp.Jobs.EmailWorker.new()
|> Oban.insert()
Experiments: force a failure and watch retry timing grow; add a second queue with lower concurrency; emit telemetry on job completion.
Project 3: Metrics Pipeline
Ingest events faster than you can write them: Broadway reads from a source, batches, and writes to a sink, with back-pressure so the queue drains on its own terms. Reading order: the producer → processors → batcher.
defmodule MyApp.MetricsPipeline do
use Broadway
def start_link(_opts) do
Broadway.start_link(__MODULE__,
name: __MODULE__,
producer: [
module: {MyApp.EventProducer, []}, # your own GenStage producer
concurrency: 1
],
processors: [
stages: 20, # parse/validate concurrently
max_demand: 100 # back-pressure: never take 1000 at once
],
batchers: [
sink: [batch_size: 500, batch_timeout: 2_000]
]
)
end
@impl true
def handle_message(_processor, %Message{data: raw} = msg, _ctx) do
msg
|> Message.update_data(&parse_event/1)
|> Message.put_batcher(:sink)
end
@impl true
# handle_batch writes 500 events in ONE database round-trip.
def handle_batch(:sink, messages, _batch_info, _ctx) do
events = Enum.map(messages, & &1.data)
MyApp.Repo.insert_all(MyApp.Event, events)
messages
end
defp parse_event(raw), do: raw |> JSON.decode!() |> normalize()
end
Experiments: log batch sizes while varying batch_size; kill a processor mid-run (Broadway re-drives the failed messages); swap the producer for an SQS source.
Project 4: Distributed Chat
The LiveView chat: one LiveView process per browser session, a Phoenix.PubSub topic
per room, and messages fanning out cluster-wide. Reading order: RoomLive →
Chat context → PubSub wiring.
defmodule MyAppWeb.RoomLive do
use MyAppWeb, :live_view
@topic "room:general"
def mount(_params, _session, socket) do
# Only subscribe once the socket is live — otherwise the client
# would miss nothing anyway, but tests would double-subscribe.
if connected?(socket), do: Phoenix.PubSub.subscribe(MyApp.PubSub, @topic)
{:ok, assign(socket, messages: [], name: nil)}
end
# Two distinct event sources, two handlers — the LiveView pattern:
def handle_event("join", %{"name" => name}, socket) do
Phoenix.PubSub.broadcast(MyApp.PubSub, @topic, {:joined, name})
{:noreply, assign(socket, name: name)}
end
def handle_event("say", %{"text" => text}, socket) do
Phoenix.PubSub.broadcast(MyApp.PubSub, @topic, {:said, socket.assigns.name, text})
{:noreply, socket}
end
def handle_info({:said, who, text}, socket) do
{:noreply, update(socket, :messages, &([&1 | ["#{who}: #{text}"]]))}
end
def handle_info({:joined, name}, socket) do
{:noreply, update(socket, :messages, &(["* #{name} joined*" | &1]))}
end
end
Experiments: run two nodes and send from one — the other receives (PubSub is cluster-aware); add presence to show connected users; cap message history in assigns to prevent memory growth.