业务系统性能改造手册 · 异步拆链

让消息投递与消费可重复

2026-08-057 min read异步拆链
摘要

用本地事件表把订单状态与待发送事件放进同一数据库事务,再以稳定事件标识、租约式投递和消费端唯一约束承接至少一次交付,让重复消息不再产生重复业务结果。

04-01 已经确定:订单提交接口返回时,核心订单事实必须可以查询,通知等后续动作允许延迟完成。接下来要补上那段最危险的空白——数据库已经提交,后续动作却还没有被可靠接走。

直接在事务方法里调用消息客户端,看起来只多了一行代码,实际上会留下两个无法同时消除的窗口:先发消息,数据库可能回滚;先提交数据库,进程可能在发消息前退出。把发送代码换成 @Async,只会再增加一个内存队列,并没有让两个资源共享一次提交。

小编对这类链路的要求并不神秘:接受重复,拒绝静默丢失。生产端把业务状态和待发送事件一起落到本地数据库;中继可能重复发布;消费端用稳定事件标识,让同一业务副作用只生效一次。

一、先承认数据库与消息系统之间的窗口

订单数据库和消息系统是两个独立资源。本系列不引入跨服务分布式事务,应用也就不能用一次本地 commit() 同时证明“订单已提交”和“消息已进入 Broker”。

text
1方案 A:先发消息,再提交订单 2 3消息已被 Broker 接受 ──进程退出/事务回滚──> 消费者看见一笔不存在的订单 4 5方案 B:先提交订单,再发消息 6 7订单提交成功 ──进程退出/发送失败──> 后续动作永远不知道这笔订单

“事务提交后监听事件”可以避开读取未提交数据,却仍有提交后、监听器执行前的进程退出窗口。消息系统提供的生产确认也只能说明它所定义的接收阶段已经完成,不能倒过来提交本地数据库,更不能证明消费者的业务副作用已经完成。

本地事件表把必须原子成立的范围缩回单个数据库:

text
1订单事务 2 ├─ 写 orders 3 └─ 写 outbox_event(PENDING, stable eventId) 45 └─ 同一次本地事务提交 6 7中继 8 └─ claim → publish(eventId) → broker confirm → mark PUBLISHED 9 10消费者 11 └─ receive → deduplicate(eventId) + local effect → commit → acknowledge

这里没有 exactly-once。中继在 Broker 确认后、标记 PUBLISHED 前退出,会再次发布;消费者在本地提交后、确认消息前退出,也会再次收到。设计的目标是让这些重复可识别、可安全重放。

二、订单和待发送事件必须同事务落盘

订单数据与待发送事件必须在同一数据库事务中落盘才能关闭数据库与消息系统之 以下 DDL 面向 MySQL/InnoDB 的示例形状。项目的 MySQL 版本尚未确定,JSON、时间精度、索引长度、字符集和在线变更方式要在目标版本验证。

sql
1CREATE TABLE outbox_event ( 2 event_id CHAR(36) CHARACTER SET ascii COLLATE ascii_bin NOT NULL, 3 aggregate_type VARCHAR(64) NOT NULL, 4 aggregate_id VARCHAR(128) NOT NULL, 5 event_type VARCHAR(128) NOT NULL, 6 event_version BIGINT NOT NULL, 7 payload JSON NOT NULL, 8 status VARCHAR(16) NOT NULL, 9 lease_owner VARCHAR(128) NULL, 10 lease_until TIMESTAMP(6) NULL, 11 created_at TIMESTAMP(6) NOT NULL, 12 published_at TIMESTAMP(6) NULL, 13 PRIMARY KEY (event_id), 14 KEY idx_outbox_claim (status, lease_until, created_at) 15) ENGINE = InnoDB; 16 17CREATE TABLE consumed_event ( 18 id BIGINT NOT NULL AUTO_INCREMENT, 19 consumer_name VARCHAR(128) NOT NULL, 20 event_id CHAR(36) CHARACTER SET ascii COLLATE ascii_bin NOT NULL, 21 consumed_at TIMESTAMP(6) NOT NULL, 22 PRIMARY KEY (id), 23 CONSTRAINT uq_consumed_event_identity UNIQUE (consumer_name, event_id) 24) ENGINE = InnoDB; 25 26CREATE TABLE notification_job ( 27 event_id CHAR(36) CHARACTER SET ascii COLLATE ascii_bin NOT NULL, 28 order_id BIGINT NOT NULL, 29 source_version BIGINT NOT NULL, 30 status VARCHAR(16) NOT NULL, 31 created_at TIMESTAMP(6) NOT NULL, 32 PRIMARY KEY (event_id) 33) ENGINE = InnoDB; 34 35-- 订单表需要保存入口幂等身份和本次受理事件身份;实际变更按现有 schema 编写。 36ALTER TABLE orders 37 ADD COLUMN request_id VARCHAR(128) NOT NULL, 38 ADD COLUMN acceptance_event_id CHAR(36) 39 CHARACTER SET ascii COLLATE ascii_bin NOT NULL, 40 ADD CONSTRAINT uq_orders_tenant_request UNIQUE (tenant_id, request_id), 41 ADD CONSTRAINT uq_orders_acceptance_event UNIQUE (acceptance_event_id);

