RabbitMQ 运行链:路由、确认、预取与副本恢复
basicPublish(exchange, routingKey, properties, body) 把消息交给交换器。交换器按 binding 选择目标队列,队列保存并交付消息,消费者再决定何时确认处理。一条发送 API 背后包含了路由、存储和消费三套不同操作。
下面使用 RabbitMQ 4.3.5 和 Java 客户端 5.33.0 的 AMQP 0-9-1 API。RabbitMQ 也支持 AMQP 1.0、MQTT、Stream 等协议,它们的会话、消息确认和客户端接口需分别查阅,不能直接套用 channel / delivery tag 的全部规则。
建立 AMQP 连接和消息拓扑
连接、通道与虚拟主机
RabbitMQ 节点
├── TCP/TLS 连接 connection
│ ├── channel 1:生产者发布与 confirm
│ └── channel 2:消费者订阅与 ACK
└── virtual host:lab
├── exchange:me12.rabbit.route,类型 topic
├── binding:order.* → queue one
├── binding:#.paid → queue two
├── queue one:库存订阅
└── queue two:审计订阅connection 是客户端与节点之间的网络连接,channel 是连接内的协议会话。多个 channel 共用一条连接,因此连接故障会同时影响它们;某个 channel 因声明错误被关闭时,同一 connection 上的其他 channel 仍可使用。
vhost 隔离拓扑名称和权限。一条 connection 登录一个 vhost,其中的队列、exchange 和 binding 使用该 vhost 的名称空间。应用可分别持有发布 channel 和消费 channel,减少生命周期混用;不要为每条消息新建 TCP 连接,也不要让多个发布线程随意共享一个 channel。Java 连接与通道 API
Exchange 选择队列,Binding 表达匹配关系
| Exchange 类型 | 路由依据 | 适合表达的关系 |
|---|---|---|
| direct | routing key 与 binding key 精确相等 | 指定业务类别、错误等级或明确目标 |
| topic | 以点分段,* 匹配一个段,# 匹配零个或多个段 | 按实体、动作、区域组合订阅 |
| fanout | 向所有绑定目标分发,忽略 routing key | 独立订阅方都处理每次发布 |
| headers | 按头字段与 x-match 条件匹配 | 多字段属性路由 |
队列声明后,会自动按队列名绑定到默认 exchange。向 exchange=""、routingKey=队列名 发布时仍经过交换器,只是这条默认 binding 省去了显式配置。复杂拓扑应创建自己的 exchange,不把默认交换器当作任意修改的业务路由表。Exchange 与 Binding
例如 order.paid 同时匹配 order.* 和 #.paid。如果两条 binding 指向两个不同队列,每个队列各得到一份消息;若两条匹配路径最终指向同一个队列,不应据此期待收到两份副本。消费者在同一个队列里竞争工作,跨队列才形成独立消费进度。
启动隔离实例并运行完整入口
下载完整实验工程,解压进入 message-event-lab。需要 Linux shell、Docker Engine、Compose v2+,由有 Docker 权限的普通宿主用户操作。项目内的 compose.yaml 提供 RabbitMQ 4.3.5、vhost lab 和无管理标签的应用用户 app;它不发布宿主端口。目录、初始化配置和数据库版本见 消息提交模型中的环境说明。这里只启动 RabbitMQ,不需要数据库。
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户'; exit 1; }
mkdir -p .m2
docker compose up -d --wait --wait-timeout 120 rabbit
docker compose exec rabbit rabbitmq-diagnostics -q ping预期 Ping succeeded。失败先看 docker compose ps 与 docker compose logs --tail=80 rabbit,排除服务启动问题后再检查账号。实验客户端使用普通 app 用户,其权限仅限 lab vhost;CLI 运维命令使用容器的 Erlang cookie,不代表 Java 应用也有集群管理权限。
rabbit 模块的核心依赖是:
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.33.0</version>
</dependency>共享 POM 还固定日志 API 与日志实现版本,编译目标为 Java 17。入口 rabbit/src/main/java/example/RabbitRuntimeLab.java 的完整代码如下:
package example;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Delivery;
import com.rabbitmq.client.GetResponse;
import com.rabbitmq.client.ShutdownSignalException;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.HashSet;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
public final class RabbitRuntimeLab {
static void check(boolean ok, String detail) {
if (!ok) throw new IllegalStateException(detail);
}
static ConnectionFactory factory() {
var f = new ConnectionFactory(); f.setHost(System.getenv().getOrDefault("RABBIT_HOST", "rabbit"));
f.setVirtualHost("lab"); f.setUsername("app"); f.setPassword("app-lab-only");
f.setAutomaticRecoveryEnabled(false); f.setConnectionTimeout(5000);
return f;
}
static void queue(Channel ch, String name) throws Exception {
ch.queueDeclare(name, true, false, false, Map.of("x-queue-type", "quorum"));
}
static void send(Channel ch, String exchange, String key, String id) throws Exception {
ch.basicPublish(exchange, key, true, new AMQP.BasicProperties.Builder()
.messageId(id).deliveryMode(2).contentType("text/plain").build(),
id.getBytes(StandardCharsets.UTF_8));
}
static GetResponse get(Channel ch, String q) throws Exception {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(8);
do {
GetResponse m = ch.basicGet(q, false);
if (m != null) return m;
Thread.sleep(25);
} while (System.nanoTime() < deadline);
throw new IllegalStateException("no delivery: " + q);
}
static int replyCode(Throwable error) {
for (Throwable t = error; t != null; t = t.getCause()) {
if (t instanceof ShutdownSignalException s && s.getReason() instanceof AMQP.Channel.Close c)
return c.getReplyCode();
}
return -1;
}
public static void main(String[] args) throws Exception {
check(args.length == 1, "mode: route|declare|prefetch|bad-ack|persist-send|persist-read");
try (var connection = factory().newConnection()) {
switch (args[0]) {
case "route" -> {
try (var ch = connection.createChannel()) {
String exchange = "me12.rabbit.route", q1 = exchange + ".one", q2 = exchange + ".two";
ch.exchangeDeclare(exchange, "topic", true);
queue(ch, q1); queue(ch, q2); ch.queuePurge(q1); ch.queuePurge(q2);
ch.queueBind(q1, exchange, "order.*"); ch.queueBind(q2, exchange, "#.paid");
var returned = new CompletableFuture<com.rabbitmq.client.Return>();
ch.addReturnListener((com.rabbitmq.client.ReturnCallback) returned::complete);
ch.confirmSelect(); send(ch, exchange, "missing.key", "unroutable");
ch.waitForConfirmsOrDie(5000);
var result = returned.get(5, TimeUnit.SECONDS);
check(result.getReplyCode() == 312, "expected NO_ROUTE");
check("unroutable".equals(result.getProperties().getMessageId()), "wrong returned event");
System.out.println("unroutable: confirm=ACK return=312 queues=0");
send(ch, exchange, "order.paid", "routed"); ch.waitForConfirmsOrDie(5000);
for (String q : List.of(q1, q2)) {
var m = get(ch, q); check("routed".equals(m.getProps().getMessageId()), "wrong route");
ch.basicAck(m.getEnvelope().getDeliveryTag(), false);
}
System.out.println("routed: event=routed copies=2 bindings=order.*,#.paid");
ch.queueDelete(q1); ch.queueDelete(q2); ch.exchangeDelete(exchange);
}
}
case "declare" -> {
String q = "me12.rabbit.declare";
var ch = connection.createChannel(); queue(ch, q);
try {
ch.queueDeclare(q, false, false, false, Map.of());
throw new IllegalStateException("incompatible declaration unexpectedly accepted");
} catch (java.io.IOException e) {
check(replyCode(e) == 406, "expected PRECONDITION_FAILED: " + e);
check(!ch.isOpen(), "failed channel remained open");
System.out.println("declare: reply=406 channelOpen=false");
}
try (var recovered = connection.createChannel()) {
queue(recovered, q); recovered.queueDelete(q);
System.out.println("recovery: newChannel=true compatibleDeclaration=true");
}
}
case "prefetch" -> {
String q = "me12.rabbit.prefetch";
try (var publisher = connection.createChannel()) {
queue(publisher, q); publisher.queuePurge(q); publisher.confirmSelect();
for (int i = 1; i <= 3; i++) send(publisher, "", q, "prefetch-" + i);
publisher.waitForConfirmsOrDie(5000);
var deliveries = new ArrayBlockingQueue<Delivery>(8);
String firstId;
try (var consumer = connection.createChannel()) {
consumer.basicQos(2);
consumer.basicConsume(q, false, (tag, delivery) -> deliveries.add(delivery), tag -> {});
Delivery first = deliveries.poll(5, TimeUnit.SECONDS);
Delivery second = deliveries.poll(5, TimeUnit.SECONDS);
check(first != null && second != null, "two deliveries missing");
check(deliveries.poll(700, TimeUnit.MILLISECONDS) == null, "prefetch exceeded 2");
check(consumer.queueDeclarePassive(q).getMessageCount() == 1, "third message not ready");
System.out.println("prefetch: delivered=2 ready=1 thirdCallback=false");
firstId = first.getProperties().getMessageId();
consumer.basicAck(first.getEnvelope().getDeliveryTag(), false);
Delivery third = deliveries.poll(5, TimeUnit.SECONDS);
check(third != null, "ACK did not release credit");
Set<String> ids = Set.of(firstId, second.getProperties().getMessageId(), third.getProperties().getMessageId());
check(ids.size() == 3, "unexpected duplicate before reconnect");
System.out.println("afterAck: totalDelivered=3 unacked=2");
}
try (var recovered = connection.createChannel()) {
var ids = new HashSet<String>();
for (int i = 0; i < 2; i++) {
var m = get(recovered, q);
check(m.getEnvelope().isRedeliver(), "unacked delivery did not return");
check(!firstId.equals(m.getProps().getMessageId()), "confirmed message returned");
ids.add(m.getProps().getMessageId()); recovered.basicAck(m.getEnvelope().getDeliveryTag(), false);
}
check(ids.size() == 2, "wrong recovered IDs");
check(recovered.basicGet(q, false) == null, "extra delivery");
System.out.println("recovery: redelivered=2 acknowledgedMessageReturned=false");
recovered.queueDelete(q);
}
}
}
case "bad-ack" -> {
var bad = connection.createChannel();
var closed = new CompletableFuture<ShutdownSignalException>();
bad.addShutdownListener(closed::complete);
bad.basicAck(99999, false);
var reason = closed.get(5, TimeUnit.SECONDS);
check(replyCode(reason) == 406 && !bad.isOpen(), "unknown delivery tag did not close channel");
try (var fresh = connection.createChannel()) {
check(fresh.isOpen(), "connection unusable after channel failure");
System.out.println("badAck: reply=406 failedChannelOpen=false newChannelOpen=true");
}
}
case "persist-send" -> {
try (var ch = connection.createChannel()) {
String q = "me12.rabbit.persist"; queue(ch, q); ch.queuePurge(q); ch.confirmSelect();
send(ch, "", q, "persisted"); ch.waitForConfirmsOrDie(5000);
System.out.println("persist: confirmed=true ready=" + ch.queueDeclarePassive(q).getMessageCount());
}
}
case "persist-read" -> {
try (var ch = connection.createChannel()) {
String q = "me12.rabbit.persist"; var m = get(ch, q);
check("persisted".equals(m.getProps().getMessageId()), "persisted ID missing");
ch.basicAck(m.getEnvelope().getDeliveryTag(), false); ch.queueDelete(q);
System.out.println("persist: recoveredEvent=persisted acknowledged=true");
}
}
default -> throw new IllegalArgumentException("unknown mode: " + args[0]);
}
}
}
}程序为每种模式使用独立 me12.rabbit.* 名称,并清理对应实验拓扑。运行前应确认这些名称没有被其他应用使用。
构建使用 Maven 3.9.12 / Temurin 25,显式传入当前 UID/GID 和可写缓存:
docker run --rm --user "$(id -u):$(id -g)" \
-e MAVEN_CONFIG=/maven-cache \
-v "$PWD:/workspace" -v "$PWD/.m2:/maven-cache" -w /workspace \
maven:3.9.12-eclipse-temurin-25 \
mvn -B -ntp '-Dmaven.repo.local=/maven-cache' '-Duser.home=/tmp' \
-pl rabbit -am clean package dependency:copy-dependencies得到 BUILD SUCCESS 后,先运行 route:
docker run --rm --user "$(id -u):$(id -g)" \
--network me12-lab_default -v "$PWD:/workspace:ro" -w /workspace \
eclipse-temurin:25.0.4_7-jdk \
java -cp 'rabbit/target/classes:rabbit/target/dependency/*' \
example.RabbitRuntimeLab routeunroutable: confirm=ACK return=312 queues=0
routed: event=routed copies=2 bindings=order.*,#.paid同一程序先发送不匹配的 missing.key,再发送匹配两个队列的 order.paid。两次发布都得到 confirm,路由结果却不同。
从路由结果读懂发布确认
Confirm ACK 可以和 Return 同时出现
mandatory=true 要求在最终没有目标队列可路由时把消息退回生产者;basic.return 中包含 reply code、exchange、routing key、消息属性和 body。示例收到了 312 NO_ROUTE,说明需要修正 binding 或 routing key,而不是重建网络连接。
publisher confirm 对另一件事作答:Broker 是否完成了本次发布所要求的处理。无路由消息被正确识别并退回,也可能收到 confirm ACK。RabbitMQ 会先发出对应 return,再发 confirm;应用应关联两者,不能只在 confirm 回调里把业务发布标记为成功。发布确认与无法路由消息
发布结果
├── 收到 confirm ACK
│ ├── 已匹配目标:按该队列类型完成确认条件
│ └── mandatory return:这次消息未进入目标队列
├── 收到 confirm NACK:发布侧未按要求完成处理
└── 未收到结果就超时/断线:结果未知,重试须保留稳定事件 ID配置 alternate exchange 后,原交换器可以把未匹配消息交给备用路由。若备用路径成功,行为就不再是示例中的“零路由直接 return”。为备用队列配置容量和观察指标同样重要;备用 exchange 的名称存在,并不能保证其下方有有效 binding。
同步等待、批量等待与在途集合
示例逐次 waitForConfirmsOrDie(5000),可以清楚定位结果,但串行发送的吞吐会受到往返时延限制。业务流量较大时,可批量发布后等待确认,或者用异步 confirm listener 维护在途集合。
channel 在 confirm 模式下为发布分配序号。发送前读取 getNextPublishSeqNo(),把序号与消息身份写入在途表,再 publish;ACK/NACK 回调根据序号处理单条或 multiple=true 的连续范围。必须先登记再发布,防止回复已经到达而表中尚无记录。
在途表还应有数量与字节上限。Broker 变慢时,生产者应暂停接收更多工作或交给有容量约束的持久发布记录,不能无限累积 Java 对象。回调中避免阻塞 I/O 和再次同步等待该 channel 的回复,以免确认处理被自身阻塞。发送异常、return 与 confirm 对同一消息的状态更新要协调,不能重复释放在途计数。
confirm 序号只属于当前 channel 的发布会话。业务 eventId 保留在消息属性或正文中,连接重建后用于关联和去重;不要把一次 channel 内的序号持久化为跨重启的事件身份。
声明检查不会替应用迁移拓扑
队列的名称、durable、exclusive、auto-delete 和 arguments 构成声明条件。重复声明兼容队列可以确认其存在,属性冲突会关闭 channel。队列类型也不能通过重声明从 classic 原地改成 quorum。队列属性与声明一致性
运行同一命令,将末尾模式改成 declare:
docker run --rm --user "$(id -u):$(id -g)" \
--network me12-lab_default -v "$PWD:/workspace:ro" -w /workspace \
eclipse-temurin:25.0.4_7-jdk \
java -cp 'rabbit/target/classes:rabbit/target/dependency/*' \
example.RabbitRuntimeLab declaredeclare: reply=406 channelOpen=false
recovery: newChannel=true compatibleDeclaration=true第一次已声明 durable quorum queue,第二次用同名非 durable 参数触发 406。恢复时修正参数并新建 channel;继续使用旧对象只会遇到关闭状态。生产迁移通常需要新队列、转移发布或订阅、验证旧队列排空后再删除,不能让所有实例启动时各自尝试“改一下声明”。
消费确认与 Prefetch 控制什么
Delivery tag 属于产生投递的 channel
手动确认的 deliveryTag 是 Broker 在 channel 上分配的投递标识。处理成功后,消费者在产生该投递的 channel 上调用 basicAck(tag, false)。它与消息属性中的 messageId 无关;同一业务事件重投时可以获得新的 delivery tag。
basicReject 处理单条投递,basicNack 还支持批量范围。requeue=true 把消息重新交给队列,适合有恢复条件的失败;持续解析失败却立即 requeue,会让同一消息高频循环。requeue=false 根据配置进入死信路径或被丢弃,死信转发的可靠性还取决于队列类型与策略。重试与死信拓扑
多个业务线程乱序完成时,要特别谨慎使用 multiple=true。如果 tag 7 尚未完成,tag 8、9 先完成,直接 ACK 9 且 multiple=true 会把尚未完成的 7 也确认。可以逐条确认,或维护连续完成的位置,只有前面的所有业务都成功才推进批量 ACK。
故意确认一个从未投递过的 tag,bad-ack 模式会得到:
badAck: reply=406 failedChannelOpen=false newChannelOpen=true该模式可用前面的 Java 容器命令运行,只将末尾参数替换为 bad-ack。异常来自通道协议状态,同一个 connection 仍可以创建新 channel。消费确认发生在 SQL 提交之前或之后的不同结果,参见 消息提交与确认实验。
Prefetch 限制已交付但未确认的数量
在 RabbitMQ 中,basicQos(2) 默认对随后建立的每个消费者设置未确认投递上限。一个 channel 上创建两个消费者时,两者可以分别持有 2 条,而不是自动共同分享总数 2。global=true 表达 channel 级上限,quorum queue 不支持这种 global QoS,应使用 per-consumer 设置。Consumer Prefetch
预取限制控制消息从 Broker 进入消费者窗口的速度。它既不是线程池大小,也不是业务数据库的连接数。单线程逐条处理却设置很大的 prefetch,会让大量消息滞留在客户端;设置太小则可能使消费者在网络交付之间空转。选择时应结合消息大小、并发工作数、平均处理时间和故障后可接受的重投批次。
prefetch 模式使用真正的 basicConsume 回调:
docker run --rm --user "$(id -u):$(id -g)" \
--network me12-lab_default -v "$PWD:/workspace:ro" -w /workspace \
eclipse-temurin:25.0.4_7-jdk \
java -cp 'rabbit/target/classes:rabbit/target/dependency/*' \
example.RabbitRuntimeLab prefetchprefetch: delivered=2 ready=1 thirdCallback=false
afterAck: totalDelivered=3 unacked=2
recovery: redelivered=2 acknowledgedMessageReturned=false先发送三条消息,回调队列只得到两条;在受控观察窗口内没有第三次回调,Broker 查询也显示还有一条 ready。ACK 第一条后,第三条回调才到达。关闭消费 channel 后,剩余两条未确认消息重投,已经确认的第一条没有再出现。
prefetch = 2
起始:队列 [1, 2, 3] 消费者未确认窗口 [空, 空]
交付:队列 [3] 消费者未确认窗口 [1, 2]
ACK1:队列 [空] 消费者未确认窗口 [3, 2]
关闭:未确认的 2 和 3 返回队列,等待新的消费交付这里的回调只把投递放入一个有界集合,真正的确认由实验主线程发出。实际业务应明确回调线程与处理线程的分工,避免把无限线程池当作流控手段。自动确认模式也不适合拿来观察这个未确认窗口:它可能在业务尚未完成时就失去重投机会。
连接恢复与通道错误要分别处理
本实验关闭自动连接恢复,是为了让每次关闭与重新订阅直接可见。生产应用使用 Java 客户端自动恢复时,网络连接恢复、channel 重建和拓扑恢复由客户端按其规则执行,但两类情况仍需要应用处理。
channel 级异常通常意味着声明、权限或投递标识使用错误。客户端不会靠网络重连修好一个不兼容声明。连接尚未建立的首次连接失败,也应由启动流程给出有上限的重试与就绪状态。
已经启用自动恢复的 channel 对象则可能是恢复代理。Broker 重建通道后会重置投递标签,Java 客户端会调整标签并抑制过期 ACK,不能把协议底层标签重置简单理解成“恢复后必然发送旧 tag 把通道关掉”。应用仍须容忍重投,并且自行保存未得到发布确认的消息;自动恢复不会可靠缓存所有恢复期间的业务发布。相关行为见前面的 Java API Guide。
Quorum Queue 怎样保存与恢复消息
队列类型决定存储方式
| 类型 | 常见用途 | 持久化与复制特点 |
|---|---|---|
| classic queue | 不要求复制的队列、部分临时通信 | 4.x 的 classic queue 为非复制队列;旧 mirrored classic 方案已移除 |
| quorum queue | 需要复制和故障恢复的业务队列 | durable,使用 Raft,多数派提交决定可用性与确认 |
| stream | 保留历史、重放和高吞吐日志读取 | 追加存储,消费位置和保留策略与普通队列不同 |
durable 标记描述队列声明的持久性,消息持久属性描述消息要求。确认条件还取决于具体存储类型:quorum queue 本身使用持久复制存储,不会因某条消息省略 persistent 属性就退化为内存队列。quorum queue 不能用作 exclusive 临时队列,也不支持上文的 global QoS。Quorum queue 特性与限制
先检验单节点进程重启。persist-send 发布一条持久消息并等待确认,随后重启当前实验 Broker,再运行 persist-read:
run_rabbit() {
docker run --rm --user "$(id -u):$(id -g)" \
--network me12-lab_default -v "$PWD:/workspace:ro" -w /workspace \
eclipse-temurin:25.0.4_7-jdk \
java -cp 'rabbit/target/classes:rabbit/target/dependency/*' \
example.RabbitRuntimeLab "$1"
}
run_rabbit persist-send
docker compose restart rabbit
docker compose up -d --wait --wait-timeout 120 rabbit
run_rabbit persist-readpersist: confirmed=true ready=1
persist: recoveredEvent=persisted acknowledged=true数据保存在该实例的命名卷中。这验证了进程重启后的保存结果,不能覆盖磁盘永久丢失或多个节点的选举与复制;后者需要实际建立多副本队列。
三副本实验先确认成员,再制造故障
工程的 compose.rabbit-ha.yaml 创建独立项目 me12-rabbit-ha,包含 rabbit1、rabbit2、rabbit3,每个节点各有数据卷,使用相同实验 cookie 和静态节点发现。它与前面的单节点项目互不替代,不发布宿主端口。三个服务分别限额 512 MiB,宿主应留出容器、Java 客户端与磁盘余量。
共享 config/rabbit-ha.conf 的集群部分为:
cluster_formation.peer_discovery_backend = classic_config
cluster_formation.classic_config.nodes.1 = rabbit@rabbit1
cluster_formation.classic_config.nodes.2 = rabbit@rabbit2
cluster_formation.classic_config.nodes.3 = rabbit@rabbit3
cluster_partition_handling = ignore节点还导入与单节点相同的普通 app / lab 定义。cluster_partition_handling=ignore 在这套只测试 quorum queue 的拓扑中把队列一致性交给 Raft;它不是对包含其他队列类型的既有集群的通用配置建议。节点名、名称解析、Erlang cookie、集群端口与数据目录是建立集群所需的条件。RabbitMQ 集群
启动并检查三个节点都属于同一集群:
docker compose -f compose.rabbit-ha.yaml up -d --wait --wait-timeout 120
docker compose -f compose.rabbit-ha.yaml exec rabbit1 rabbitmqctl cluster_statusRunning Nodes 应同时列出 rabbit@rabbit1、rabbit@rabbit2、rabbit@rabbit3。若只显示一个,先排查节点发现和 cookie,不能继续把三个独立 Broker 当成三副本队列。
入口 RabbitQuorumLab 建队列时显式给出三成员初始大小:
ch.queueDeclare("me12.rabbit.quorum", true, false, false,
Map.of("x-queue-type", "quorum", "x-quorum-initial-group-size", 3));定义运行命令并发布第一条:
run_quorum() {
docker run --rm --user "$(id -u):$(id -g)" -e RABBIT_HOST=rabbit1 \
--network me12-rabbit-ha_default -v "$PWD:/workspace:ro" -w /workspace \
eclipse-temurin:25.0.4_7-jdk \
java -cp 'rabbit/target/classes:rabbit/target/dependency/*' \
example.RabbitQuorumLab "$1"
}
run_quorum prepare
docker compose -f compose.rabbit-ha.yaml exec rabbit1 \
rabbitmq-queues quorum_status --vhost lab me12.rabbit.quorum应看到三个 voter,其中一个为 leader,其余为 follower。日志索引随实际操作变化,不需要与固定数值相等;应检查成员、运行状态和提交进展。
多数派丢失后,确认超时仍可能留下消息
停止 rabbit3,只保留两个在线成员,随后发布:
docker compose -f compose.rabbit-ha.yaml stop rabbit3
run_quorum one-down预期 quorum: event=one-down confirmed=true。三个成员中的两个仍构成多数派,可以继续提交。再停止 rabbit2,只留下 rabbit1:
docker compose -f compose.rabbit-ha.yaml stop rabbit2
run_quorum minority预期出现:
minority: confirmTimeout=true outcome=unknown doNotAssignNewEventId=true等待 4 秒仍未收到 confirm,程序将发送结果记为未知。少数派无法形成新的多数派提交;消息是否到达或会在恢复后被提交,要继续检查恢复结果。
失败时清理同样需要超时。少数派队列可能让 channel 的正常关闭继续等待,实验最终使用 connection.abort(1000) 限制关闭等待。这个主动断开可能附带 Socket closed 警告;关键结果仍是 confirm 未完成以及已关闭连接,不能把警告改写成业务明确拒绝。
恢复两个节点并读取队列:
docker compose -f compose.rabbit-ha.yaml up -d --wait --wait-timeout 120 rabbit2 rabbit3
run_quorum drain一次真实结果为:
recovered: confirmedIdsPresent=true minorityObserved=true total=3两条已确认消息全部保留,少数派期间没有得到确认的第三条也在恢复后出现。程序强制检查已确认的两条存在;第三条是否出现按实际结果报告,不能要求所有故障时序都得到同样的 minorityObserved。这正是发布端需要稳定 eventId 和消费去重的原因。
队列成员数量与集群节点数量不是同一个设置。扩容集群后,已有 quorum queue 的成员不会仅因新节点上线就自然平均分布到所有节点;应按成员管理和再平衡规则操作。生产通常以跨独立故障域的奇数成员形成多数派;把三个容器放在同一宿主,能够验证协议故障路径,但宿主故障仍会同时影响三者。
完成后停止这套集群,保留数据可重新启动:
docker compose -f compose.rabbit-ha.yaml stop只有确认数据全部可丢弃时,才执行 docker compose -f compose.rabbit-ha.yaml down --volumes。单节点项目则使用不带 -f 的 docker compose stop 或独立清理命令,避免误删另一个项目。
连接、容量与故障定位
心跳发现失联,资源告警限制发布
心跳用于缩短死连接识别时间,TCP 连接存在不代表对端应用仍能有效通信。超时值过小可能在短暂网络拥塞或调度停顿时造成误判,进而放大重连和重投;应结合网络特征和客户端支持配置,并观察连接关闭原因。Heartbeats
Broker 触发内存或磁盘告警后,可能阻塞发布连接以限制进一步资源占用。生产者看到发送变慢时,应先查看告警和 blocked 通知,而不是立即增加重试线程。发布和消费共用同一连接还会使流控与恢复更难分析,重要业务通常将两类连接分开。内存与磁盘告警
docker compose exec rabbit rabbitmq-diagnostics -q alarms
docker compose exec rabbit rabbitmqctl list_connections name state channels
docker compose exec rabbit rabbitmqctl list_queues -p lab \
name type messages_ready messages_unacknowledged consumers无告警时 alarms 不应列出资源警报。队列 ready 高且 consumers=0,先恢复订阅;unacknowledged 长期居高,优先看消费者业务耗时、死锁、ACK 位置与 prefetch;publish 速率持续高于消费完成速率,则需要控制入口或增加有效处理能力。指标的含义和采集方式见 RabbitMQ Monitoring。
沿错误所在对象恢复
| 现象 | 先看的对象 | 处理与再验证 |
|---|---|---|
| confirm ACK 但收到 312 | exchange、routing key、bindings | rabbitmqctl list_bindings -p lab,修正绑定后用 route 模式验证目标队列各收到正确 ID |
| 406 声明冲突 | queue/exchange 原有属性 | 查 name type durable,修正配置并新建 channel;类型迁移使用新资源 |
| 406 unknown delivery tag | 发出 ACK 的 channel 与标签 | 核对是否跨通道、重复 ACK 或批量范围越过未完成工作;用 bad-ack 对照实际通道关闭 |
| 消费停止在一批消息 | unacked、处理线程、下游事务 | 检查线程与数据库等待,修复后确认已有结果,观察窗口恢复流动 |
| 重连后重复处理 | 断开时的确认与 DB 状态 | 保留相同 eventId,按数据库幂等结果确认,避免重复扣减 |
| quorum queue 不再确认 | 队列成员与在线多数派 | 查询 quorum_status / cluster_status,恢复足够成员后核对已确认消息,未知发布按去重重试 |
| 发送等待但无明显网络错误 | 内存/磁盘告警、blocked 连接、在途容量 | 先解除资源告警与入口过载,再观察确认延迟回落,禁止无限扩大在途集合 |
发布前应固定拓扑命名、队列类型、消息大小、确认策略和角色权限,并让应用就绪状态反映关键连接是否可用。声明与收发权限分离、TLS、备份与升级可在 RabbitMQ 部署与使用中继续配置;业务消费的数据库幂等与重复处理则见 投递、幂等与顺序。
权威资料与规范地址
| 资料 | 完整地址 |
|---|---|
| Java Client API Guide:连接、通道、消费与自动恢复 | https://www.rabbitmq.com/client-libraries/java-api-guide |
| Exchanges:类型、binding、默认与备用交换器 | https://www.rabbitmq.com/docs/exchanges |
| Publisher Confirms / Consumer Acknowledgements | https://www.rabbitmq.com/docs/confirms |
| Queues:声明属性、顺序与生命周期 | https://www.rabbitmq.com/docs/queues |
| Consumer Prefetch:消费者与 channel 上限 | https://www.rabbitmq.com/docs/consumer-prefetch |
| Quorum Queues:多数派、成员与功能条件 | https://www.rabbitmq.com/docs/quorum-queues |
| Clustering Guide:节点、名称解析与集群建立 | https://www.rabbitmq.com/docs/clustering |
| Heartbeats:失联识别与配置 | https://www.rabbitmq.com/docs/heartbeats |
| Memory and Disk Alarms:发布连接流控 | https://www.rabbitmq.com/docs/alarms |
| Monitoring:指标与诊断对象 | https://www.rabbitmq.com/docs/monitoring |
