elixirnewbie

elixirnewbie

How can I get the message queue length of a GenServer? I have tried

def handle_call({:add,number},_from, state) do
  
  {_,num}=Process.info(self(), :message_queue_len)
 IO.inspect num
  {:reply,reply state}
end

but getting 0.

Showing Posts 1 to 10

blatyo

blatyo

Conduit Core Team

How are you sending messages to that process?

I notice that you’re doing this in a handle_call which is synchronous. So, if you only have one process doing a GenServer.call to this process it’ll block until it gets a response. That means, you’ll never have a queue length, because this handle_call is consuming the only message that has been sent.

In short, what you’re doing is correct, but the way you’re testing it isn’t.

blatyo

blatyo

Conduit Core Team

Try something like this to test:

1..100
|> Task.async_stream(fn _ -> 
  # Thing that calls your process
end)
|> Stream.run()
elixirnewbie

elixirnewbie OP

Thanks @blatyo

KronicDeth

KronicDeth

If you want to observe the size and not assert on it, :observer is specifically for this. :observer.start() will open the Observer WX GUI, go to the processes tab and you can sort on the MsgQ column. Right-clicking will allow you to get a dump of the current messages, but the Processes tab will refresh itself every so often, so you can use it to monitor for message queues growing, which usually means you have a bug with an unmatched message. I have found at least 3 bugs in production using only the MsgQ column and clicking to sort it. It is super useful.

elixirnewbie

elixirnewbie OP

@KronicDeth thanks…but I want the GenServer to take actions based on the queue length.Imagine 3 GenServers spun up by a DynamicSupervisor and they are all sending messages to each other but they don’t accept new messages if the message count is over 50 in their respective queues and send a reply back to the caller ..something like {:error, :queue_full}

stefanchrobot

stefanchrobot

Sounds like you want to control the ingestion of data in your system. If that’s the case, rather than fiddling with the queue length, I’d suggest taking a look at GenStage.

peerreynders

peerreynders

That is irrelevant given that the process can only control removal of messages from the mailbox and nothing else - i.e. a process can’t stop accepting messages; it can only stop accepting work which would be tracked inside the process state.

something like {:error, :queue_full}

a response like that happens when a received message is requesting more work but your work backlog indicates that you already have enough - i.e. the response should have nothing to do with messages remaining in the process mailbox but everything with information you currently have in the process state.

elixirnewbie

elixirnewbie OP

Thanks @stefanchrobot @peerreynders
@peerreynders basically I want to know the total remaining messages in the mailbox that haven’t been processed

dom

dom

A better way is to use ETS to track the queue length, then you can tell that the queue is full before messaging the process. See “Bounded Queues” under Handling Overload

Edit: just realized you might be beginning with Elixir, so this is probably overkill. What kind of work do your processes do? There’s likely a way to rearrange the problem so you don’t need bounded queues. For instance, using a process per message / task and limiting concurrency via DynamicSupervisor’s max children instead.

peerreynders

peerreynders

You already know the answer - you just aren’t creating the right conditions to see the effects.

defmodule Demo do

  defp do_it(state) do
    {_, num} = Process.info(self(), :message_queue_len)
    IO.puts "queue length: #{inspect num}"

    reply = {:reply, :hello, state}
    cond do
      num > 0 ->
        reply
      true ->
        send(self(), :done)
        reply
    end
  end

  def init(_args) do
    send(self(), :block) # get the first message in the mailbox
    {:ok, []}
  end

  def handle_call(:hi, _from, state),
    do: do_it(state)

  def handle_info(:block, state) do
    Process.sleep(1000) # block for one second
    {:noreply, state}
  end
  def handle_info(:done, state) do
    {:stop, :normal, state}
  end

  def terminate(reason, state) do
    IO.puts "terminate: #{inspect reason} #{inspect state}"
  end

end

{:ok, pid} = GenServer.start_link(Demo,[])
f = fn ->
  GenServer.call(pid, :hi)
  :ok
end

(for _ <- 1..4, do: Task.async(f))
|> Enum.map(&Task.await(&1))
$ elixir demo.exs
queue length: 3
queue length: 2
queue length: 1
queue length: 0
terminate: :normal []
$

Where Next? Top

Trending in Questions Top

RSP87
I’m working on a project that simulates the bumbl example in the programming phoenix book. It acts almost like an email client. We have a...
New
kszambelanczyk
Hello! Could someone please give me a help/sample code, how to delete a file from s3 using waffle/waffle_ecto from Phoenix app. I creat...
New
RemyXRenard
I’m seeing that a list inside a Kino.DataTable will be interpreted as a charlist, even if the Kino.configure() is set to charlists: :as_l...
New
velrest
So my question is quite simple and i have found no conclusive answer on forum, google or AI. Should we use :erlang.float for Integer to ...
New
samoloth
Hi, I’ve just set up an application with ash_authentication. There is only magic link strategy for now, so there is no confirmation add o...
New
FlyingNoodle
If a change or preparation module uses Ash.Changeset.get_argument/2 or Ash.Query.get_argument/2 (or any of the other get_argument functio...
New
psy-q
I’m trying to set up Emacs with elixir-ls via lsp-mode and credo via Flycheck. This should mostly be preconfigured as Flycheck picks up c...
New

Other Trending Topics Top

mudasobwa
I am happy to introduce the very α version of the new programming language compiled to BEAM. Welcome Cure. It has literally three kille...
New
garrison
Hobbes is a low-level distributed database for the Elixir programming language. Hobbes provides a simple, safe, and scalable storage lay...
New
marciok
Hi there! We created Gust: A task orchestrator inspired by Airflow. For those who have never heard about Aiflow, it’s a Python-based wor...
New
jimsynz
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
Dmk
Xamal is a deployment tool for Elixir apps that deploys native releases to bare metal servers over SSH. It’s a port of GitHub - basecamp/...
New
Damirados
Hello everyone. After busy few months I am happy to announce v0.1.0 of Emerge &amp; Solve. They are GUI (Emerge) and State management (S...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews