投递、幂等与顺序:至少一次系统怎样守住业务不变量
用失败窗口推导至多一次、至少一次和局部恰好一次,建立消费幂等、版本检查、乱序缓冲与并发隔离模型。
Kafka 生产者文档说明幂等生产只抑制单一生产会话中的协议重试重复;RabbitMQ 与 RocketMQ 文档都把消费确认和重试设计为独立机制。跨业务存储的幂等仍由应用状态承担。
投递语义是失败窗口的结果,不是配置名称
至多一次通常在处理前推进消费进度:不会由 Broker 主动重复,但进程崩溃会丢失未完成业务。至少一次在业务提交后确认:不会轻易静默丢失,但确认前崩溃会重投。所谓恰好一次必须先说明观察范围;对 Kafka 日志可用生产序号或事务抑制重复,不代表数据库扣款、邮件发送和第三方 HTTP 都只发生一次。
选择至少一次后,重复就不再是异常边角,而是正常输入。测试必须主动制造“业务已提交、确认未到达”“确认响应丢失”“消费者在提交后被杀死”等窗口。只有这些窗口下结果仍满足不变量,系统才真正具备可靠消费能力。
幂等记录必须与业务结果共享提交命运
最小模型以 eventId 或业务命令 id 为唯一键,记录 PROCESSING、SUCCEEDED 与可重试失败。若先在 Redis 写“已处理”再提交数据库,数据库回滚会造成永久跳过;若先提交数据库再异步标记 Redis,崩溃会再次执行。可靠做法是把幂等记录、目标状态变化与结果摘要放在同一个本地事务,使用唯一约束解决并发首写。
对天然状态赋值,例如把订单改为 PAID,仍需版本或前置状态校验,不能简单认为重复无害。第二次事件可能携带旧金额或旧版本。对非幂等副作用,如发送短信、调用支付机构,要把“准备发送”写成本地任务,再由有幂等键的适配器执行;没有幂等能力的外部系统只能通过查询、对账和人工补偿收敛。
顺序应缩小到业务实体而不是整个 Topic
同一聚合根使用稳定 key 进入同一队列或 partition,可以获得 Broker 层局部顺序;消费者内部并发、异步线程池和重试仍可能打乱提交顺序。处理器应按 key 串行化,或用 expectedVersion 做乐观检查。收到未来版本可短暂缓冲,收到旧版本直接判为重复或过期,缺口超过预算则转入修复队列。
不要为了一个少数状态机要求把全部消息放进单队列。那会让任意毒消息阻塞全局。顺序域应与业务冲突域一致,例如 orderId、accountId 或 aggregateId;跨聚合流程用 Saga 状态和因果关系协调,而不是假设全局时钟。
用所有权转移图读懂可靠性
失败矩阵比“保证不丢”更可执行
工程选择的核心是把不可避免的不确定窗口转成可重复、可查询、可补偿的状态。不要用无限重试掩盖未知结果,也不要用提前确认换取表面低 lag。
并发首写和处理中租约需要不同约束
两个消费者同时发现幂等记录不存在时,先查后写无法阻止双执行;必须依赖数据库唯一键或条件更新决定唯一赢家。PROCESSING 状态可以带 owner、leaseUntil 和 attempt,允许进程崩溃后接管,但接管者还需 fencing version,防止旧 owner 在暂停恢复后继续写入。
SUCCEEDED 状态应保存足以回答重复请求的结果摘要和业务版本。FAILED 不能只有一个布尔值:可重试失败允许租约到期后继续,永久失败要求人工或新版本事件介入。幂等记录的保留期必须覆盖 Broker 最大重投、死信重放和上游重试窗口;过早清理会让历史重复重新获得首写资格。
用两个 Java 状态模型固定不变量
下面的模型不连接真实 Broker。它们把本篇最容易混淆的状态、位置或所有权变化压缩成确定输出;随后再用 RabbitMQ、Kafka 或 RocketMQ 的集成环境验证协议、持久化、重平衡和故障时序。
javac --release 17 -Xlint:all -Werror examples/backend-development/message-event/delivery-idempotency-order/IdempotentConsumerDemo.java examples/backend-development/message-event/delivery-idempotency-order/PartitionOrderDemo.java
java -cp examples/backend-development/message-event/delivery-idempotency-order IdempotentConsumerDemo
java -cp examples/backend-development/message-event/delivery-idempotency-order PartitionOrderDemoeventId=evt-42 attempts=2 businessWrites=1 duplicateReturned=true
key=order-7 acceptedVersions=[3, 4, 5] rejectedVersions=[2, 4] finalVersion=5模型输出变化时,应定位是哪条业务不变量被修改,而不是只更新示例字符串。真实集成测试还要注入进程崩溃、响应丢失、Broker 切换和下游超时,因为这些窗口无法由纯内存模型模拟。
契约、权限与数据生命周期不能交给默认值
用分层测试证明协议与业务同时成立
把设计决定写成可以被证伪的承诺
指标、日志与容量要围绕状态迁移
容量模型至少计算平均/峰值每秒消息数、平均/p99 字节、保留时间、副本倍数、重试放大、批处理密度、消费者处理时长与下游并发。恢复测试要在稳定生产流量仍存在时注入积压,验证净排空速率,而不是暂停生产后测一个理想峰值。
上线与演练以可恢复为验收标准
故障演练依次注入生产响应丢失、Broker 节点切换、消费者提交后崩溃、下游超时、永久坏消息和热点 key。每次演练都回答三件事:业务不变量是否保持、重复或延迟是否有持久证据、系统是否在容量预算内自动或人工收敛。