event_id 在业务事务开始前生成,之后投递和消费始终沿用同一个值。Broker 自己生成的消息 ID、每次发送的请求 ID,都不能替代这份业务事件身份。event_version 表示订单事实的版本;payload 保存 04-01 已确认的最小事实快照,不塞 JPA 实体、请求对象或可变的“当前值”。

下面的标准 JDBC 代码演示订单与事件共用同一个 Connection 和本地事务。表中的订单列只是最小示例,真实 schema 与权限校验仍需接入案例工程。

java
1import java.sql.Connection; 2import java.sql.PreparedStatement; 3import java.sql.SQLException; 4import java.sql.Timestamp; 5import java.time.Instant; 6import java.util.Optional; 7import javax.sql.DataSource; 8 9record SubmitOrder( 10 String requestId, 11 String eventId, 12 long orderId, 13 long tenantId, 14 long version, 15 String eventPayload) {} 16 17record AcceptedOrder(long orderId, String eventId) {} 18 19enum TransactionOutcome { 20 COMMITTED, 21 ROLLED_BACK, 22 UNKNOWN 23} 24 25record SubmissionResult( 26 TransactionOutcome outcome, 27 String requestId, 28 long orderId, 29 String eventId, 30 Throwable failure) {} 31 32final class OrderWithEventWriter { 33 private static final String INSERT_ORDER = """ 34 INSERT INTO orders( 35 id, tenant_id, request_id, acceptance_event_id, 36 version, status, created_at) 37 VALUES (?, ?, ?, ?, ?, 'ACCEPTED', ?) 38 """; 39 private static final String INSERT_EVENT = """ 40 INSERT INTO outbox_event( 41 event_id, aggregate_type, aggregate_id, event_type, 42 event_version, payload, status, created_at) 43 VALUES (?, 'ORDER', ?, 'ORDER_ACCEPTED', ?, CAST(? AS JSON), 'PENDING', ?) 44 """; 45 private static final String FIND_BY_REQUEST = """ 46 SELECT id, acceptance_event_id 47 FROM orders 48 WHERE tenant_id = ? AND request_id = ? 49 """; 50 51 private final DataSource dataSource; 52 53 OrderWithEventWriter(DataSource dataSource) { 54 this.dataSource = dataSource; 55 } 56 57 SubmissionResult submit(SubmitOrder command) throws SQLException { 58 Instant now = Instant.now(); 59 try (Connection connection = dataSource.getConnection()) { 60 if (!connection.getAutoCommit()) { 61 throw new SQLException("example requires a locally managed connection"); 62 } 63 connection.setAutoCommit(false); 64 try { 65 insertOrder(connection, command, now); 66 insertEvent(connection, command, command.eventId(), now); 67 } catch (SQLException | RuntimeException beforeCommitFailure) { 68 try { 69 connection.rollback(); 70 } catch (SQLException rollbackFailure) { 71 beforeCommitFailure.addSuppressed(rollbackFailure); 72 return result(TransactionOutcome.UNKNOWN, command, beforeCommitFailure); 73 } 74 return result(TransactionOutcome.ROLLED_BACK, command, beforeCommitFailure); 75 } 76 77 try { 78 connection.commit(); 79 return result(TransactionOutcome.COMMITTED, command, null); 80 } catch (SQLException commitFailure) { 81 // 服务端可能已经提交,只是确认没有到达客户端。这里不能安全重放, 82 // 也不能用随后一次 rollback 把结果改写成 ROLLED_BACK。 83 return result(TransactionOutcome.UNKNOWN, command, commitFailure); 84 } 85 } 86 } 87 88 Optional<AcceptedOrder> reconcile(long tenantId, String requestId) throws SQLException { 89 try (Connection connection = dataSource.getConnection(); 90 PreparedStatement statement = connection.prepareStatement(FIND_BY_REQUEST)) { 91 statement.setLong(1, tenantId); 92 statement.setString(2, requestId); 93 try (var rows = statement.executeQuery()) { 94 if (!rows.next()) { 95 return Optional.empty(); 96 } 97 return Optional.of(new AcceptedOrder(rows.getLong(1), rows.getString(2))); 98 } 99 } 100 } 101 102 private SubmissionResult result( 103 TransactionOutcome outcome, SubmitOrder command, Throwable failure) { 104 return new SubmissionResult( 105 outcome, 106 command.requestId(), 107 command.orderId(), 108 command.eventId(), 109 failure); 110 } 111 112 private void insertOrder( 113 Connection connection, SubmitOrder command, Instant now) throws SQLException { 114 try (PreparedStatement statement = connection.prepareStatement(INSERT_ORDER)) { 115 statement.setLong(1, command.orderId()); 116 statement.setLong(2, command.tenantId()); 117 statement.setString(3, command.requestId()); 118 statement.setString(4, command.eventId()); 119 statement.setLong(5, command.version()); 120 statement.setTimestamp(6, Timestamp.from(now)); 121 statement.executeUpdate(); 122 } 123 } 124 125 private void insertEvent( 126 Connection connection, 127 SubmitOrder command, 128 String eventId, 129 Instant now) throws SQLException { 130 try (PreparedStatement statement = connection.prepareStatement(INSERT_EVENT)) { 131 statement.setString(1, eventId); 132 statement.setString(2, Long.toString(command.orderId())); 133 statement.setLong(3, command.version()); 134 statement.setString(4, command.eventPayload()); 135 statement.setTimestamp(5, Timestamp.from(now)); 136 statement.executeUpdate(); 137 } 138 } 139}

JDBC 在关闭自动提交后,会把语句归入由 commit()rollback() 结束的事务。Java SE JDBC Connection 案例工程若由 Spring 管理事务,应使用同一个事务管理器和数据源执行两次写入,不能在内部手工提交连接。

requestId 是入口幂等身份,eventId 是这次业务事件的稳定身份;同一次调用及其重试不能重新生成。返回 COMMITTED 时订单和事件都已提交;只有提交前失败且 rollback() 已确认时才返回 ROLLED_BACKcommit() 异常或回滚异常返回 UNKNOWN,上层不得直接重放,而应在数据库连接恢复后通过主库按 tenantId + requestId 调用 reconcile():查到记录就沿用其中的 acceptance_event_id,查不到也要在能够排除复制延迟和连接故障后,才用同一组稳定身份再次提交。

JDBC 文档规定 commit() 使事务修改持久化,rollback() 撤销当前事务修改,但没有承诺“commit 抛异常后再 rollback 就能证明之前未提交”。因此代码把提交异常保留为待对账结果,而不是普通 SQL 失败。这是入口请求的结果核对边界,不等于消息消费去重。

三、中继用租约认领,发布确认后再标记

多个实例都可能扫描 outbox_event。只执行 SELECT ... WHERE status='PENDING',两个实例会同时拿到同一批;先把状态改成 PUBLISHED 再发送,则会在发送失败时永久丢失。

一个可操作的顺序是:

  1. 在短数据库事务内锁定候选行,写入 PUBLISHING + lease_owner + lease_until 后提交。
  2. 在事务外发布,避免网络等待期间一直持有数据库行锁和连接。
  3. 收到目标 Broker 所定义的成功确认后,带 event_id + lease_owner 条件标记 PUBLISHED
  4. 进程中断时让租约到期,其他实例可以重新认领;重复发布由消费者承接。

MySQL 8.4 文档说明,SKIP LOCKED 不等待已锁行,并明确提到它适合多个会话访问队列式表;它返回的是不一致视图,不适合普通事务查询。MySQL 8.4 Locking Reads 目标 MySQL 版本、隔离级别和复制格式未确认前,不能直接照搬。

核心 SQL 如下:

sql
1START TRANSACTION; 2 3SELECT event_id, event_type, payload 4FROM outbox_event 5WHERE status = 'PENDING' 6 OR (status = 'PUBLISHING' AND lease_until < CURRENT_TIMESTAMP(6)) 7ORDER BY created_at, event_id 8LIMIT :batch_size 9FOR UPDATE SKIP LOCKED; 10 11UPDATE outbox_event 12SET status = 'PUBLISHING', 13 lease_owner = :worker_id, 14 lease_until = :lease_until 15WHERE event_id IN (:claimed_event_ids); 16 17COMMIT; 18 19-- 每条消息在事务外发布并获得目标 Broker 的确认后执行: 20UPDATE outbox_event 21SET status = 'PUBLISHED', published_at = CURRENT_TIMESTAMP(6), 22 lease_owner = NULL, lease_until = NULL 23WHERE event_id = :event_id 24 AND status = 'PUBLISHING' 25 AND lease_owner = :worker_id;

实际代码不能把命名参数原样交给 JDBC,也不能在空 ID 集合上拼 IN ();应使用固定上限的占位符或批量更新。认领事务必须关闭自动提交,否则 FOR UPDATE 锁会随语句结束而释放。租约应长于目标发布确认的正常上界,并观察发布时长;租约过期与旧发布同时进行时,重复消息是预期结果。

PUBLISHED 只说明目标 Broker 的生产确认已经满足适配器契约,不代表任何消费者完成。MQ 类型、客户端版本和确认等级尚未确定,发布接口要把语义写清楚:

java
1record OutboxMessage(String eventId, String eventType, String payload) {} 2 3interface ConfirmedMessagePublisher { 4 /** 5 * 仅在目标 Broker 按已配置确认等级接受消息后返回; 6 * 拒绝、超时或结果未知时抛出异常,调用方不得标记 PUBLISHED。 7 */ 8 void publish(OutboxMessage message) throws Exception; 9}

扫描周期、退避、最大尝试次数、失败终态和积压告警会影响长期运行,但属于 04-03。本篇只要求租约过期后仍能重新发现事件,且每次状态变化可以按 event_id 追踪。

四、消费去重和本地副作用放进同一事务

消费去重记录与本地副作用要在同一事务内完成使重复消息只产生一次业务结果 生产端的唯一键不能阻止 Broker 重投,也不能阻止多个消费实例同时处理。消费者要为“哪个消费逻辑处理过哪个事件”建立命名唯一约束 uq_consumed_event_identity(consumer_name, event_id)。同一个事件可以被通知和报表两个消费者分别处理,但同一消费者只取得一次本地生效资格。

下面用 ORDER_ACCEPTED 创建一条本地通知任务。代码使用普通 INSERT;只有目标版本适配器确认错误来自 uq_consumed_event_identity 时,才返回“已经处理”。SQLState 23000 表示一类完整性约束错误,可能是主键、唯一键或外键,单独检查它不够。Java SE SQLIntegrityConstraintViolationException

java
1import java.sql.Connection; 2import java.sql.PreparedStatement; 3import java.sql.SQLException; 4import java.sql.Timestamp; 5import java.time.Instant; 6import javax.sql.DataSource; 7 8record OrderAcceptedEvent(String eventId, long orderId, long sourceVersion) {} 9 10interface ConsumedEventDuplicateClassifier { 11 /** 只识别 uq_consumed_event_identity,不能把任意 SQLState 23000 当成重复。 */ 12 boolean isConsumedEventIdentityDuplicate(SQLException failure); 13} 14 15final class OrderAcceptedConsumer { 16 private static final String CLAIM_EVENT = """ 17 INSERT INTO consumed_event(consumer_name, event_id, consumed_at) 18 VALUES ('notification-planner', ?, ?) 19 """; 20 private static final String CREATE_NOTIFICATION = """ 21 INSERT INTO notification_job( 22 event_id, order_id, source_version, status, created_at) 23 VALUES (?, ?, ?, 'PENDING', ?) 24 """; 25 26 private final DataSource dataSource; 27 private final ConsumedEventDuplicateClassifier duplicateClassifier; 28 29 OrderAcceptedConsumer( 30 DataSource dataSource, 31 ConsumedEventDuplicateClassifier duplicateClassifier) { 32 this.dataSource = dataSource; 33 this.duplicateClassifier = duplicateClassifier; 34 } 35 36 void handle(OrderAcceptedEvent event) throws SQLException { 37 try (Connection connection = dataSource.getConnection()) { 38 if (!connection.getAutoCommit()) { 39 throw new SQLException("example requires a locally managed connection"); 40 } 41 connection.setAutoCommit(false); 42 Instant now = Instant.now(); 43 try { 44 claim(connection, event.eventId(), now); 45 } catch (SQLException claimFailure) { 46 rollbackOrThrow(connection, claimFailure); 47 if (duplicateClassifier.isConsumedEventIdentityDuplicate(claimFailure)) { 48 return; 49 } 50 throw claimFailure; 51 } 52 53 try { 54 createNotification(connection, event, now); 55 } catch (SQLException | RuntimeException failure) { 56 try { 57 connection.rollback(); 58 } catch (SQLException rollbackFailure) { 59 failure.addSuppressed(rollbackFailure); 60 } 61 throw failure; 62 } 63 64 try { 65 connection.commit(); 66 } catch (SQLException commitFailure) { 67 // 可能已经提交,只是确认丢失。抛错让适配层不要 ack; 68 // 重投后由唯一约束核对,不在这里用 rollback 宣称未生效。 69 throw new SQLException( 70 "consumer commit outcome unknown; message must not be acknowledged", 71 commitFailure); 72 } 73 } 74 } 75 76 private void claim(Connection connection, String eventId, Instant now) throws SQLException { 77 try (PreparedStatement statement = connection.prepareStatement(CLAIM_EVENT)) { 78 statement.setString(1, eventId); 79 statement.setTimestamp(2, Timestamp.from(now)); 80 if (statement.executeUpdate() != 1) { 81 throw new SQLException("consumed event insert changed an unexpected row count"); 82 } 83 } 84 } 85 86 private void rollbackOrThrow(Connection connection, SQLException original) 87 throws SQLException { 88 try { 89 connection.rollback(); 90 } catch (SQLException rollbackFailure) { 91 rollbackFailure.addSuppressed(original); 92 throw rollbackFailure; 93 } 94 } 95 96 private void createNotification( 97 Connection connection, OrderAcceptedEvent event, Instant now) throws SQLException { 98 try (PreparedStatement statement = connection.prepareStatement(CREATE_NOTIFICATION)) { 99 statement.setString(1, event.eventId()); 100 statement.setLong(2, event.orderId()); 101 statement.setLong(3, event.sourceVersion()); 102 statement.setTimestamp(4, Timestamp.from(now)); 103 statement.executeUpdate(); 104 } 105 } 106}

MySQL 8.4 中普通重复条目错误是 ER_DUP_ENTRY(错误码 1062、SQLState 23000),但 23000 还覆盖其他完整性错误。MySQL 8.4 Server Error Reference ConsumedEventDuplicateClassifier 必须由确定版本的 MySQL 驱动适配器实现,并同时核对 1062 和命名约束 uq_consumed_event_identity;约束名如何暴露在异常中要用目标驱动的集成测试锁定。外键错误、数据截断、连接异常、其他唯一键冲突和无法识别的异常全部失败,不进入“已处理”分支。

消息适配层只在 handle() 提交成功或严格确认目标唯一键已经存在后,才向 Broker 确认消费。提交后、消息确认前退出,会再次进入 handle(),唯一约束让第二次调用直接返回。消费事务的 commit() 异常同样按 UNKNOWN 处理:不确认消息,等待重投后通过该唯一约束核对。消息确认发生在本地提交之前,则可能在进程退出时丢掉尚未生效的副作用。

这段代码保护的是同一 MySQL 内的 consumed_eventnotification_job。若消费者直接调用短信、支付或其他外部服务,本地去重记录与远端副作用之间又出现一个双写窗口:远端成功、本地提交前退出,重投后仍可能再次调用。此时需要远端接受稳定幂等键,或把外部调用建模成可查询、可核对的状态机;不能宣称一张去重表解决了所有副作用。

五、用故障注入证明“可重复”

正常请求跑通只证明快乐路径。验证记录要在每个提交窗口主动中断,并按事件 ID 核对数据库、Broker 观察和最终业务结果:

yaml
1repeatableDeliveryVerification: 2 experimentId: ${EXPERIMENT_ID} 3 versions: 4 application: ${VERSION} 5 mysqlAndDriver: ${VERSIONS} 6 brokerAndClient: ${VERSIONS} 7 event: 8 requestId: ${STABLE_REQUEST_ID} 9 eventId: ${STABLE_EVENT_ID} 10 orderId: ${ORDER_ID} 11 eventTypeAndVersion: ${TYPE_AND_VERSION} 12 injections: 13 - point: ${BEFORE_ORDER_TRANSACTION_COMMIT} 14 expected: ${ORDER_AND_EVENT_BOTH_ABSENT_OR_UNKNOWN_RECONCILED} 15 observed: ${EVIDENCE} 16 - point: ${COMMIT_SUCCEEDS_BUT_CLIENT_RECEIVES_EXCEPTION} 17 expected: ${UNKNOWN_THEN_RECONCILED_BY_REQUEST_ID_NO_BLIND_REPLAY} 18 observed: ${OUTCOME_AND_RECONCILIATION_EVIDENCE} 19 - point: ${AFTER_COMMIT_BEFORE_RELAY_CLAIM} 20 expected: ${EVENT_REMAINS_DISCOVERABLE} 21 observed: ${EVIDENCE} 22 - point: ${AFTER_BROKER_CONFIRM_BEFORE_MARK_PUBLISHED} 23 expected: ${MESSAGE_MAY_REPEAT_NO_DUPLICATE_LOCAL_EFFECT} 24 observed: ${EVIDENCE} 25 - point: ${DURING_CONCURRENT_CLAIM} 26 expected: ${LEASE_OWNERSHIP_TRACEABLE_DUPLICATE_TOLERATED} 27 observed: ${EVIDENCE} 28 - point: ${AFTER_CONSUMER_COMMIT_BEFORE_MESSAGE_ACK} 29 expected: ${REDELIVERY_DEDUPLICATED} 30 observed: ${EVIDENCE} 31 - point: ${NON_TARGET_INTEGRITY_OR_DATA_ERROR_WHILE_CLAIMING} 32 expected: ${FAILED_NOT_CLASSIFIED_AS_ALREADY_CONSUMED} 33 observed: ${SQLSTATE_VENDOR_CODE_CONSTRAINT_AND_ACK_EVIDENCE} 34 counts: 35 committedOrders: ${COUNT} 36 outboxEventsByStatus: ${COUNTS} 37 brokerDeliveriesForEvent: ${COUNT_OR_BROKER_EVIDENCE} 38 consumerHandleAttempts: ${COUNT} 39 consumedEventRows: ${COUNT} 40 notificationJobs: ${COUNT} 41 invariants: 42 committedOrderHasOutboxEvent: ${PASS_FAIL} 43 unknownCommitReconciledByStableRequestId: ${PASS_FAIL} 44 noEffectWithoutCommittedOrder: ${PASS_FAIL} 45 onlyTargetUniqueConflictMeansAlreadyConsumed: ${PASS_FAIL} 46 oneLocalEffectPerConsumerAndEvent: ${PASS_FAIL} 47 everyNonTerminalStateIsQueryable: ${PASS_FAIL}

“最终只看到一条通知任务”还不够。测试要确认重复消息真的发生过,否则消费去重分支可能从未执行。并发认领也要覆盖租约过期和旧 Worker 晚确认的情形,观察条件更新是否拒绝覆盖新租约所有者。

吞吐和延迟仍要测,但不能通过减小扫描周期、扩大批次或无限并发把数据库拖回 03-03 已解决的排队状态。发布批次、连接占用、行锁等待、从事件创建到 Broker 确认的耗时,以及从接收到本地副作用提交的耗时,都应进入同一实验窗口。

六、上线边界与回退

上线前至少保留这些查询入口:按 event_id 查看订单版本、Outbox 状态、租约所有者和发布时间;按 consumer_name + event_id 查看消费记录;按业务键核对最终副作用。日志里只有一句“发送失败”,无法支撑故障恢复。

回退发送中继时,不能删除尚未发布的事件;停止新认领,等待当前租约结束,再保留表和查询工具。回退消费者时也要保留 consumed_event,否则旧消息再次出现时会失去去重历史。表数据保留周期必须覆盖 Broker 可能重投和人工恢复的窗口,具体时长留给目标 MQ 与运维策略确认。

这套实现适合业务状态与事件表位于同一事务型数据库、消费者的本地副作用也能与去重记录共用事务的场景。跨数据库写入、外部不可查询副作用和必须同步完成的核心事实不在它的保证范围内。

可靠链路的证据不是“消息通常只来一次”。订单和事件一起提交;发送确认前不提前标记;重复投递后,本地业务结果仍只有一份;每个中间状态都能按稳定 ID 查到。下一篇再处理消息持续失败、重试放大和积压超过预算时怎么办。