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

