You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Elixir中使用GenServer实现消息代理的队列系统?

用Elixir GenServer实现带长期运行流的队列方案

核心设计思路

用队列管理器GenServer维护「队列名称→队列进程PID」的映射,每个队列是独立的GenServer进程,负责自身的订阅者管理、消息广播和流维护。全程遵循OTP规范,替代手动spawn和send的原始并发方式,保证可靠性和可维护性。

分步实现代码

1. 启动依赖组件(Registry + 队列管理器)

先启动Registry用于队列名称的全局注册,避免进程名称冲突;同时启动队列管理器维护队列映射:

defmodule MyApp.Application do
  use Application

  @impl true
  def start(_type, _args) do
    children = [
      # 注册唯一队列名称的Registry
      {Registry, keys: :unique, name: QueueRegistry},
      # 队列管理器GenServer
      QueueManager
    ]

    Supervisor.start_link(children, strategy: :one_for_one)
  end
end

2. 队列管理器GenServer

负责创建队列、查询队列PID,状态存储队列名称与进程的映射:

defmodule QueueManager do
  use GenServer

  # 客户端API
  def start_link(_opts) do
    GenServer.start_link(__MODULE__, %{}, name: __MODULE__)
  end

  # 创建队列(存在则返回已有PID)
  def create_queue(queue_name) do
    GenServer.call(__MODULE__, {:create_queue, queue_name})
  end

  # 根据名称获取队列PID
  def get_queue_pid(queue_name) do
    GenServer.call(__MODULE__, {:get_queue_pid, queue_name})
  end

  # 服务器回调
  @impl true
  def init(state), do: {:ok, state}

  @impl true
  def handle_call({:create_queue, queue_name}, _from, state) do
    case Map.get(state, queue_name) do
      nil ->
        {:ok, pid} = Queue.start_link(queue_name)
        {:reply, {:ok, pid}, Map.put(state, queue_name, pid)}
      pid ->
        {:reply, {:ok, pid}, state}
    end
  end

  @impl true
  def handle_call({:get_queue_pid, queue_name}, _from, state) do
    {:reply, Map.get(state, queue_name), state}
  end
end

3. 单个队列GenServer

负责维护订阅者列表、消息广播,支持长期运行的流:

defmodule Queue do
  use GenServer

  # 启动队列进程,通过Registry注册名称
  def start_link(queue_name) do
    GenServer.start_link(__MODULE__, {queue_name, []}, name: via_tuple(queue_name))
  end

  # 客户端API:订阅队列
  def subscribe(queue_name, subscriber_pid) do
    GenServer.cast(via_tuple(queue_name), {:subscribe, subscriber_pid})
  end

  # 客户端API:向队列发布消息(广播给所有订阅者)
  def publish(queue_name, message) do
    GenServer.cast(via_tuple(queue_name), {:publish, message})
  end

  # 获取当前订阅者列表(可选)
  def get_subscribers(queue_name) do
    GenServer.call(via_tuple(queue_name), :get_subscribers)
  end

  # 生成Registry的via元组,用于通过名称查找队列进程
  defp via_tuple(queue_name) do
    {:via, Registry, {QueueRegistry, queue_name}}
  end

  # 服务器回调
  @impl true
  def init({queue_name, subscribers}), do: {:ok, {queue_name, subscribers}}

  # 处理订阅请求,监控订阅者进程避免僵尸订阅
  @impl true
  def handle_cast({:subscribe, pid}, {queue_name, subscribers}) do
    Process.monitor(pid)
    new_subscribers = if pid in subscribers, do: subscribers, else: [pid | subscribers]
    {:noreply, {queue_name, new_subscribers}}
  end

  # 处理消息发布,广播给所有订阅者
  @impl true
  def handle_cast({:publish, message}, {queue_name, subscribers}) do
    Enum.each(subscribers, fn pid -> send(pid, {:queue_message, queue_name, message}) end)
    {:noreply, {queue_name, subscribers}}
  end

  @impl true
  def handle_call(:get_subscribers, _from, state = {_name, subs}) do
    {:reply, subs, state}
  end

  # 订阅者进程退出时,自动从列表移除
  @impl true
  def handle_info({:DOWN, _ref, :process, pid, _reason}, {name, subs}) do
    {:noreply, {name, List.delete(subs, pid)}}
  end

  # 忽略未知消息
  @impl true
  def handle_info(_msg, state), do: {:noreply, state}
end

关键问题解决

  • 通过名称访问队列:两种方式可选:
    1. 调用QueueManager.get_queue_pid/1获取PID后通信;
    2. 直接用Queue.via_tuple/1生成的元组,通过GenServer API直接调用队列进程(更高效)。
  • PID存储与安全:队列管理器的状态是进程内的映射,GenServer串行处理请求,避免了手动维护映射的并发问题。
  • 长期运行保障:队列进程由Supervisor监控,崩溃后自动重启;同时通过Process.monitor/1清理失效订阅者,保证流的稳定性。
  • 替代原始并发方式:用GenServer的call/cast替代手动send和spawn,符合OTP规范,自带错误处理和进程生命周期管理。

内容的提问来源于stack exchange,提问作者cecuca

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 11:07:45