Code Review: Batching GenServer

I can’t claim I am an authority but here’s what I am using as an idiom when I need the following:

  1. A GenServer accumulating items to process inside its state;
  2. A GenServer that processes the currently accumulated items after a certain deadline, regardless of when you last queued an item – meaning that if you fire it up with initial 10 items and it has a periodic flush of 5 seconds set up and you don’t queue up anything after the initial 10 items, it will process them after 5 seconds. If you queue up another 100 items before the 5 seconds expire then it will process all 110 items. The point being, it will attempt to process its outstanding items each 5 seconds sharp, non-negotiable and independent on when do you enqueue items.

So here’s the code (76 lines):

defmodule YYY.BatchingWorker do
  use GenServer, restart: :permanent

  @wait_before_processing_millis 5_000

  def start_link(inputs) do
    GenServer.start_link(
      __MODULE__,
      inputs,
      name: __MODULE__
    )
  end

  def child_spec(inputs) do
    %{
      id: __MODULE__,
      start: {__MODULE__, :start_link, [inputs]},
      restart: :permanent
    }
  end

  @impl GenServer
  def init(inputs) do
    {:ok, %{timer_ref: make_timer(self()), inputs: inputs}}
  end

  @impl GenServer
  def handle_call({:enqueue, input}, _from, %{inputs: inputs} = state) do
    new_state = Map.put(state, :inputs, inputs ++ [input])
    {:reply, :ok, new_state}
  end

  @impl GenServer
  def handle_call(:get_state, _from, state) do
    {:reply, state, state}
  end

  @impl GenServer
  def handle_info(:process, %{timer_ref: timer_ref, inputs: inputs}) do
    if [] != inputs do
      process_batch(inputs)
    end

    cancel_timer(timer_ref)
    new_timer_ref = make_timer(self())

    {:noreply, %{timer_ref: new_timer_ref, inputs: []}}
  end

  def enqueue(input) do
    GenServer.call(__MODULE__, {:enqueue, input})
  end

  def get_state() do
    GenServer.call(__MODULE__, :get_state)
  end

  def process() do
    send(__MODULE__, :process)
  end

  def stop() do
    GenServer.stop(__MODULE__)
  end

  def process_batch(inputs) do
    IO.puts("Processing #{inspect(inputs)}")
  end

  defp make_timer(pid) when is_pid(pid) do
    Process.send_after(pid, :process, @wait_before_processing_millis)
  end

  defp cancel_timer(nil), do: false
  defp cancel_timer(ref) when is_reference(ref), do: Process.cancel_timer(ref)
end

Notice the following:

  • You can just put YYY.BatchingWorker inside your supervision tree without options, or you can start it up manually whenever you wish (if you specify name: YourChosenModuleName though! Otherwise the convenience functions I specified below will crash – they rely on the GenServer being named after the module they are put in).
  • Accumulating items to process does not have to be synchronous (it’s currently handled by handle_call); you can migrate it to handle_cast if desired so queuing items will never wait on anything.
  • The message to process stuff is NOT a GenServer message; it’s a lower-level process message and that’s why it’s handled by handle_info – this is needed because we want to be able to both (1) manually trigger processing and (2) have a timer trigger processing for us.
  • I have added utility / convenience functions:
    • enqueue (enqueue item or items)
    • get_state (for debugging, feel free to remove)
    • process (manually start processing before the deadline expires)
    • stop (to kill the GenServer: again for debugging purposes and again feel free to remove).
  • You can of course change the name of the module, I’ve used __MODULE__ everywhere so you just have to change the name once in the defmodule definition.
  • process_batch is obviously not set in stone, you might supply a function in the GenServer’s initial state even, it does not have to live inside the GenServer module after all.
  • You cannot modify the handle_info(:process) handler to return something different than :ok if the current batch to process is empty. I left that out because to me that’s not an error; I’d presume it’s possible that no items were enqueued within the deadline (if the deadline is short enough). But if you want to differentiate between both cases (no items vs. there are items), it’s doable – exercise for the reader. :smiley:

I’ve used the above code in more than one contract and personal projects and it works pretty well. You can just put it in a toy project and play with it:

{:ok, _} = GenServer.start_link(YYY.BatchingWorker, [1, 2], name: YYY.BatchingWorker)
YYY.BatchingWorker.process()
YYY.BatchingWorker.enqueue([1,2,3,4,5])
:timer.sleep(5000)
# etc.

Hope that’s helpful. I wouldn’t get paralyzed on what’s idiomatic per se, I’d more look into the proper OTP building blocks to model my problem.