RocketMQ 运行链:存储索引、消费确认与事务回查
RocketMQ 将消息追加到 Broker 的存储,再通过 Topic 与队列组织消费。普通消息关注收发和确认,FIFO 消息增加顺序约束,事务消息则先保存暂不可消费的记录,等本地事务结果确定后再决定是否对消费者可见。三种消息共享基础设施,但处理条件各有不同。
部署角色与首次消息收发
NameServer、Broker 与 Proxy
Java SDK 5.2.1
└── gRPC → Proxy :8081
├── 向 NameServer :9876 查询路由
└── Remoting → Broker :10911
├── Topic / MessageQueue
├── CommitLog 与消费索引
├── 消费进度、重投状态
└── 事务半消息与回查
传统 Remoting 客户端
├── 向 NameServer 查询路由
└── 直接连接路由中的 BrokerNameServer 保存 Broker 注册的路由信息,业务消息存放在 Broker。现代 Java SDK 通过 Proxy 使用 gRPC;传统 DefaultMQProducer、DefaultMQPushConsumer 属于另一套客户端接口。配置 SDK endpoint 时,需要先知道正在使用哪套协议,不能把 NameServer 的 9876 端口填进 gRPC endpoint。Java SDK 与服务端要求
下面固定服务端镜像 apache/rocketmq:5.5.0、现代客户端 rocketmq-client-java:5.2.1。服务端镜像内置自己的 Java 运行环境;实验 Maven 用 Java 25 构建到 Java 17 字节码,客户端使用 Temurin 17.0.20。服务端、Proxy 和客户端版本应分别记录,不按相同数字推断兼容。官方 Docker 部署
下载工程并准备权限
下载完整实验工程,解压进入 message-event-lab。使用 Linux shell、Docker Engine、Compose v2+,预留约 3 GiB 容器内存,由有 Docker 权限的普通宿主用户操作:
message-event-lab/
├── compose.rocket.yaml
├── init-rocket.sh
├── config/
│ ├── rocket-broker.conf
│ ├── start-rocket-broker.sh
│ ├── rocket-proxy.json
│ └── init.sql
└── rocket/
├── pom.xml
└── src/main/
├── java/example/RocketRuntimeLab.java
└── resources/rocketmq.logback.xmlCompose 为 NameServer、Broker、Proxy 和 PostgreSQL 分别设置内存上限。Broker 数据放在 rocket-store 命名卷;镜像中的运行用户为 UID/GID 3000。一次性的 volume-init 容器以 root 为这个指定卷根目录设置所有权,后续 Broker 始终使用镜像的普通用户。不能把这个初始化动作描述成所有容器都不需要 root。
数据库为 PostgreSQL 18.6,仅事务实验使用普通应用角色 app 和 schema lab。初始化角色需要数据库管理权限,Java 不使用引导超级用户。实验口令与明文端口仅存在于专用 Docker 网络,不应直接开放到公网。
Broker 使用 1 GiB 容器上限、384 MiB 最大堆和 256 MiB 直接内存;Proxy 上限为 768 MiB。两个进程都通过 ActiveProcessorCount=4 与 io.netty.allocator.numDirectArenas=2 控制本机小实验的线程/分配池规模。容器能看到很多宿主 CPU 时,只缩小堆与直接内存而保留默认分配池,连续收发可能出现 OutOfDirectMemoryError,客户端表现为 ACK 50001 或超时。排查时应查 Broker 的 pop.log 和具体内存错误,而非只看 Docker OOMKilled。堆、直接内存、线程栈、原生库和文件映射应一起纳入容器预算。Netty 分配池参数
Broker 地址必须与实际回查路径一致
核心存储配置是:
brokerClusterName=Me12Cluster
brokerName=me12-broker
brokerId=0
brokerRole=ASYNC_MASTER
flushDiskType=SYNC_FLUSH
namesrvAddr=namesrv:9876
listenPort=10911
autoCreateTopicEnable=false
autoCreateSubscriptionGroup=false
storePathRootDir=/home/rocketmq/store
storePathCommitLog=/home/rocketmq/store/commitlog
fileReservedTime=24
transactionTimeOut=6000
transactionCheckInterval=5000brokerId=0 表示这里的主 Broker。只有这一份 Broker 数据,SYNC_FLUSH 要求本地刷盘,仍没有另一台机器的副本保护。读写队列数也不表示副本数量。
启动脚本将只读基础配置复制到容器临时目录,取得当前容器 IP 并追加 brokerIP1:
#!/bin/sh
set -eu
cp /home/rocketmq/broker.conf /tmp/me12-broker.conf
broker_address=$(hostname -i | awk '{print $1}')
test -n "$broker_address"
printf '\nbrokerIP1=%s\n' "$broker_address" >> /tmp/me12-broker.conf
exec sh mqbroker -c /tmp/me12-broker.conf这个地址只供同一 Docker 网络的 Proxy 和运维客户端使用,不复制到宿主配置里,也不在文中固定某台机器的 IP。生产网络应设置可由实际调用方访问、且与服务端回查信息一致的地址。
在 5.5.0 的 cluster Proxy 中,事务服务按地址字符串维护 Broker 名称映射。将 brokerIP1 写成服务名 broker 时,普通发送可以成功,但回查里的数值地址可能无法命中这份映射,出现 add transaction data failed。检查路由的地址与回查请求中的 brokerAddr 是否一致,可从 ClusterTransactionService 源码 追到具体查找过程。
显式创建 Topic 和消费组
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户'; exit 1; }
mkdir -p .m2
sh init-rocket.sh脚本先启动 NameServer、数据库和 Broker,轮询 mqadmin clusterList,只有实际注册信息包含 me12-broker 后才创建业务资源和启动 Proxy。若超时,它输出服务日志并非零退出,不会继续把后续连接错误当作消息投递失败。
创建规则如下,完整循环在脚本中:
| Topic | message.type | 消费组 | consumeMessageOrderly |
|---|---|---|---|
| me12_normal | NORMAL | me12-me12_normal-group | false |
| me12_fifo | FIFO | me12-me12_fifo-group | true |
| me12_transaction | TRANSACTION | me12-me12_transaction-group | false |
例如 FIFO 需要同时配置 Topic 类型与消费组:
docker compose -f compose.rocket.yaml exec -T broker sh mqadmin updateTopic \
-n namesrv:9876 -c Me12Cluster -t me12_fifo -r 1 -w 1 -a '+message.type=FIFO'
docker compose -f compose.rocket.yaml exec -T broker sh mqadmin updateSubGroup \
-n namesrv:9876 -c Me12Cluster -g me12-me12_fifo-group -o truemqadmin 部分子命令即使打印异常也可能以 0 退出。初始化脚本同时核对 create topic to ... success. 和 create subscription group to ... success.,不能只看 shell 返回值。消费组描述应包含 consumeMessageOrderly=true。-r 1 -w 1 为一个逻辑读写队列,便于观察,不产生副本。已有 Topic 的类型变更涉及产品约束,应创建新的适用 Topic,而不是靠自动建 Topic 掩盖类型错误。
构建与运行:
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 rocket -am clean package dependency:copy-dependencies
run_rocket() {
docker run --rm --user "$(id -u):$(id -g)" --network me12-rocket_default \
-v "$PWD:/workspace:ro" -w /workspace eclipse-temurin:17.0.20_8-jdk \
java '-Duser.home=/tmp' \
'-Drocketmq.logback.configurationFile=/workspace/rocket/target/classes/rocketmq.logback.xml' \
-cp 'rocket/target/classes:rocket/target/dependency/*' example.RocketRuntimeLab "$@"
}
run_rocket normal先看到 BUILD SUCCESS,随后应看到:
normal: body=event-1 acknowledged=true代码验证收到的正文为 event-1,再提交 ACK,并检查该组没有额外可见消息。SDK 日志配置显式写到控制台,避免无 home 目录的容器 UID 尝试在根目录创建日志文件。
完整 Java 入口
同一消费组始终订阅 *,每次运行产生的随机 tag 只用于标记本轮消息,事务 ID 也带有本轮前缀。不要随运行次数修改同组订阅条件:Broker 的过滤状态和消费进度会跨客户端实例保存。每次应等当前模式成功结束后再运行下一模式;异常退出留下的消息先按 ID 核对,不能假定下一轮队列为空。
package example;
import java.nio.charset.StandardCharsets;
import java.sql.DriverManager;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.rocketmq.client.apis.ClientConfiguration;
import org.apache.rocketmq.client.apis.ClientServiceProvider;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.SimpleConsumer;
import org.apache.rocketmq.client.apis.message.MessageView;
import org.apache.rocketmq.client.apis.producer.TransactionResolution;
public final class RocketRuntimeLab {
static final ClientServiceProvider API=ClientServiceProvider.loadService();
static final ClientConfiguration CONFIG=ClientConfiguration.newBuilder()
.setEndpoints("proxy:8081").enableSsl(false).setRequestTimeout(Duration.ofSeconds(10)).build();
static void check(boolean ok,String why) { if(!ok) throw new IllegalStateException(why); }
static SimpleConsumer consumer(String topic) throws Exception {
return API.newSimpleConsumerBuilder().setClientConfiguration(CONFIG)
.setConsumerGroup("me12-"+topic+"-group").setAwaitDuration(Duration.ofSeconds(1))
.setSubscriptionExpressions(Map.of(topic,new FilterExpression("*"))).build();
}
static MessageView receive(SimpleConsumer consumer) throws Exception {
long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(60);
while(System.nanoTime()<deadline) {
var messages=consumer.receive(1,Duration.ofSeconds(10));
if(!messages.isEmpty()) return messages.get(0);
}
throw new IllegalStateException("message not received within 60 seconds");
}
static String body(MessageView message) {
var buffer=message.getBody(); var bytes=new byte[buffer.remaining()];buffer.get(bytes);
return new String(bytes,StandardCharsets.UTF_8);
}
static java.sql.Connection database() throws Exception {
return DriverManager.getConnection("jdbc:postgresql://postgres:5432/lab","app","app-lab-only");
}
static void status(String id,String state) throws Exception {
try(var c=database();var s=c.createStatement()) {
s.execute("CREATE TABLE IF NOT EXISTS rocket_local_tx(tx_id text PRIMARY KEY,state text NOT NULL)");
s.execute("CREATE TABLE IF NOT EXISTS rocket_business(tx_id text PRIMARY KEY,amount int NOT NULL)");
c.setAutoCommit(false);
try {
if(state.equals("COMMITTED")) {
try(var p=c.prepareStatement("INSERT INTO rocket_business VALUES (?,10) ON CONFLICT DO NOTHING")) {
p.setString(1,id);p.executeUpdate();
}
}
try(var p=c.prepareStatement("INSERT INTO rocket_local_tx VALUES (?,?) ON CONFLICT(tx_id) DO UPDATE SET state=excluded.state")) {
p.setString(1,id);p.setString(2,state);p.executeUpdate();
}
c.commit();
} catch (Exception | Error failure) {
try { c.rollback(); }
catch (java.sql.SQLException rollbackFailure) { failure.addSuppressed(rollbackFailure); }
throw failure;
}
}
}
public static void main(String[] args) throws Exception {
check(args.length==1,"mode: normal|redelivery|fifo|transaction");
String mode=args[0],tag="run"+UUID.randomUUID().toString().replace("-","");
if(mode.equals("transaction")) { transactions(tag);return; }
check(List.of("normal","redelivery","fifo").contains(mode),"unknown mode");
String topic=mode.equals("fifo")?"me12_fifo":"me12_normal";
try(var c=consumer(topic);var p=API.newProducerBuilder().setClientConfiguration(CONFIG).setTopics(topic).build()) {
int count=mode.equals("fifo")?3:1;
for(int i=1;i<=count;i++) {
var builder=API.newMessageBuilder().setTopic(topic).setTag(tag).setKeys(tag)
.setBody(("event-"+i).getBytes(StandardCharsets.UTF_8));
if(mode.equals("fifo")) builder.setMessageGroup("order-42");
p.send(builder.build());
}
var first=receive(c);check(body(first).equals("event-1"),"first value mismatch");
if(mode.equals("redelivery")) {
var again=receive(c);
check(first.getMessageId().equals(again.getMessageId()),"different message returned");
check(again.getDeliveryAttempt()>=2,"delivery attempt did not increase");
c.ack(again);
System.out.println("redelivery: sameMessageId=true attempt="+again.getDeliveryAttempt()+" acknowledged=true");
} else if(mode.equals("fifo")) {
check(c.receive(1,Duration.ofSeconds(10)).isEmpty(),"same FIFO group advanced before ACK");
var values=new ArrayList<String>();values.add(body(first));c.ack(first);
for(int i=2;i<=count;i++) {var next=receive(c);values.add(body(next));c.ack(next);}
check(values.equals(List.of("event-1","event-2","event-3")),"FIFO sequence changed: "+values);
System.out.println("fifo: group=order-42 blockedBeforeAck=true values="+values);
} else { c.ack(first);System.out.println("normal: body=event-1 acknowledged=true"); }
check(c.receive(1,Duration.ofSeconds(10)).isEmpty(),"unexpected extra matching message");
}
}
static void transactions(String tag) throws Exception {
String topic="me12_transaction";var checks=new AtomicInteger();var failed=new AtomicInteger();
String committed=tag+"-commit",rolledBack=tag+"-rollback",recover=tag+"-recover";
status(committed,"PENDING");status(rolledBack,"PENDING");status(recover,"PENDING");
try(var c=consumer(topic);var p=API.newProducerBuilder().setClientConfiguration(CONFIG).setTopics(topic)
.setTransactionChecker(message->{
checks.incrementAndGet();
try(var db=database();var query=db.prepareStatement("SELECT state FROM rocket_local_tx WHERE tx_id=?")) {
query.setString(1,message.getKeys().iterator().next());
try(var rows=query.executeQuery()) {
if(!rows.next()) return TransactionResolution.UNKNOWN;
return switch(rows.getString(1)) {
case "COMMITTED" -> TransactionResolution.COMMIT;
case "ROLLED_BACK" -> TransactionResolution.ROLLBACK;
default -> TransactionResolution.UNKNOWN;
};
}
} catch(Exception error) {failed.incrementAndGet();return TransactionResolution.UNKNOWN;}
}).build()) {
for(String id:List.of(committed,rolledBack,recover)) {
var tx=p.beginTransaction();
var message=API.newMessageBuilder().setTopic(topic).setTag(tag).setKeys(id)
.setBody(id.getBytes(StandardCharsets.UTF_8)).build();
p.send(message,tx);
if(id.equals(committed)) {status(id,"COMMITTED");tx.commit();}
else if(id.equals(rolledBack)) {status(id,"ROLLED_BACK");tx.rollback();}
else {status(id,"COMMITTED");} // Deliberately omit the second-phase request.
}
var values=new ArrayList<String>();
for(int i=0;i<2;i++) {var message=receive(c);values.add(body(message));c.ack(message);}
check(values.size()==2&&values.contains(committed)&&values.contains(recover),"unexpected transaction visibility");
check(checks.get()>0&&failed.get()==0,"database-backed transaction checker not successful");
check(c.receive(1,Duration.ofSeconds(10)).isEmpty(),"rolled-back message visible");
try(var db=database();var query=db.prepareStatement("SELECT count(*) FROM rocket_business WHERE tx_id IN (?,?,?)")) {
query.setString(1,committed);query.setString(2,rolledBack);query.setString(3,recover);
try(var rows=query.executeQuery()) {check(rows.next()&&rows.getInt(1)==2,"business rows differ from committed decisions");}
}
System.out.println("transaction: directCommit=true rollbackInvisible=true recoveredByCheck=true checks="+checks.get());
}
}
}receive 每次最多取得一条消息,在本地 60 秒窗口内开始轮询;单次正在执行的 SDK 调用仍受请求超时控制,函数总耗时可能略超过这个窗口。网络失败沿 SDK 异常返回,等待窗口结束仍没有消息时,程序也会非零退出。排查时应区分连接失败与订阅成功但暂时无消息。
从 CommitLog 找到消费队列
物理追加与逻辑索引
Broker 的 CommitLog
├── 物理偏移 A:Topic orders,queue 0,正文与属性
├── 物理偏移 B:Topic payments,queue 1,正文与属性
└── 物理偏移 C:Topic orders,queue 0,正文与属性
orders / queue 0 的 ConsumeQueue
├── 逻辑位置 0 → 物理偏移 A + 长度 + tag 信息
└── 逻辑位置 1 → 物理偏移 C + 长度 + tag 信息
payments / queue 1 的 ConsumeQueue
└── 逻辑位置 0 → 物理偏移 B + 长度 + tag 信息CommitLog 保存完整消息记录,消费索引按 Topic 与 queue 定位相应记录。不同 Topic 的消息可以共享同一物理追加文件,客户端按消费队列的逻辑位置读取,Broker 再定位到完整数据。
5.5.0 标准 ConsumeQueue 实现中的一条索引占 20 字节:物理偏移 8 字节、消息长度 4 字节、tag 哈希或扩展索引信息 8 字节。它是索引条目大小,不是完整消息大小。标准 ConsumeQueue 源码
BatchConsumeQueue 则采用另一种格式,增加存储时间、批次基础位置、批次数量等字段;5.5.0 该实现的单位大小为 46 字节。分析存储容量和二进制文件时,先确认实际队列类型,不能对所有 RocketMQ 存储套用 20 字节公式。BatchConsumeQueue 源码
物理位置、逻辑位置与业务身份
| 标识 | 对应对象 | 主要用途 |
|---|---|---|
| CommitLog 物理 offset | Broker 文件中的记录位置 | 存储定位、恢复和底层诊断 |
| queue offset | Topic/queue 内的逻辑读取位置 | 读取与消费进度 |
| message ID | 具体客户端与协议定义的消息身份 | 查找消息、关联收发 |
| message key | 应用提供的可检索业务键 | 按订单号等检索消息 |
| 业务 event ID / tx ID | 应用的一次事件或操作 | 去重、事务结果查询 |
一次重发可能生成新的 message ID,业务 event ID 应按业务重复规则保持稳定。SimpleConsumer 的某次交付还携带确认所需的 receipt handle;应用应拿当前收到的 MessageView 执行 ack,不自行用旧消息 ID 构造确认。
查看真实队列与组进度:
docker compose -f compose.rocket.yaml exec -T broker sh mqadmin topicStatus \
-n namesrv:9876 -t me12_normal
docker compose -f compose.rocket.yaml exec -T broker sh mqadmin consumerProgress \
-n namesrv:9876 -g me12-me12_normal-group最大位置随重复执行增长,不固定比对某一个 offset。ACK 通常推进该组的消费状态,物理消息仍按存储保留规则回收;消费结束不表示 CommitLog 已立即缩小。消息存储与清理
确认、不可见时间与 FIFO 消费
SimpleConsumer 把控制权交给应用
PushConsumer 由 SDK 驱动监听器并处理收取、并发及结果提交;SimpleConsumer 则由应用主动 receive、处理、ack,并管理不可见时间。选择 SimpleConsumer 时,线程、处理超时和在途数量都要有明确限制。消费者类型
receive(1, Duration.ofSeconds(10)) 中的十秒表示取到消息后的不可见时间。业务超过这个时间仍未确认时,其他消费者可能重新得到它。较长任务可以调用相应的不可见时间变更 API,但延长应有总时限,不能用无限续期隐藏卡死任务。
run_rocket redelivery预期关键结果:
redelivery: sameMessageId=true attempt=2 acknowledged=true程序第一次读取后不 ACK,继续等待下一次可读取的消息。它比较两次真实 message ID,并断言 delivery attempt 至少为 2;最后使用新交付的 MessageView 确认。具体尝试次数可能因环境延迟更多,核心条件是确实发生重投,并完成后续确认。
这个不可见时间实验中观察到同一 message ID,不能据此概括所有重试实现都保留同一传输 ID。业务 Inbox 应使用信封里的稳定事件身份,操作窗口见 投递、幂等与顺序。
FIFO 同时依赖生产与消费设置
run_rocket fifo程序由同一个生产者依次发送 event-1、event-2、event-3,使用同一 message group order-42。消费者先收到 event-1,故意暂不确认并再调用 receive,检查同组后续消息没有提前到达;ACK 后继续读取:
fifo: group=order-42 blockedBeforeAck=true values=[event-1, event-2, event-3]这条路径同时依赖 FIFO Topic、消费组 -o true、同一 message group 的串行发送,以及消费者按收到—处理—确认的顺序执行。只给消息设置 message group 而使用普通并发消费组,后续交付仍可能继续推进。FIFO 的生产与消费条件
SimpleConsumer 可以一批拉取多条消息,业务仍要负责批内顺序。实验把最大条数设为 1,便于观察 ACK 对推进的影响;实际应用若批量异步执行,不应把“存储顺序”直接当作“业务完成顺序”。
一个 message group 的失败会拖住后续相关消息,重试上限结束后的处理也要结合业务决定。订单可以按订单号分组,不宜让所有订单共享一个巨大 group。不同 group 的事件没有全局先后关系;确实需要跨对象排序时,应设计更高层的业务协调。
用数据库结果回答事务回查
半消息等候本地事务结果
发送半消息 ──→ 本地数据库事务
├── 业务行 + COMMITTED 状态共同提交
│ └── 二阶段 commit → 消费者可见
├── 本地操作取消,记录 ROLLED_BACK
│ └── 二阶段 rollback → 不向消费者交付
└── 二阶段请求缺失
└── Broker 回查 → Producer 查询持久状态
├── COMMITTED → COMMIT
├── ROLLED_BACK → ROLLBACK
└── PENDING / 暂不可查 → UNKNOWN半消息先由 Broker 保存,正常消费者暂时无法消费。生产者完成本地业务事务后通知 Broker 结果;缺失二阶段结果时,Broker 发起检查,客户端 TransactionChecker 根据本地持久记录作答。
这套协议把“本地事务已经提交,却没有成功通知 Broker”的情况变成可恢复查询。消费者自己的数据库更新仍发生在消息可见之后,需要单独处理重复、失败和补偿;没有把生产者数据库与所有消费者放进同一个原子事务。事务消息行为
运行三种真实决定路径
run_rocket transaction代码准备三个独立事务 ID。直接提交与待回查事务都在同一个 PostgreSQL 事务中写入 rocket_business 业务行和 rocket_local_tx=COMMITTED 状态。取消事务只保留 ROLLED_BACK 决定,没有业务行。第三个事务故意不调用 SDK 的二阶段 commit/rollback。
预期:
transaction: directCommit=true rollbackInvisible=true recoveredByCheck=true checks=1checks 来自 TransactionChecker 的实际回调次数,可以大于 1。程序还验证收到的两条正文分别属于直接提交和回查恢复的事务,未收到回滚事务;最后查询数据库,确认三个事务 ID 对应的业务行恰好有两条。
回查通过消息 key 取得 tx ID,再查询 PostgreSQL。记录不存在、仍为 PENDING 或数据库暂时不可用时返回 UNKNOWN;查询失败会计数,实验要求失败计数为 0。生产处理应为 UNKNOWN 设置监控与人工核对时限,不应把未知结果一律当成回滚。
超时、回查次数与记录保存期
示例把 transactionCheckInterval 缩短为五秒以便观察。这个周期与第一次允许检查的超时、Proxy 的生产者心跳、客户端会话状态共同影响实际回查时刻,不能承诺每条消息刚好五秒后回查。Broker 事务参数源码
回查存在最大次数与保留限制。业务事务记录的保存期应覆盖最长回查、重复投递和人工恢复窗口;提前删除 COMMITTED 记录后,客户端只能回答未知,Broker 无法替应用重建事实。
生产者停机时,需要由持有对应事务订阅、能访问同一份持久状态的实例接管查询。只把结果存在 JVM Map 或只在原进程保留请求参数,进程重启后就失去了回答能力。
生产运行与失败处理
持久化、复制和队列数分别配置
同步刷盘影响本地确认成本;同步或异步复制影响节点失效时的恢复条件;副本选主与元数据角色又影响接管流程。单 Broker 实验只能支撑当前节点上的收发、确认与状态恢复判断,不能用它的 API 成功推导跨机器容灾能力。
消息大小、发送批次、队列数量、磁盘延迟、消费并发和下游容量会一起影响吞吐。先确定哪一个队列最慢、最老未完成消息等待多久,再判断应增加队列、增加消费者还是优化业务处理。库存服务已经被数据库锁等待限制时,增加消息消费线程可能只会增加连接争用。
对照现象查到下一步
| 现象 | 优先检查 | 下一步 |
|---|---|---|
| SDK 无法建立连接 | 使用 gRPC 还是 Remoting;endpoint 对应哪一角色 | 从客户端所在网络验证 Proxy 8081,查看 Proxy 日志 |
| Topic 不存在或消息类型被拒绝 | Topic 名称与 NORMAL/FIFO/TRANSACTION 属性 | 显式创建正确资源,检查权限,不开启自动建资源掩盖拼写错误 |
| 正常发送可用,回查持续失败 | 回查 brokerAddr、路由地址、生产者事务订阅 | 对照 Proxy 地址映射和 producer 心跳 |
| FIFO 后续记录提前收到 | 消费组 consumeMessageOrderly、批量与线程处理 | 确认 -o true,按同 group 串行完成 |
| 不可见时间结束后重复处理 | 处理时长、ACK 失败和当前 receipt handle | 延长合理窗口或分拆任务,保持业务幂等 |
| 事务半消息长期未决定 | 本地状态、checker 查询失败、在线生产者与回查次数 | 查询真实业务结果,恢复可查询状态,不盲目重发 |
| 初次启动卷目录拒绝访问 | Broker UID/GID 与指定卷所有权 | 修复该卷权限,不将整个宿主目录开放写入 |
| 初次收发正常,重复执行后 ACK 50001 / 超时 | Broker pop.log 中直接内存分配错误、分配池与容器预算 | 调整具体内存与池配置后复测,不把所有 50001 都归为订阅错误 |
结束时 docker compose -f compose.rocket.yaml stop 保留消息与事务记录。确认 me12-rocket 全部为可丢弃实验数据后,可执行:
docker compose -f compose.rocket.yaml down --volumes它删除本项目的容器、网络、Broker 卷与 PostgreSQL 卷,不删除共享源码和 Maven 缓存。真实业务系统不能用删卷清理“事务未完成”的现象,应先确定那些事务最终发生了什么。
