消息投递、幂等与顺序:让重复请求只改变一次业务
余额增加 10 元的消息被收到两次,正确结果仍应只增加 10 元。把两次交付识别为同一业务事件只是第一步;去重记录和余额更新还必须一起提交,否则去重本身也可能挡住一次尚未完成的业务处理。
先决定重复交付应产生什么结果
三种投递语义描述哪一段行为
at-most-once 允许丢失但避免恢复时重复投递,例如先确认消息再执行操作。at-least-once 在未确认时允许重新交付,应用需要承受重复。exactly-once 必须说明约束覆盖哪些记录、状态和参与者:Kafka 事务可以协调 Kafka 输出与输入位点,单个数据库事务可以协调表内更新,外部支付平台则有自己的请求与确认规则。Kafka 交付语义
实际业务通常更关心“同一个业务动作生效几次”。消费者收到两次事件、两次都查询了数据库,但账户只改变一次,是可接受的幂等处理。事件送达次数、代码执行次数与持久业务变更次数应分开记录。
ACK 丢失也会导致重复。例如消费者已经提交余额事务,刚发送确认就断开连接;Broker 没有记录到这次确认时会再次交付。RabbitMQ 的 redelivered 标志可帮助诊断,业务去重仍应使用稳定事件 ID,因为上游重新发布的相同业务事件可能作为全新的消息进入队列。RabbitMQ 消费确认
事件 ID、业务键与消费方名称各自做什么
eventId = payment-recorded-42-1 一次已发生的业务事件
aggregateId = order-42 事件所属订单
aggregateVersion = 1 此订单的事件顺序
consumer = balance-v1 当前处理职责
deliveryTag / offset / messageId 传输系统中的交付或存储标识事件 ID 在重发和重试时保持不变。库存和审计可以独立处理同一事件,因此 Inbox 的唯一键经常是 (consumer, event_id)。只按 event ID 建全局唯一约束,会让不同处理职责错误地互相去重;仅按订单号去重,又会把同一订单后续的合法支付、退款等事件全部挡住。
业务请求 ID 与事件 ID 也可能不同。一次下单请求可以产生订单创建、库存预留等多个事件;重试下单请求应返回既有订单结果,而这些事件各自拥有独立 ID。对客户端提供幂等键时,要同时约束用户或租户、操作类型和请求内容:同一键搭配不同金额,应拒绝冲突,而不是沿用第一次结果。
CloudEvents 将 source 与 id 的组合用于事件唯一性。采用其他信封格式时也要明确 ID 的唯一范围、生成方和重发规则,不只在 JSON 中增加一个没有约束的字符串字段。CloudEvents 事件身份
选择能约束最终更新的存储
数据库唯一约束适合与本地业务更新放进一个事务。进程内 Set 在重启后丢失,跨实例也无法仲裁;Redis 短期标记可以拦截短时间内的重复流量,但它的过期、持久化与故障切换条件会影响保护范围。涉及账务等不可重复动作时,应把持久去重规则落实到实际业务写入所在的系统。
如果调用外部支付、短信或文件系统,本地 Inbox 事务无法原子提交那个外部副作用。优先使用对方支持的幂等请求键与结果查询接口;缺少这些能力时,需要持久操作状态、对账和人工处理入口。数据库已经标记完成,却无法确认外部是否执行成功时,不应重新换一个请求 ID 盲目重做。
把 Inbox 和业务更新放进同一个事务
两张表共同保存一次处理
实验用一张余额表和一张 Inbox。余额只是可精确断言的业务计数,省略真实货币所需的币种、小数精度与账本设计。
CREATE TABLE delivery_account (
id text PRIMARY KEY,
balance int NOT NULL,
version int NOT NULL
);
CREATE TABLE delivery_inbox (
consumer text NOT NULL,
event_id text NOT NULL,
PRIMARY KEY (consumer, event_id)
);正常事务先尝试登记事件,只有新事件才修改余额。唯一冲突使用 ON CONFLICT DO NOTHING 处理,并以受影响行数判断是否继续。两个消费者同时插入同一个唯一键时,PostgreSQL 会协调冲突行;失败方不能凭插入前的一次 SELECT 认定自己拥有处理资格。PostgreSQL INSERT 与 ON CONFLICT
BEGIN
INSERT inbox ... ON CONFLICT DO NOTHING
├── 插入 1 行 → UPDATE account ... WHERE version=expected
│ ├── 更新 1 行 → COMMIT → ACK
│ └── 更新 0 行 → ROLLBACK,Inbox 也回滚
└── 插入 0 行 → 已完成过,结束事务 → ACK这里可以把“查到 Inbox”解释为处理完成,因为 Inbox 和业务更新遵循同一提交约定。如果表中还允许单独写入 PROCESSING 状态,就必须读取状态并处理租约、过期和恢复,不能仍将“存在即成功”用于所有记录。PostgreSQL 事务原子性
准备 Linux、Broker 与数据库
下载完整实验工程,解压进入 message-event-lab。使用 Linux shell、Docker Engine 和 Compose v2+,由有 Docker 权限的普通宿主用户操作。Compose 启动 RabbitMQ 4.3.5 与 PostgreSQL 18.6,应用以 vhost lab 的用户 app 连接 RabbitMQ,以数据库普通角色 app 在 schema lab 建表。专用网络不发布宿主端口,账号口令仅用于可丢弃实验。
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户'; exit 1; }
mkdir -p .m2
docker compose up -d --wait --wait-timeout 120 rabbit postgres
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-dependenciesMaven 使用显式 UID/GID 和当前用户可写的缓存目录,预期输出 BUILD SUCCESS。共享根 POM、RabbitMQ 配置和数据库角色初始化见 消息与事件实验环境。rabbit 模块使用 AMQP 客户端 5.33.0、pgJDBC 42.7.13,编译目标 Java 17;下面直接运行 Java 25 镜像。
入口 rabbit/src/main/java/example/DeliveryLab.java 的完整代码如下。每个模式都会清理 balance-v1 的实验 Inbox 并将 order-42 的余额与版本重置为 0,只能用于这份独立实验数据,不能并行启动多个模式。
package example;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;
import java.sql.DriverManager;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
public final class DeliveryLab {
record Event(String id,int version,int amount) {}
enum Result { APPLIED,DUPLICATE }
static final class VersionConflict extends Exception {}
static void check(boolean ok,String why) {if(!ok)throw new IllegalStateException(why);}
static java.sql.Connection database() throws Exception {
return DriverManager.getConnection("jdbc:postgresql://postgres:5432/lab","app","app-lab-only");
}
static void prepare() throws Exception {
try(var c=database();var s=c.createStatement()) {
s.execute("CREATE TABLE IF NOT EXISTS delivery_account(id text PRIMARY KEY,balance int NOT NULL,version int NOT NULL)");
s.execute("CREATE TABLE IF NOT EXISTS delivery_inbox(consumer text NOT NULL,event_id text NOT NULL,PRIMARY KEY(consumer,event_id))");
s.execute("DELETE FROM delivery_inbox WHERE consumer='balance-v1'");
s.execute("INSERT INTO delivery_account VALUES('order-42',0,0) ON CONFLICT(id) DO UPDATE SET balance=0,version=0");
}
}
static boolean claim(java.sql.Connection c,String id) throws Exception {
try(var p=c.prepareStatement("INSERT INTO delivery_inbox VALUES('balance-v1',?) ON CONFLICT DO NOTHING")) {
p.setString(1,id);return p.executeUpdate()==1;
}
}
static Result apply(Event event) throws Exception {
try(var c=database()) {
c.setAutoCommit(false);
try {
if(!claim(c,event.id())) {c.commit();return Result.DUPLICATE;}
try(var p=c.prepareStatement("UPDATE delivery_account SET balance=balance+?,version=? WHERE id='order-42' AND version=?")) {
p.setInt(1,event.amount());p.setInt(2,event.version());p.setInt(3,event.version()-1);
if(p.executeUpdate()!=1)throw new VersionConflict();
}
c.commit();return Result.APPLIED;
} catch (Exception | Error failure) {
try { c.rollback(); }
catch (java.sql.SQLException rollbackFailure) { failure.addSuppressed(rollbackFailure); }
throw failure;
}
}
}
static void state(int balance,int version,int inbox) throws Exception {
try(var c=database();var s=c.createStatement()) {
try(var r=s.executeQuery("SELECT balance,version FROM delivery_account WHERE id='order-42'")) {
check(r.next()&&r.getInt(1)==balance&&r.getInt(2)==version,"account state differs");
}
try(var r=s.executeQuery("SELECT count(*) FROM delivery_inbox WHERE consumer='balance-v1'")) {
check(r.next()&&r.getInt(1)==inbox,"inbox count differs");
}
}
}
static com.rabbitmq.client.GetResponse receive(Channel ch,String queue) throws Exception {
long end=System.nanoTime()+TimeUnit.SECONDS.toNanos(5);
while(System.nanoTime()<end) {var m=ch.basicGet(queue,false);if(m!=null)return m;Thread.sleep(20);}
throw new IllegalStateException("delivery did not arrive");
}
public static void main(String[] args) throws Exception {
check(args.length==1,"mode: duplicate|split|concurrent|order");prepare();
if(args[0].equals("concurrent")) {
var start=new CountDownLatch(1);var pool=Executors.newFixedThreadPool(2);
try {
var one=pool.submit(()->{start.await();return apply(new Event("same-event",1,10));});
var two=pool.submit(()->{start.await();return apply(new Event("same-event",1,10));});
start.countDown();var results=List.of(one.get(10,TimeUnit.SECONDS),two.get(10,TimeUnit.SECONDS));
check(results.contains(Result.APPLIED)&&results.contains(Result.DUPLICATE),"concurrent claim did not serialize");
state(10,1,1);System.out.println("concurrent: applied=1 duplicate=1 balance=10 inbox=1");
} finally {pool.shutdownNow();}
return;
}
if(args[0].equals("order")) {
try {apply(new Event("event-2",2,20));throw new IllegalStateException("gap accepted");}
catch(VersionConflict expected) {state(0,0,0);}
check(apply(new Event("event-1",1,10))==Result.APPLIED,"first failed");
check(apply(new Event("event-2",2,20))==Result.APPLIED,"gap did not recover");
check(apply(new Event("event-2",2,20))==Result.DUPLICATE,"duplicate accepted");
try {apply(new Event("stale-copy",1,10));throw new IllegalStateException("old version accepted");}
catch(VersionConflict expected) {state(30,2,2);}
System.out.println("order: gapRejected=true recoveredVersion=2 duplicateSkipped=true staleRejected=true balance=30 inbox=2");
return;
}
check(args[0].equals("duplicate")||args[0].equals("split"),"unknown mode");
String queue="me12.delivery."+args[0],id="event-"+args[0];var event=new Event(id,1,10);
var factory=new ConnectionFactory();factory.setHost("rabbit");factory.setVirtualHost("lab");
factory.setUsername("app");factory.setPassword("app-lab-only");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);
publisher.confirmSelect();publisher.basicPublish("",queue,true,new AMQP.BasicProperties.Builder()
.messageId(id).deliveryMode(2).build(),"1:10".getBytes(StandardCharsets.UTF_8));
publisher.waitForConfirmsOrDie(5000);
try(var ch=connection.createChannel()) {
var message=receive(ch,queue);check(message.getProps().getMessageId().equals(id),"wrong initial ID");
if(args[0].equals("split")) {
try(var c=database()){check(claim(c,id),"claim failed");} // Auto-commit: deliberately wrong split.
try(var c=database();var s=c.createStatement()) {
c.setAutoCommit(false);s.executeUpdate("UPDATE delivery_account SET balance=10,version=1 WHERE id='order-42'");c.rollback();
}
} else check(apply(event)==Result.APPLIED,"first apply failed");
// Close without ACK after the database transaction has already finished.
}
try(var ch=connection.createChannel()) {
var message=receive(ch,queue);
check(message.getEnvelope().isRedeliver()&&id.equals(message.getProps().getMessageId()),"expected same event redelivery");
check(apply(event)==Result.DUPLICATE,"redelivery not deduplicated");
ch.basicAck(message.getEnvelope().getDeliveryTag(),false);
check(ch.queueDeclarePassive(queue).getMessageCount()==0,"unexpected ready messages");
}
int balance=args[0].equals("split")?0:10,version=args[0].equals("split")?0:1;
state(balance,version,1);
System.out.println(args[0]+": redelivered=true duplicateSkipped=true balance="+balance+" inbox=1");
}
}
}claim 和 UPDATE 使用同一个 JDBC Connection,并关闭自动提交。业务处理或提交失败时,先尝试 rollback,再抛出原异常;如果回滚也失败,将它附加到原异常的 suppressed 列表,避免丢失最初的失败原因。消费者只有确认数据库已提交,或确认该事件此前完整提交过,才执行 ACK。
VersionConflict 表示预期旧版本不满足,并不会在程序里伪装成成功。事务回滚也删除本次尚未提交的 Inbox 行,因此修复顺序之后仍能处理原事件。
提交后关闭 Channel,再消费同一事件
定义运行函数:
run_delivery() {
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 '-Dorg.slf4j.simpleLogger.defaultLogLevel=warn' \
-cp 'rabbit/target/classes:rabbit/target/dependency/*' example.DeliveryLab "$@"
}
run_delivery duplicate预期:
duplicate: redelivered=true duplicateSkipped=true balance=10 inbox=1程序先发布一条持久消息并等待真实 publisher confirm;消费时提交数据库事务,然后故意关闭尚未 ACK 的 channel。新的 channel 重新读取,检查 redelivered=true 且事件 ID 相同,再执行相同业务方法。第二次遇到 Inbox 唯一键,因此余额保持 10。
这个故障注入隔离的是“数据库已提交、Broker 尚未确认”的窗口。关闭 channel 足以让未确认交付返回队列;它没有额外证明断电、磁盘损坏或多节点失效后的持久性,那些需要另外的存储与副本演练。
错误拆分为何会让去重吞掉业务
run_delivery split预期:
split: redelivered=true duplicateSkipped=true balance=0 inbox=1负例使用 JDBC 自动提交单独写入 Inbox,再用另一个事务更新余额并 rollback。消息重投后,Inbox 已存在,于是处理被跳过。消费者仍可能记录“去重成功”,但真实余额一直为 0。
遇到这类历史数据,直接删 Inbox 后重放可能再次执行某些已经成功的外部操作。应先按事件 ID 查询业务表、操作流水和外部结果,确认需要补写哪些变化,再执行受控修复。设计阶段保持同事务约束,可以减少这种昂贵的事后判断。
两个 JDBC 会话同时争抢同一事件
run_delivery concurrent预期:
concurrent: applied=1 duplicate=1 balance=10 inbox=1两个工作线程各建立独立 JDBC Connection,在同一计数门上同时出发。实验检查返回值集合中恰有一次 APPLIED 和一次 DUPLICATE,并再次查询数据库确认余额、版本与 Inbox 条数。去重来自 PostgreSQL 主键约束,并非 Java 的 synchronized 或内存锁。
如果取得唯一键的事务随后回滚,其他等待者可以继续尝试插入并完成业务。不要把等待唯一冲突的连接永久占住:为生产 SQL 配置合理的锁等待和语句超时,回滚后再按业务重试策略处理。PostgreSQL 锁与并发
用业务版本处理乱序与重复
分区顺序需要配合执行顺序
选择同一个 Kafka key 或 RocketMQ message group 可以建立局部顺序。消费者把记录交给多个线程后,实际完成顺序仍可能变化;RabbitMQ 中多个消费者、重投和优先级也会影响观察到的处理顺序。对同一订单串行执行是最直接的办法,跨订单则可并行。
应用版本号使处理规则能够独立检查。例如当前余额版本为 1,版本 2 的增量允许执行,版本 3 则存在缺口,版本 1 的新事件 ID 则属于旧数据或生产错误。RocketMQ FIFO 的分组顺序
版本判断需要使用数据库条件更新。先 SELECT 版本为 1,随后无条件 UPDATE,会让两个并发请求都根据旧值推进;WHERE version=1 将判断与更新绑定在同一 SQL 中。更新影响 0 行后,再读取当前版本以区分记录不存在、版本过旧与版本超前,决定下一步。
增量事件不能直接跳到较新版本
run_delivery order程序按以下顺序调用真实 SQL:
当前版本 0
收到 event-2,版本 2,增加 20 → 缺少版本 1,回滚,Inbox 仍为 0
收到 event-1,版本 1,增加 10 → 余额 10,版本 1
重放 event-2,版本 2,增加 20 → 余额 30,版本 2
再次收到 event-2 → 重复跳过
新 ID 携带版本 1 → 旧版本拒绝,Inbox 不增加预期最终输出:
order: gapRejected=true recoveredVersion=2 duplicateSkipped=true staleRejected=true balance=30 inbox=2UPDATE ... WHERE version=event.version-1 对增量事件尤其重要。若收到版本 3 就直接把版本号设成 3,版本 2 中的余额变化可能永远无法补回。
快照事件的规则可以不同。如果版本 3 携带完整余额快照,消费者可能允许以较新快照覆盖旧状态,再忽略更旧的快照;前提是快照内容完整、版本可比较且来源明确。不能把这种“保留最新”策略用于逐条扣款、发送通知等不可丢失的动作。
缺口、重试与跨对象关系
顺序缺口出现后,可以短暂保留事件等待前序,查询权威状态重建当前投影,或暂停对应业务对象并告警。必须限制等待时间和缓存量:一直等待不存在的版本会形成永久积压,直接跳过又可能破坏业务。
同一 Topic 中不同订单可各自递增版本;跨订单转账等关系通常需要数据库事务、Saga 或其他业务协调方式。给所有事件加一个全局序号会引入集中排序成本,也不能替代跨系统提交。
消息重试应保留原事件 ID、版本与原始发生时间,另记本次尝试次数和重试时间。重新生成 ID 会让每次重试都绕过 Inbox;反过来,修正业务内容形成新的合法事件时,则应赋予新事件身份并关联原事件,便于审计两次变化。
从去重命中率追查真正的失败
需要一起观察的数据库与消费状态
| 观察结果 | 可能原因 | 下一步检查 |
|---|---|---|
| Inbox 增加,业务记录未变化 | 两者拆分提交;状态机拒绝后仍保留 Inbox | 追踪同一事件的事务 Connection、commit 和 rollback |
| 余额重复增加,Inbox 没有冲突 | ID 每次重发改变;consumer 名称每个实例不同 | 比对事件信封与唯一键组成 |
| 去重命中持续升高 | ACK 丢失、上游重复发布、处理超时重投 | 关联发送确认、重投标志和处理耗时 |
| 版本超前持续堆积 | 前序事件未发布、保留窗口不足、错误扩分区 | 查询生产者事实与对应版本的发布记录 |
| 更新影响 0 行但没有错误 | 业务代码忽略 SQL 返回值 | 明确已完成、版本冲突与缺失记录的分支 |
| 两个消费者都认为自己是首次 | 先查后写;唯一约束未部署到真实表 | 检查数据库约束,不只查看 ORM 注解 |
实验中可用普通应用身份直接查询表。容器内的运维 shell 不等于数据库超级用户授权;下面显式使用 app:
docker compose exec -T -e PGPASSWORD=app-lab-only postgres \
psql -h 127.0.0.1 -U app -d lab -v ON_ERROR_STOP=1 \
-c "SELECT id,balance,version FROM lab.delivery_account WHERE id='order-42'" \
-c "SELECT consumer,event_id FROM lab.delivery_inbox WHERE consumer='balance-v1' ORDER BY event_id"紧接 order 模式运行时,应看到余额 30、版本 2,以及 event-1、event-2 两条 Inbox。运行其他模式会重置这些专属数据,需按最近一次模式解释结果。
清理期限由可重复窗口决定
Inbox 记录通常会持续增长。清理期限需要覆盖上游重放、重试、灾备恢复和业务人工补发的窗口;如果 Broker 仍可重放一年前的事件,而 Inbox 只保留七天,恢复旧位点就可能重新产生业务变化。
对长期不能重复的业务,可以将幂等约束沉淀到订单号、支付流水号等永久业务唯一键。投影重建使用独立 consumer 名称和独立目标表,避免用线上处理方的 Inbox 阻止重放,也避免重建任务重做短信、支付等外部操作。RabbitMQ 可靠性交付与去重
停止实验用 docker compose stop,保留数据库和队列。确认项目全为可丢弃实验数据后,可执行 docker compose down --volumes 删除本项目容器、网络与命名卷。生产修复先导出受影响事件范围,不使用实验程序的重置逻辑。
权威资料与规范地址
| 查阅内容 | 官方地址 |
|---|---|
| Kafka 交付语义 | Design:https://kafka.apache.org/43/design/design/ |
| 消费确认与重投 | RabbitMQ Confirms:https://www.rabbitmq.com/docs/confirms |
| 事件身份与唯一范围 | CloudEvents 1.0.2:https://raw.githubusercontent.com/cloudevents/spec/v1.0.2/cloudevents/spec.md |
| 插入与唯一冲突处理 | PostgreSQL INSERT:https://www.postgresql.org/docs/18/sql-insert.html |
| 事务原子性 | PostgreSQL Transactions:https://www.postgresql.org/docs/18/tutorial-transactions.html |
| 锁等待与并发控制 | PostgreSQL Explicit Locking:https://www.postgresql.org/docs/18/explicit-locking.html |
| 按消息组保持顺序 | RocketMQ FIFO:https://rocketmq.apache.org/docs/featureBehavior/03fifomessage/ |
| 重投与幂等处理 | RabbitMQ Reliability:https://www.rabbitmq.com/docs/reliability |
