消息与事件:业务含义、提交与消费确认
消费者把订单写入数据库,再向消息系统确认处理完成。这两步通过不同连接调用,修改的也是不同数据:数据库保存订单结果,消息系统保存投递或消费进度。程序可能在两步之间退出,网络也可能只送达其中一方的回复。
理解异步处理,可以从三个具体对象开始:消息描述什么,消息系统如何保存和交付它,接收方怎样把处理结果提交到自己的存储。
消息载体与业务含义
命令、事件、通知和任务
消息是传输的数据单元。它的正文可以承载执行请求、业务事实、变更提示或待运行的工作。接收方怎样处理重复、超时和拒绝,取决于这些业务含义。
| 类型 | 示例 | 接收方承担的工作 | 需要说明的条件 |
|---|---|---|---|
| 命令 | ReserveStock:为订单预留 2 件商品 | 执行动作,成功或拒绝 | 目标能力、业务幂等键、数量单位、截止条件、结果获取方式 |
| 事件 | StockReserved:订单的库存已预留 | 根据既有事实更新自身状态 | 事实来源、事件身份、实体版本、数据含义 |
| 通知 | OrderChanged:订单状态有变化 | 按通知内容处理,或重新查询订单 | 是否携带完整状态、如何补漏、查询权限 |
| 任务 | GenerateInvoice:生成发票文件 | 调度和执行一项工作 | 任务 ID、输入、尝试次数、执行期限、取消与终止状态 |
命令描述期望发生的动作,执行方可以因库存不足而拒绝。事件描述发生过的事实:消费者无法通过拒绝接收 StockReserved 撤销库存预留;业务要释放库存,应再执行释放动作并产生相应结果。一个动作也可能产生多个不同事件,例如支付入账和风控状态变化,事件类型应让接收方知道各自描述哪项事实。
通知的可靠性由业务要求决定。页面刷新提示可以通过下次查询补上,必须送达的账户安全通知则需要持久存储、投递记录和重试。是否允许遗漏,不能仅从 Notification 这个类名判断。任务还多了调度状态:已创建、可执行、运行中、完成或取消;Broker 保存的消息只是安排执行的一种方式,长任务通常仍需独立的任务记录。
事件事实、上下文和消息载体的区分也见于 CloudEvents 1.0.2 规范。业务可以使用自己的消息格式;采用标准信封时,则要遵守该格式的字段与协议绑定规则。
竞争消费、发布订阅和可重放日志
同一条消息由谁接收,首先取决于存储和订阅关系。
竞争消费
一个工作队列
├─ 消费者 A 处理其中一部分消息
└─ 消费者 B 处理另一部分消息
发布订阅
一次发布
├─ 库存订阅:有自己的待处理记录和进度
└─ 审计订阅:有自己的待处理记录和进度
可重放日志
一份保留的分区日志
├─ 消费组 A 保存自己的读取位置
└─ 消费组 B 保存自己的读取位置在工作队列里加两个消费者,是分担同一批工作。库存和审计都必须处理每个订单时,需要独立订阅;让它们竞争同一个队列,会导致一部分订单只进入库存,另一部分只进入审计。
RabbitMQ 的常见队列模型会在消费确认后移除相应消息,重新读取已确认历史通常要依赖另存的记录或其他存储类型。优先级、重投和并发消费者会改变接收方观察到的顺序,不能仅从“队列”推断业务完成顺序。RabbitMQ 队列文档
Kafka 把记录保留在分区日志中,各消费组维护独立位置;提交消费位点后,记录仍可在保留条件允许时被重新读取。分区同时提供局部顺序和组内分工,数据清理由保留或压缩策略管理。Kafka 日志与消费设计
这些是通信模型,不是三个互斥的产品分类。一个产品可以提供多种存储或订阅能力。具体接入分别见 RabbitMQ 运行链、Kafka 运行链和 RocketMQ 运行链。
异步通信改变了等待方式
同步 HTTP 调用让调用方等待一次请求的返回,便于立即展示结果;下游不可用时,请求延迟和错误也会直接传回。异步消息把处理时间拆开:生产者完成发布后可以释放当前请求,下游按自身速度消费。突发流量先进入存储,再逐步处理,生产者与消费者也能独立部署。
代价是应用需要表达中间状态。订单接口返回“已受理”后,库存可能仍在排队,前端需要查询或订阅最终结果。消息积压超出保留时间、消费者长期失败、旧版本事件无法解析,都可能让业务停在中间状态;这些情况需要可查询的业务状态和恢复动作。
选择通信方式时,可先明确两个问题:调用方是否必须立刻得到执行结果,以及业务是否允许把完成时间延后。一个只需快速查询且依赖链很短的接口,用同步调用往往更清楚。要削峰、独立订阅、异步任务或重放事件时,消息存储才提供了直接收益。用消息模拟请求—响应也可实现,但 correlation ID、回复地址、超时和孤立回复都要由应用处理,工作量不会因换了协议而消失。
一条消息经过哪些独立提交
发布、交付和处理各修改什么
以“订单创建后,库存服务保存一条处理结果”为例,可以观察到以下位置:
| 位置 | 发起操作的组件 | 实际改变的数据 | 应观察什么 |
|---|---|---|---|
| 订单事务提交 | 订单服务 | 订单数据库中的业务记录 | 事务结果、可查询的订单状态 |
| 消息发布 | Producer 客户端 | 客户端缓冲、网络输出及 Broker 中的记录 | 发送异常、发布确认、路由结果 |
| 消息交付 | Broker / Consumer 客户端 | 消费者获得一次投递或一批记录 | 消息身份、分区位置或 delivery tag |
| 库存事务提交 | 库存服务 | 库存数据库中的处理结果 | 数据库事务和实际业务行 |
| 消费确认或位点提交 | Consumer 客户端 | 队列投递状态或消费组位置 | Broker 的确认处理、已提交位置 |
表中的“发布”还可能包含多个阶段。客户端 send() 返回,有时只是把记录放进本地缓冲;等待的 Future 或 confirm 回调才对应 Broker 回复。不同 SDK 的同步、异步和单向发送有不同返回含义,要以具体 API 为准。
RabbitMQ 的 publisher confirm 和 consumer acknowledgement 分属发布侧与消费侧。前者反馈 Broker 对发布的处理,后者告诉 Broker 某次投递可以完成。消费者数据库提交发生在应用自己的 JDBC 连接上,不在这两种 AMQP 确认协议之内。RabbitMQ 发布确认与消费确认
下图采用“先提交消费者数据库,再消费确认”的普通队列处理方式。生产者与自身数据库的协调另有双写问题,图中不假设它们已自动组成分布式事务。
数据库事务可以把多条 SQL 合成一次提交。ROLLBACK 撤销当前事务中的数据库修改,已提交结果可由其他连接查询;AMQP ACK 不属于这组 SQL,数据库回滚也不会把它撤回。PostgreSQL 事务
超时后的结果可能需要查询
一次发送可能已到达 Broker,但确认回复在网络中丢失。此时客户端知道自己没有收到确认,却不知道 Broker 是否保存了消息。直接重发可以恢复未送达的情况,也可能造成同一业务事件再次发布。
数据库 commit() 遇到连接中断也有类似问题:连接失败不能单独说明事务一定回滚。重试前查询稳定的业务键,或者让重复请求进入幂等处理,才能把未知结果转成可判断的状态。为重试生成全新的事件 ID,会让消费者失去识别同一次业务事实的依据。
需要区分三类结果:明确成功、明确拒绝,以及因超时或连接断开而未知。拒绝可能要求修复权限或消息格式;未知通常要求查询、重试和去重共同处理。RabbitMQ 可靠性说明
两个消费失败窗口
| 操作顺序与故障位置 | 数据库 | 队列 | 恢复方向 |
|---|---|---|---|
| 先 ACK,再写数据库;写入前或提交前失败 | 无业务结果 | 消息已确认移除 | 按业务记录补发、查询或对账,不能等这次投递自动回来 |
| 数据库回滚,尚未 ACK;消费连接关闭 | 无业务结果 | 未确认消息可重新投递 | 再次执行事务,成功后确认 |
| 数据库已提交,ACK 前消费连接关闭 | 已有业务结果 | 消息可能重新投递 | 根据稳定事件 ID 去重,避免再次产生业务效果 |
后两行通常更容易恢复,因此业务消费者常在数据库提交后确认。第三行也说明了为什么需要 投递、幂等与顺序处理:调整两次调用的顺序解决不了重复业务效果,必须让数据库更新能识别已经处理的事件。
生产端也有一对独立操作:保存订单和发布消息。订单已经提交而发布失败时,消费者甚至还没有消息可重投。把待发布事件与订单一起写入数据库,再由发布器恢复发送,是 Outbox 与 CDC处理的问题。
运行一个提交后确认的消费者
Linux、容器与实验账号
下载完整消息实验源码 ZIP,解压后进入 message-event-lab。宿主需要 Docker Engine、Compose v2+ 和解压工具,无需安装 Java。下面所有 shell 命令都在这个解压目录执行,同一实验项目一次只运行一个模式。
| 组件 | 版本 | 用途 |
|---|---|---|
| RabbitMQ | 4.3.5 management 镜像 | 保存真实 quorum queue 和消费确认状态 |
| PostgreSQL | 18.6 | 保存消费者的业务结果 |
| Java / Maven | Temurin 25 / Maven 3.9.12 | 在构建容器编译,生成 Java 17 字节码 |
| RabbitMQ Java 客户端 | 5.33.0 | AMQP 0-9-1 连接、发布确认与消费确认 |
| pgJDBC | 42.7.13 | JDBC 事务与 SQL 查询 |
RabbitMQ 服务端版本入口见 官方下载页,Java 客户端及运行要求见 官方客户端说明。实验固定这些版本,切换 Broker 或客户端版本后应重新运行三种模式。
RabbitMQ 实验使用的目录如下,完整 ZIP 还包含其他消息实验入口:
message-event-lab/
├── compose.yaml
├── config/
│ ├── rabbitmq.conf
│ ├── rabbit-definitions.json
│ └── init.sql
├── pom.xml
└── rabbit/
├── pom.xml
└── src/main/java/example/MessageModelLab.java宿主操作者须有运行 Docker 的权限。Maven 和 Java 容器用宿主 UID/GID,源码目录及缓存由该用户写入;这不会改变 Docker daemon 或服务镜像初始化时的身份。PostgreSQL 的 bootstrap 用于初始化,Java 只使用普通数据库角色 app。RabbitMQ 的 app 没有管理标签,仅能访问实验 vhost lab。
Compose 不向宿主发布端口。Java 容器加入专用网络后,通过 rabbit:5672 与 postgres:5432 访问服务;这里的服务名由容器网络解析,不能在宿主上直接当作 DNS 主机名使用。
compose.yaml 的两个服务如下:
name: me12-lab
services:
rabbit:
image: rabbitmq:4.3.5-management
hostname: rabbit
mem_limit: 512m
volumes:
- rabbit-data:/var/lib/rabbitmq
- ./config/rabbitmq.conf:/etc/rabbitmq/rabbitmq.conf:ro
- ./config/rabbit-definitions.json:/etc/rabbitmq/definitions.json:ro
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "check_running"]
interval: 3s
timeout: 5s
retries: 30
postgres:
image: postgres:18.6
mem_limit: 256m
environment:
POSTGRES_DB: lab
POSTGRES_USER: bootstrap
POSTGRES_PASSWORD: bootstrap-lab-only
volumes:
- postgres-data:/var/lib/postgresql
- ./config/init.sql:/docker-entrypoint-initdb.d/01-lab.sql:ro
healthcheck:
test: ["CMD-SHELL", "pg_isready -h 127.0.0.1 -U bootstrap -d lab"]
interval: 3s
timeout: 5s
retries: 30
volumes:
rabbit-data:
postgres-data:PostgreSQL 18 镜像的命名卷挂载在 /var/lib/postgresql。配置文件使用只读挂载,业务数据放命名卷。实验凭据只用于这两个隔离容器,不用于已有实例或生产环境。
config/rabbitmq.conf 导入本地定义:
definitions.import_backend = local_filesystem
definitions.local.path = /etc/rabbitmq/definitions.jsonconfig/rabbit-definitions.json 创建普通应用用户和 vhost:
{
"users": [{"name": "app", "password": "app-lab-only", "tags": []}],
"vhosts": [{"name": "lab"}],
"permissions": [{"user": "app", "vhost": "lab", "configure": ".*", "write": ".*", "read": ".*"}]
}应用可管理 lab 中的实验队列。生产接入应进一步限制队列、exchange 的名称表达式,并分离拓扑管理与普通收发身份。config/init.sql 为应用建立自己的数据库 schema:
CREATE ROLE app LOGIN PASSWORD 'app-lab-only';
CREATE SCHEMA lab AUTHORIZATION app;
ALTER DATABASE lab SET search_path TO lab, public;启动并等待两个健康检查通过:
docker compose up -d --wait --wait-timeout 120 rabbit postgres
docker compose ps两项服务应为 healthy。如果等待超时,先执行 docker compose logs --tail=80 rabbit postgres。认证失败与服务启动失败要分开判断:健康检查通过并不验证 Java 使用的账号。可直接检查账号和数据库 schema:
docker compose exec rabbit rabbitmqctl list_users
docker compose exec -e PGPASSWORD=app-lab-only postgres \
psql -h 127.0.0.1 -U app -d lab \
-c 'SELECT current_user, current_schema();'应看到 RabbitMQ 用户 app 的标签为 [],数据库结果为 app | lab。已有数据卷不会因修改 init.sql 自动重新初始化;出现角色或密码不符时,先确认是否复用了旧实验数据,别把删卷当成通用排障步骤。
依赖与完整入口
根 pom.xml 定义 Java 17 编译目标。下面是只保留 rabbit 模块的完整最小配置;共享 ZIP 中的根 POM 还会声明其他产品模块,运行 -pl rabbit -am 只选择当前模块及其构建依赖:
<project xmlns="http://maven.apache.org/POM/4.0.0">
<modelVersion>4.0.0</modelVersion>
<groupId>example</groupId><artifactId>message-event-lab</artifactId><version>1.0.0</version>
<packaging>pom</packaging>
<modules><module>rabbit</module></modules>
<properties>
<maven.compiler.release>17</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<build><plugins>
<plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-compiler-plugin</artifactId><version>3.14.1</version></plugin>
<plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-dependency-plugin</artifactId><version>3.8.1</version></plugin>
</plugins></build>
</project>rabbit/pom.xml 引入 AMQP 和 JDBC 客户端。日志 API 与实现显式使用同一版本,避免依赖传递引入的日志绑定警告干扰输出:
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent><groupId>example</groupId><artifactId>message-event-lab</artifactId><version>1.0.0</version></parent>
<artifactId>rabbit-lab</artifactId>
<dependencies>
<dependency><groupId>com.rabbitmq</groupId><artifactId>amqp-client</artifactId><version>5.33.0</version></dependency>
<dependency><groupId>org.postgresql</groupId><artifactId>postgresql</artifactId><version>42.7.13</version></dependency>
<dependency><groupId>org.slf4j</groupId><artifactId>slf4j-api</artifactId><version>2.0.17</version></dependency>
<dependency><groupId>org.slf4j</groupId><artifactId>slf4j-simple</artifactId><version>2.0.17</version></dependency>
</dependencies>
</project>MessageModelLab.java 有三个模式。先运行 normal:发布到独立队列,等待 confirm,取出一条未自动确认的消息,提交数据库,再 ACK。early-ack 与 unacked 只改变失败位置。
程序会重置对应 me12.model.<mode> 实验队列及 event-<mode> 数据行,确保单次结果可比较。它只适合这个隔离环境,不应连接真实订单库或已有业务队列。
package example;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.GetResponse;
import java.nio.charset.StandardCharsets;
import java.sql.DriverManager;
import java.util.Map;
public final class MessageModelLab {
private static final String JDBC = "jdbc:postgresql://postgres:5432/lab";
private static final String USER = "app";
private static final String PASSWORD = "app-lab-only";
private static void check(boolean ok, String message) {
if (!ok) throw new IllegalStateException(message);
}
private static java.sql.Connection database() throws java.sql.SQLException {
return DriverManager.getConnection(JDBC, USER, PASSWORD);
}
private static void prepareRow(String id) throws java.sql.SQLException {
try (var c = database(); var s = c.createStatement()) {
s.execute("CREATE TABLE IF NOT EXISTS model_receipt(event_id text PRIMARY KEY, order_id text NOT NULL)");
try (var p = c.prepareStatement("DELETE FROM model_receipt WHERE event_id=?")) {
p.setString(1, id); p.executeUpdate();
}
}
}
private static void writeRow(java.sql.Connection c, String id) throws java.sql.SQLException {
try (var p = c.prepareStatement("INSERT INTO model_receipt(event_id,order_id) VALUES (?,?)")) {
p.setString(1, id); p.setString(2, "order-42"); p.executeUpdate();
}
}
private static int rows(String id) throws java.sql.SQLException {
try (var c = database(); var p = c.prepareStatement("SELECT count(*) FROM model_receipt WHERE event_id=?")) {
p.setString(1, id);
try (var r = p.executeQuery()) { r.next(); return r.getInt(1); }
}
}
private static GetResponse receive(Channel ch, String queue, String id) throws Exception {
long deadline = System.nanoTime() + 5_000_000_000L;
do {
GetResponse message = ch.basicGet(queue, false);
if (message != null) {
check(id.equals(message.getProps().getMessageId()), "unexpected event ID");
check("order-42".equals(new String(message.getBody(), StandardCharsets.UTF_8)), "unexpected body");
return message;
}
Thread.sleep(20);
} while (System.nanoTime() < deadline);
throw new IllegalStateException("no message within 5 seconds: " + queue);
}
private static void ackAndObserve(Channel ch, String queue, GetResponse message) throws Exception {
ch.basicAck(message.getEnvelope().getDeliveryTag(), false);
// queue.declare-ok is a synchronous reply after the ACK on this channel.
check(ch.queueDeclarePassive(queue).getMessageCount() == 0, "unexpected ready message");
}
public static void main(String[] args) throws Exception {
String mode = args.length == 1 ? args[0] : "";
check(mode.equals("normal") || mode.equals("early-ack") || mode.equals("unacked"),
"usage: MessageModelLab normal|early-ack|unacked");
String queue = "me12.model." + mode;
String id = "event-" + mode;
prepareRow(id);
var factory = new ConnectionFactory();
factory.setHost("rabbit"); factory.setVirtualHost("lab");
factory.setUsername(USER); factory.setPassword(PASSWORD);
factory.setConnectionTimeout(5000); factory.setAutomaticRecoveryEnabled(false);
try (var connection = factory.newConnection(); var publisher = connection.createChannel()) {
publisher.queueDeclare(queue, true, false, false, Map.of("x-queue-type", "quorum"));
publisher.queuePurge(queue); // Only this mode's isolated lab queue.
publisher.confirmSelect();
publisher.basicPublish("", queue, true, new AMQP.BasicProperties.Builder()
.messageId(id).deliveryMode(2).contentType("text/plain").build(),
"order-42".getBytes(StandardCharsets.UTF_8));
publisher.waitForConfirmsOrDie(5000);
System.out.println("published: event=" + id + " confirmed=true");
try (var consumer = connection.createChannel(); var c = database()) {
var message = receive(consumer, queue, id);
c.setAutoCommit(false);
if (mode.equals("early-ack")) {
ackAndObserve(consumer, queue, message);
writeRow(c, id); c.rollback();
} else if (mode.equals("unacked")) {
writeRow(c, id); c.rollback();
// Channel close returns this unacknowledged delivery to the queue.
} else {
try { writeRow(c, id); c.commit(); }
catch (Exception | Error failure) {
try { c.rollback(); }
catch (java.sql.SQLException rollbackFailure) { failure.addSuppressed(rollbackFailure); }
throw failure;
}
ackAndObserve(consumer, queue, message);
}
}
try (var observer = connection.createChannel()) {
if (mode.equals("unacked")) {
check(rows(id) == 0, "rolled-back row remained");
var message = receive(observer, queue, id);
check(message.getEnvelope().isRedeliver(), "expected broker redelivery");
System.out.println("before-recovery: rows=0 redelivered=true sameEventId=true");
try (var c = database()) {
c.setAutoCommit(false);
try { writeRow(c, id); c.commit(); }
catch (Exception | Error failure) {
try { c.rollback(); }
catch (java.sql.SQLException rollbackFailure) { failure.addSuppressed(rollbackFailure); }
throw failure;
}
}
ackAndObserve(observer, queue, message);
}
check(observer.basicGet(queue, false) == null, "unexpected additional message");
int count = rows(id);
check(count == (mode.equals("early-ack") ? 0 : 1), "unexpected committed row count");
System.out.println("result: mode=" + mode + " rows=" + count + " ready=0");
}
}
}
}basicGet(queue, false) 每次主动取一条,便于逐步观察消息;生产中的持续消费通常使用 basicConsume 回调。这里没有利用 basicGet 验证 prefetch,预取上限应在真实订阅消费者上观察。连接、channel、同步取消息和资源关闭的 API 见 RabbitMQ Java API Guide。
queueDeclare 建立 durable quorum queue,deliveryMode(2) 表达持久消息属性。Quorum queue 的复制与确认由其存储实现完成;当前环境只有一个 RabbitMQ 节点,能够观察进程内外的交付与提交,但没有多节点容灾能力。Quorum queue 文档
构建并得到第一次成功
源码目录与缓存使用当前 Linux 用户的 UID/GID。先建立可写缓存,再编译模块并复制运行依赖:
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户运行实验'; exit 1; }
mkdir -p .m2
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预期 Maven 返回 BUILD SUCCESS,并生成 rabbit/target/classes 和 rabbit/target/dependency。下载失败先检查 Maven 输出中的仓库地址与网络错误;若提示本地 repository 不可写,检查 .m2 权限以及 -Dmaven.repo.local 是否完整传入,不能只修改 Java 运行容器身份。
启动 Java 容器,源码和产物只读挂载:
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.MessageModelLab normal结果为:
published: event=event-normal confirmed=true
result: mode=normal rows=1 ready=0第一行来自真正完成的 confirm 等待。第二行同时依赖新的数据库连接查询和 Broker 取消息结果:业务表有一行,队列没有另一条待取消息。Java 中的 check 是普通运行时判断,不依赖 -ea;任何结果不符都会抛出异常并非零退出。
再查看队列与数据库的当前状态:
docker compose exec rabbit rabbitmqctl list_queues -p lab \
name messages_ready messages_unacknowledged
docker compose exec -e PGPASSWORD=app-lab-only postgres \
psql -h 127.0.0.1 -U app -d lab -c 'TABLE lab.model_receipt;'me12.model.normal 的 ready 和 unacknowledged 均为 0,业务表含 event-normal | order-42。ready=0 单独看仍有歧义:消息也可能已交付但尚未确认,因此排查时要连同 unacknowledged 查看。
先确认,再回滚数据库
把运行命令最后一个参数换成 early-ack:
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.MessageModelLab early-ack程序先发送 ACK,并在同一 channel 执行同步查询,再写入数据库后主动回滚。这样产生的是可控的“确认已先行、业务提交失败”窗口:
published: event=event-early-ack confirmed=true
result: mode=early-ack rows=0 ready=0数据库中没有 event-early-ack,队列也取不到该事件。回滚影响 SQL,无法撤销之前的消息确认。关闭并重新打开消费者不会产生新的这次投递,应用需要从订单记录、发布记录或对账结果确定是否补发。
这个模式正常退出表示负例产生了预期状态,绝不表示业务处理成功。监控若只统计发送成功或队列已经排空,就会漏掉这里的业务缺失。
未确认的消息怎样回来
unacked 在收到消息后写库、回滚,然后关闭消费 channel,没有发送 ACK:
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.MessageModelLab unacked另一个 channel 从 Broker 重新取到该事件,检查 redelivered 与稳定 ID,再执行事务并确认:
published: event=event-unacked confirmed=true
before-recovery: rows=0 redelivered=true sameEventId=true
result: mode=unacked rows=1 ready=0恢复前的数据库没有结果,重新投递后才建立一行。这里关闭的是 channel;真实连接丢失也可能触发未确认消息重投,区别在于 Broker 何时确认连接失效。网络故障不一定立刻被发现,因此应用等待必须有超时,并观察连接与心跳状态。
三个模式运行后,数据库应有 event-normal 和 event-unacked 两行,缺少 event-early-ack。各队列的 ready / unacknowledged 都为 0。它们共同说明了确认位置如何改变可恢复数据,尚未覆盖“数据库已提交但 ACK 丢失”的重复写入问题;接入实际业务时必须把该窗口纳入幂等设计。
从错误输出继续定位
| 第一个可观察结果 | 优先核对 | 操作与后续判断 |
|---|---|---|
UnknownHostException: rabbit | Java 容器是否加入正确网络 | docker network inspect me12-lab_default;确认 Java 命令使用该网络,Compose 项目名改变时同步修改 |
| 连接被拒绝或超时 | 服务是否健康、容器是否退出 | docker compose ps 与 docker compose logs --tail=80 rabbit postgres;先恢复服务再运行消息模式 |
AMQP ACCESS_REFUSED | 用户、密码、vhost 和权限 | docker compose exec rabbit rabbitmqctl list_permissions -p lab;应含 app 的实验权限 |
| JDBC 认证失败或 schema 不存在 | 连接身份与已有卷初始化 | 重跑前面的 SELECT current_user, current_schema();检查 init.sql 与已有角色,不直接删业务数据 |
PRECONDITION_FAILED | 同名队列原有属性 | docker compose exec rabbit rabbitmqctl list_queues -p lab name type durable;声明需与现有类型一致 |
| 等待消息超过 5 秒 | 发布/路由与其他消费者 | 先看 confirm 是否成功,再看对应队列的 ready / unacknowledged;同一实验不能并发运行同一模式 |
| rows 数量不符 | SQL 异常、重复执行或实验状态混用 | 查询 lab.model_receipt 与 eventId,保留异常;不要继续发 ACK 掩盖失败 |
停止服务并保留数据:
docker compose stop再次启动用 docker compose up -d --wait rabbit postgres。仅当确认本项目全为可丢弃实验数据时,删除本组容器、网络和命名卷:
docker compose down --volumes这会删除实验队列和业务表,源码与宿主 .m2 缓存仍保留。不要把该命令用于已有业务 Compose 项目。
为消息定义可长期使用的契约
把事件身份和业务实体分开
同一订单会产生创建、支付、取消等多个事件。orderId 标识业务实体,eventId 标识某一次事实;用 orderId 去重所有消息,会把该订单后续的正常变化一起丢掉。一次事件重试投递时继续使用同一 eventId,而消费投递的 delivery tag 可能变化。
一个自定义订单事件可以拆成下面几部分:
{
"eventId": "evt-order-paid-42-v3",
"eventType": "OrderPaid",
"schemaVersion": 1,
"source": "payment-service",
"aggregateId": "order-42",
"aggregateVersion": 3,
"payload": {
"amountMinor": 12900,
"currency": "CNY"
}
}事件信封
├── 身份:eventId + source 的唯一性约定
├── 含义:eventType、schemaVersion
├── 业务实体:aggregateId、aggregateVersion
├── 可选上下文:发生时间、correlationId、causationId
└── payload:事件发生时需要传播的业务数据
传输元数据
├── RabbitMQ:exchange、routing key、delivery tag
├── Kafka:topic、partition、offset
└── 客户端观察:重投标志、尝试次数、接收时间示例刻意把金额写成 amountMinor 并给出币种,避免消费者把整数当元或把小数精度当作默认规则。schemaVersion 表示消息格式/契约版本,aggregateVersion 表示订单业务状态的演进次序;两者可以分别变化。
correlation ID 用于关联同一业务过程中的多次交互,causation ID 指向直接触发当前事件的原因。需要跨服务排查时可加入,简单单向消息不必为凑字段而强行生成。事件发生时间、Broker 写入时间和消费者观察时间也应分别命名;排队延迟计算需要考虑时间来源与时钟误差。
若选择 CloudEvents,上面的自定义字段不能原样当成标准。CloudEvents 1.0.2 的必需上下文字段为 id、source、specversion、type,数据放在 data 或格式规定的位置;time 等字段有自己的可选规则。使用标准的好处在于跨组件能识别一致的事件上下文,业务金额、实体版本和拒绝策略仍需单独约定。
序列化和变更需要消费者一起验证
常用消息体包括 JSON、Avro、Protocol Buffers 和产品支持的其他二进制格式。JSON 便于直接查看,二进制格式通常依赖 schema 或生成代码。Broker 能接受字节,不会替库存服务验证“数量必须为正数”或“币种应与账户一致”;消费者应在业务写入前完成大小、类型、字段和业务条件检查。
| 修改 | 旧消费者可能遇到什么 | 更稳妥的做法 |
|---|---|---|
| 新增可选字段 | 宽容解析器可忽略,严格解析器可能拒绝 | 用现有消费者代码读取新样本,验证忽略或默认行为 |
| 字段改名或删除 | 读取不到原字段 | 保持过渡期兼容,或使用新事件类型/明确版本 |
| 金额从分改成元 | JSON 仍可解析,但金额扩大或缩小 | 使用不同字段与单位约定,禁止静默改变含义 |
| 新增枚举值 | 反序列化或业务分支失败 | 定义未知值策略,升级消费者后再扩大生产值集合 |
| 同一事件类型改为不同业务事实 | 去重、重试与状态更新方式可能失效 | 重新定义事件类型和消费契约 |
事件携带多少数据也影响可靠性。只有实体 ID 的通知很小,但每个消费者都要有权限回查,而且查询拿到的可能是实体最新状态。携带事实快照可以减少回查,并保留事件发生时的必要信息,却会增加隐私扩散、消息大小和格式演进成本。需要精确恢复历史事实时,不能只发一个 ID,再默认未来查询必然还原当时的数据。
权限、保留和运行记录
生产者需要向指定主题或 exchange 写入的权限;消费者需要订阅或读取,以及对应业务存储的权限。管理账号通常还能修改拓扑、删除队列或重置位置,应与应用收发账号分离。跨不可信网络使用 TLS,并校验服务端身份;凭据由运行环境注入,不写进消息正文和公开日志。
序列化选择也涉及安全。消费者应限制单条消息大小、允许的事件类型和嵌套复杂度,不执行消息携带的类名、脚本或任意对象反序列化指令。解析失败应隔离并保留可定位的错误,不能反复立即重投占满整个消费组。延迟、重试与死信处理
消息数据常会复制到日志、重试队列、死信队列和归档。保留期应同时考虑业务重放需求、个人信息删除和消费故障最长恢复时间。去重记录若先被清理,稍后重放的旧事件可能再次产生业务效果;恢复窗口与去重窗口需要共同设计。
运行记录至少要能回答:哪个事件发布到哪里、消费者处理了哪个实体、事务是否已提交、确认或位点推进到了哪里。日志可以带 eventId 和错误类型,指标则按 topic、queue、consumer group 等有限维度聚合,不把每个 eventId 变成一个时间序列标签。
排队数量、最老未完成时间与业务成功数承担不同用途。队列排空可能来自正常处理,也可能来自提前确认、过期或错误路由;要判断业务是否完成,仍要查看目标状态。出现持续积压时,进一步比较入口与处理速率、分区分配和下游耗时,见 积压、扩容与可观测性。
权威资料与规范地址
消息模型与事件格式
| 资料 | 查阅内容与完整地址 |
|---|---|
| CloudEvents 1.0.2 | 事件上下文、必需字段及消息格式:https://raw.githubusercontent.com/cloudevents/spec/v1.0.2/cloudevents/spec.md |
| Kafka 4.3 设计 | 分区日志、消费位置与处理语义:https://kafka.apache.org/43/design/design/ |
RabbitMQ 的存储与确认
| 资料 | 查阅内容与完整地址 |
|---|---|
| Queues | 队列属性、顺序和消息生命周期:https://www.rabbitmq.com/docs/queues |
| Acknowledgements and Confirms | 发布与消费确认:https://www.rabbitmq.com/docs/confirms |
| Reliability Guide | 连接故障、确认缺失与重投:https://www.rabbitmq.com/docs/reliability |
| Quorum Queues | 复制存储与确认条件:https://www.rabbitmq.com/docs/quorum-queues |
| Java Client API Guide | 连接、channel、取消息与资源管理:https://www.rabbitmq.com/client-libraries/java-api-guide |
| Java Client | 客户端版本和 Java 要求:https://www.rabbitmq.com/client-libraries/java-client |
| Download | 服务端版本与分发入口:https://www.rabbitmq.com/docs/download |
数据库事务
PostgreSQL 18 Transactions,提交、回滚与事务可见性:https://www.postgresql.org/docs/18/tutorial-transactions.html。
