RabbitMQ 本身就是用 Erlang 写的,这让它成为 Erlang 生态中最「原生」的消息中间件——它的队列、连接、通道都映射为 BEAM 进程,其集群能力直接受益于分布式 Erlang。对 Erlang/Elixir 应用而言,RabbitMQ 既是解耦服务的手段,也是削峰填谷的缓冲层。但消息中间件用错模式会导致消息丢失、重复消费、连接风暴、队列堆积等一连串问题,且这些问题往往在生产环境才暴露。本文将从 AMQP 协议模型出发,逐层讲解交换机路由、确认机制、可靠投递、幂等设计、连接池与背压,最后给出集群运维要点。
一、AMQP 模型与交换机类型
1.1 AMQP 0-9-1 的核心概念
AMQP 把消息系统抽象为几个正交概念:
| 概念 | 作用 | 生命周期 |
|---|---|---|
| Connection | TCP 连接,承载多路复用 | 应用级,长连接 |
| Channel | 连接内的轻量逻辑通道 | 每次操作集,短 |
| Exchange | 接收消息并按规则路由 | 声明式,持久 |
| Queue | 存储消息供消费 | 声明式,持久 |
| Binding | Exchange 到 Queue 的路由规则 | 声明式,持久 |
| Virtual Host | 逻辑隔离命名空间 | 运维创建 |
消息流:Producer → Exchange → (Binding 路由) → Queue → Consumer。
1.2 四种交换机类型
%% 声明交换机
amqp_channel:call(Channel, #'exchange.declare'{
exchange = <<"orders">>,
type = <<"topic">>,
durable = true,
auto_delete = false
}).
| 类型 | 路由规则 | 典型场景 |
|---|---|---|
direct | routing key 精确匹配 | 点对点任务分发 |
fanout | 广播到所有绑定队列 | 事件通知、缓存失效 |
topic | 通配符匹配(* 一词,# 多词) | 按类别订阅 |
headers | 按消息头匹配 | 复杂条件路由(少用) |
1.3 topic 路由实战
%% 绑定:订阅所有订单相关事件
amqp_channel:call(Channel, #'queue.bind'{
queue = <<"order_events">>,
exchange = <<"events">>,
routing_key = <<"order.#">> % 匹配 order.created / order.paid.succeeded
}).
%% 绑定:只订阅支付完成
amqp_channel:call(Channel, #'queue.bind'{
queue = <<"payment_events">>,
exchange = <<"events">>,
routing_key = <<"order.*.succeeded">>
}).
%% 发布
amqp_channel:cast(Channel, #'basic.publish'{
exchange = <<"events">>,
routing_key = <<"order.paid.succeeded">>},
#amqp_msg{payload = jsx:encode(#{order_id => 42})}).
1.4 默认交换机与队列直连
每个 vhost 有一个默认交换机(名字为空串),它把消息路由到同名队列,适合简单场景:
%% 直接发到队列(不经过自定义交换机)
amqp_channel:cast(Channel, #'basic.publish'{
exchange = <<>>,
routing_key = <<"task_queue">>},
#amqp_msg{payload = <<"job">>}).
设计建议:生产系统不要依赖默认交换机,显式声明 exchange 与 binding,把路由规则纳入版本管理,避免「隐式约定」导致的路由错误。
二、生产消费与确认机制
2.1 建立连接与通道
-module(mq_client).
-export([connect/0, publish/2, consume/1]).
-define(EXCHANGE, <<"events">>).
connect() ->
{ok, Conn} = amqp_connection:start(#amqp_params_network{
host = "rabbit.internal",
port = 5672,
username = <<"app">>,
password = <<"secret">>,
virtual_host = <<"/prod">>,
heartbeat = 30, % 心跳,防半开连接
connection_timeout = 5000
}),
{ok, Channel} = amqp_connection:open_channel(Conn),
{Conn, Channel}.
2.2 发布消息
publish(Channel, Payload) ->
Method = #'basic.publish'{
exchange = ?EXCHANGE,
routing_key = <<"order.created">>,
mandatory = true % 路由不到队列时返回 basic.return
},
Props = #'P_basic'{
content_type = <<"application/json">>,
delivery_mode = 2, % 2 = 持久化
message_id = uuid(),
timestamp = erlang:system_time(second)
},
amqp_channel:cast(Channel, Method, #amqp_msg{props = Props, payload = Payload}).
2.3 消费与手动确认
自动确认(no_ack = true)在消息投递后立即删除,进程崩溃就丢消息;生产必须用手动确认:
consume(Channel) ->
amqp_channel:subscribe(Channel, #'basic.consume'{
queue = <<"order_events">>,
no_ack = false % 手动 ack
}, self()),
receive
{#'basic.consume_ok'{}, _} -> loop(Channel)
end.
loop(Channel) ->
receive
{#'basic.deliver'{delivery_tag = Tag}, #amqp_msg{payload = Payload}} ->
case handle(Payload) of
ok ->
amqp_channel:cast(Channel, #'basic.ack'{delivery_tag = Tag});
{error, _Reason} ->
%% 拒绝并重新入队(或进死信)
amqp_channel:cast(Channel, #'basic.nack'{
delivery_tag = Tag, requeue = false})
end,
loop(Channel)
end.
| 确认方式 | 语义 | 风险 |
|---|---|---|
basic.ack | 成功,删除消息 | 无 |
basic.nack + requeue=true | 失败,重新入队 | 可能死循环 |
basic.nack + requeue=false | 失败,丢弃或进死信 | 需配 DLX |
basic.reject | 单条拒绝 | 同 nack |
2.4 prefetch 与公平分发
默认 RabbitMQ 会把消息轮询推给消费者,不关心消费者是否繁忙。用 basic.qos 限制未确认消息数:
%% 每个消费者最多 10 条未确认消息
amqp_channel:call(Channel, #'basic.qos'{prefetch_count = 10}).
prefetch 调优:设太小(如 1)吞吐上不去;设太大则单消费者堆积。经验值是「单条处理时间 × 期望并发」的 2~3 倍,配合压测调整。
三、可靠投递与幂等设计
3.1 消息丢失的三个环节
消息可能在三个地方丢失,必须逐层防护:
| 环节 | 丢失原因 | 防护 |
|---|---|---|
| 生产者 → Broker | 网络中断、Broker 崩溃 | Publisher Confirm |
| Broker 存储 | 未持久化 | delivery_mode=2 + durable 队列 |
| Broker → 消费者 | 自动确认、消费中崩溃 | 手动 ack |
3.2 Publisher Confirm
开启 confirm 模式后,Broker 会异步确认每条消息已落盘:
%% 开启 confirm 模式
amqp_channel:call(Channel, #'confirm.select'{}),
amqp_channel:register_confirm_handler(Channel, self()),
%% 发布
publish_with_confirm(Channel, Payload) ->
#'basic.publish'{exchange = <<"events">>} = Method,
amqp_channel:cast(Channel, Method, #amqp_msg{payload = Payload}),
receive
{#'basic.ack'{delivery_tag = Seq}, _} ->
{ok, Seq};
{#'basic.nack'{delivery_tag = Seq}, _} ->
{error, {rejected, Seq}}
after 5000 ->
{error, confirm_timeout}
end.
3.3 死信队列
被拒绝或过期的消息进入死信交换机(DLX),便于后续排查与补偿:
%% 声明业务队列时绑定死信交换机
amqp_channel:call(Channel, #'queue.declare'{
queue = <<"order_events">>,
durable = true,
arguments = [
{<<"x-dead-letter-exchange">>, longstr, <<"dlx">>},
{<<"x-dead-letter-routing-key">>, longstr, <<"order.failed">>},
{<<"x-message-ttl">>, long, 300000} % 5 分钟未消费进死信
]
}).
3.4 幂等消费
「至少一次」投递意味着消息可能重复,消费端必须幂等:
%% 方案一:业务唯一键去重(推荐)
handle_message(Payload) ->
#{message_id := MsgId, order_id := OrderId} = jsx:decode(Payload, [return_maps]),
case ets:insert_new(processed, {MsgId, erlang:system_time(second)}) of
false ->
{ok, duplicate_ignored}; % 已处理过
true ->
apply_order(OrderId) % 真正处理
end.
%% 方案二:数据库唯一约束
%% INSERT INTO orders (...) ON CONFLICT (message_id) DO NOTHING;
关键点:幂等键应来自业务语义(订单号 + 操作类型),而非消息 ID。因为上游重发时可能生成新的消息 ID。
四、连接池与背压
4.1 为什么需要连接池
RabbitMQ 的 channel 不是线程安全的,且单个 channel 上的确认是串行的。高并发发布需要多个 channel:
%% 用 poolboy 管理 channel 池
{poolboy, [
{name, {local, mq_pool}},
{worker_module, mq_worker},
{size, 16},
{max_overflow, 8}
]}.
%% mq_worker.erl
-module(mq_worker).
-behaviour(gen_server).
-behaviour(poolboy_worker).
init([Conn]) ->
{ok, Channel} = amqp_connection:open_channel(Conn),
{ok, #state{channel = Channel}}.
%% 借出 channel 发布
publish(Payload) ->
poolboy:transaction(mq_pool, fun(Worker) ->
gen_server:call(Worker, {publish, Payload})
end).
4.2 连接级 vs 通道级池化
| 粒度 | 优点 | 缺点 |
|---|---|---|
| 连接池 | 隔离彻底,故障域小 | 资源开销大(每连接一个 TCP) |
| 通道池 | 复用 TCP,开销小 | 单连接故障影响全部通道 |
| 混合 | 少量连接 + 每连接多通道 | 需管理两级生命周期 |
生产推荐:每个消费者独立连接,生产者共享连接池 + 通道池。
4.3 背压与限流
Broker 端堆积时,生产者必须感知并降速,否则内存会被消息撑爆:
%% 1. 监控队列深度,超过阈值时拒绝新请求
check_backpressure() ->
case queue_depth(<<"order_events">>) of
N when N > 100_000 -> {error, overloaded};
_ -> ok
end.
%% 2. 发布侧限流:令牌桶
publish_limited(Payload) ->
case ets:update_counter(rate_limit, tokens, -1, {tokens, 100}) of
N when N >= 0 -> publish(Payload);
_ ->
ets:update_counter(rate_limit, tokens, 1),
{error, rate_limited}
end.
4.4 心跳与重连
%% 断线自动重连(用 supervisor 包裹连接进程)
-module(mq_conn_sup).
-behaviour(supervisor).
%% 连接进程崩溃时由 supervisor 重启,配合指数退避
init([]) ->
{ok, {#{strategy => one_for_one, intensity => 5, period => 30},
[#{id => mq_conn,
start => {mq_conn, start_link, []},
restart => permanent,
shutdown => 5000,
type => worker}]}}.
心跳的必要性:网络设备会在空闲时静默断开 TCP 连接,心跳(默认 60s,建议 30s)能让双方及时发现死连接。消费端还应处理
basic.cancel通知(队列被删除时会收到)。
五、集群与运维
5.1 集群架构
RabbitMQ 集群由多个节点组成,队列数据默认只存在声明它的节点(除非是 quorum queue):
| 队列类型 | 复制方式 | 一致性 | 推荐 |
|---|---|---|---|
| classic | 镜像队列(已弃用) | 最终一致 | 迁移 |
| quorum | Raft 多数派 | 强一致 | 生产首选 |
| stream | 追加日志 | 顺序一致 | 大吞吐日志 |
%% 声明 quorum 队列
amqp_channel:call(Channel, #'queue.declare'{
queue = <<"orders">>,
durable = true,
arguments = [{<<"x-queue-type">>, longstr, <<"quorum">>}]
}).
5.2 关键运维指标
%% 通过管理插件查看
rabbitmqctl list_queues name messages consumers memory
rabbitmqctl list_connections name state channels
rabbitmqctl list_channels name number pending_acks
rabbitmqctl list_queues name messages | awk '$2 > 100000' % 积压告警
| 指标 | 含义 | 告警阈值 |
|---|---|---|
messages_ready | 待消费消息数 | 持续增长 |
messages_unacknowledged | 已投递未确认 | 接近 prefetch × 消费者数 |
consumers | 消费者数量 | 为 0 且队列非空 |
memory | 队列占用内存 | 接近 watermark |
disk_free | 磁盘剩余 | 低于 watermark 阻塞生产者 |
5.3 内存与磁盘水位
%% RabbitMQ 达到内存水位会阻塞所有生产者(blocked 状态)
rabbitmqctl set_vm_memory_high_watermark 0.6 # 60% 内存
rabbitmqctl set_disk_free_limit 5GB # 磁盘下限
流控现象:生产者连接进入
blocked状态时,basic.publish会挂起而非报错。应用必须设置发布超时,否则业务线程会集体卡死。
5.4 优雅停机与消费端
%% 停机时先停止消费、处理完在途消息、再关闭连接
terminate(_Reason, #state{channel = Channel, conn = Conn}) ->
amqp_channel:call(Channel, #'basic.cancel'{consumer_tag = <<"ctag">>}),
drain_inflight(), % 等待在途消息处理完
amqp_channel:close(Channel),
amqp_connection:close(Conn),
ok.
配合 OTP 监督树(见 https://plumephp.com/erlang-otp-framework/)可以让连接进程随应用生命周期启停,避免重启时的连接泄漏。
六、总结
RabbitMQ 的可靠性不来自中间件本身,而来自生产者、Broker、消费者三方的协同约定。本文的要点:
- AMQP 模型:Exchange 决定路由,Queue 决定存储,Binding 决定关系,四类交换机各有适用场景;
- 确认机制:消费端必须手动 ack,用
basic.qos控制 prefetch,避免「推给忙消费者」; - 可靠投递:Publisher Confirm 防生产端丢失,
delivery_mode=2防 Broker 丢失,手动 ack 防消费端丢失,三者缺一不可; - 幂等:至少一次投递意味着必须去重,幂等键取自业务语义而非消息 ID;
- 池化与背压:用连接池 + 通道池提升并发,用水位监控与令牌桶做背压,防内存被消息撑爆;
- 集群:优先 quorum 队列,盯紧
messages_ready、disk_free、blocked三个信号。
把这些约定固化到代码与运维流程中,消息中间件才能成为系统的稳定缓冲层,而不是新的故障源。消息的序列化与协议解析细节可延伸阅读 https://plumephp.com/erlang-bit-syntax-binaries/,集群网络与分区处理可参考 https://plumephp.com/erlang-distributed-programming/。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。