RabbitMQ 与消息中间件:AMQP 模型与可靠投递

以 Erlang 客户端为主线,讲解 AMQP 模型与交换机类型、生产消费与确认、可靠投递与幂等设计、连接池与背压、集群与运维实践。

RabbitMQ 本身就是用 Erlang 写的,这让它成为 Erlang 生态中最「原生」的消息中间件——它的队列、连接、通道都映射为 BEAM 进程,其集群能力直接受益于分布式 Erlang。对 Erlang/Elixir 应用而言,RabbitMQ 既是解耦服务的手段,也是削峰填谷的缓冲层。但消息中间件用错模式会导致消息丢失、重复消费、连接风暴、队列堆积等一连串问题,且这些问题往往在生产环境才暴露。本文将从 AMQP 协议模型出发,逐层讲解交换机路由、确认机制、可靠投递、幂等设计、连接池与背压,最后给出集群运维要点。

一、AMQP 模型与交换机类型

1.1 AMQP 0-9-1 的核心概念

AMQP 把消息系统抽象为几个正交概念:

概念作用生命周期
ConnectionTCP 连接,承载多路复用应用级,长连接
Channel连接内的轻量逻辑通道每次操作集,短
Exchange接收消息并按规则路由声明式,持久
Queue存储消息供消费声明式,持久
BindingExchange 到 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
}).
类型路由规则典型场景
directrouting 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镜像队列(已弃用)最终一致迁移
quorumRaft 多数派强一致生产首选
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/。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「erlang」更多文章

  1. rebar3 构建与发布:Erlang 工程化的完整工具箱
  2. Erlang 数据库集成与连接池实战
  3. Dialyzer 与类型规范:Erlang 静态分析的工程实践