尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Java 设计模式之幂等消费者(Idempotent Consumer):微服务消息可靠处理的实战解析

Java 设计模式之幂等消费者(Idempotent Consumer):微服务消息可靠处理的实战解析 示例工程教程【免费下载链接】java-design-patternsDesign patterns implemented in Java项目地址https://gitcode.com/GitHub_Trending/ja/java-design-patterns点击查看免费下载本指南围绕 java-design-patterns 仓库中的 microservices-idempotent-consumer 模块展开深入剖析幂等消费者模式Idempotent Consumer的核心思想、适用场景与 Java 实现方式。你将掌握如何借助「消息唯一标识 状态机」的组合让同一个消息被重复消费时依然只产生一次业务副作用从而在支付、订单等对一致性要求极高的微服务场景中构建可靠的消息处理链路。模式概述幂等消费者模式Idempotent Consumer也被称为Idempotent Subscriber幂等订阅者、Repeatable Message Consumer可重复消息消费者或Safe Consumer安全消费者。其核心意图Intent是确保在微服务架构中同一个消息被消费多次时不会引发非预期的副作用。在分布式系统中消息的「恰好一次投递」几乎不可能被网络层保证消费者可能先处理成功、在发送 ACK 前崩溃导致消息被重新投递生产者重试也可能让同一条消息进入队列两次。幂等消费者模式不试图消灭重复而是让重复变得无害——无论消息到达多少次业务状态始终等同于只处理一次。真实世界场景为什么需要幂等消费在一个支付处理系统中保证支付消息的幂等性可以防止重复交易。例如如果用户的支付消息被意外处理了两次系统应识别第二条消息为重复消息并阻止其再次执行。通过为每条已处理消息存储唯一标识符例如交易 ID系统可以跳过任何重复消息。这确保了用户不会为同一笔交易被扣款两次从而维护系统完整性和客户满意度。通俗地说幂等消费者模式通过确保「同一消息被多次处理的结果与处理一次相同」来防止重复消息引发非预期副作用这让分布式系统中常见的重复投递变得安全。维基百科对幂等Idempotence的定义也直接点明了数学基础在计算机科学中幂等是某些运算的一种性质它们可以被多次应用而结果不会超出首次应用发生改变。工作流程上图仓库中 microservices-idempotent-consumer/etc/microservices-idempotent-consumer-flowchart.png展示了该模式的典型流程消息携带唯一标识进入消费者消费者先查询该标识是否已被处理已处理则直接跳过幂等命中未处理则执行业务逻辑并记录标识。下面我们结合仓库源码看这个流程在 Java 中如何落地。Java 程序化示例以订单创建与状态流转为例仓库的 microservices-idempotent-consumer 模块构建了一个幂等服务负责创建订单Request并推进其状态。其中create方法天然幂等——用相同订单 ID 多次调用返回相同结果、不会产生重复记录对于状态变更如启动、完成订单服务会校验状态迁移是否合法非法迁移抛出异常。RequestStateMachine保证订单状态按合法序列前进例如PENDING → STARTED → COMPLETED。领域模型 Request订单实体定义在 Request.java它是一个 JPA 实体以UUID作为主键标识状态枚举包含PENDING、STARTED、COMPLETED三态默认构造时状态为PENDINGEntity NoArgsConstructor Data public class Request { enum Status { PENDING, STARTED, COMPLETED } Id private UUID uuid; private Status status; public Request(UUID uuid) { this(uuid, Status.PENDING); } public Request(UUID uuid, Status status) { this.uuid uuid; this.status status; } }这里的uuid就是幂等性的关键载体它是每条消息/每次业务请求的唯一标识也恰好是持久化表的主键天然具备「重复写入即覆盖、查询即去重」的语义。数据访问层由 RequestRepository.java 提供直接继承 Spring Data JPA 的JpaRepositoryRequest, UUID开箱即用地获得findById、save、count等 CRUD 能力。RequestService幂等创建与受控的状态迁移核心服务类 RequestService.java 是幂等逻辑的第一层防线。create方法先按 UUID 查询命中则直接返回既有记录不再落库未命中才保存新记录——这正是「先查后写」的幂等实现Service public class RequestService { RequestRepository requestRepository; RequestStateMachine requestStateMachine; public RequestService( RequestRepository requestRepository, RequestStateMachine requestStateMachine) { this.requestRepository requestRepository; this.requestStateMachine requestStateMachine; } /** * Creates a new Request or returns an existing one by its UUID. This operation is idempotent: * performing it once or several times successively leads to an equivalent result. */ public Request create(UUID uuid) { OptionalRequest optReq requestRepository.findById(uuid); return optReq.orElseGet(() - requestRepository.save(new Request(uuid))); } public Request start(UUID uuid) { OptionalRequest optReq requestRepository.findById(uuid); if (optReq.isEmpty()) { throw new RequestNotFoundException(uuid); } return requestRepository.save(requestStateMachine.next(optReq.get(), Request.Status.STARTED)); } public Request complete(UUID uuid) { OptionalRequest optReq requestRepository.findById(uuid); if (optReq.isEmpty()) { throw new RequestNotFoundException(uuid); } return requestRepository.save(requestStateMachine.next(optReq.get(), Request.Status.COMPLETED)); } }三个方法的分工值得注意create(uuid)幂等写入。重复调用返回同一个对象仓库中记录数始终为 1。start(uuid)/complete(uuid)先查后转。查询不存在时抛出 RequestNotFoundException.java消息形如Request {uuid} not found存在则委托状态机做合法迁移迁移失败抛出InvalidNextStateException。RequestStateMachine把「重复」变成「异常」的状态机状态机类 RequestStateMachine.java 是幂等消费者模式的第二层防线。它不直接判断「是否重复」而是判断「本次迁移在当前状态下是否合法」任何对已推进状态的再次推进都会触发InvalidNextStateException从而在语义上拒绝重复Component public class RequestStateMachine { public Request next(Request req, Request.Status nextStatus) { String transitionStr String.format(Transition: %s - %s, req.getStatus(), nextStatus); switch (nextStatus) { case PENDING - throw new InvalidNextStateException(transitionStr); case STARTED - { if (Request.Status.PENDING.equals(req.getStatus())) { return new Request(req.getUuid(), Request.Status.STARTED); } throw new InvalidNextStateException(transitionStr); } case COMPLETED - { if (Request.Status.STARTED.equals(req.getStatus())) { return new Request(req.getUuid(), Request.Status.COMPLETED); } throw new InvalidNextStateException(transitionStr); } default - throw new InvalidNextStateException(transitionStr); } } }合法的迁移表如下当前状态目标状态结果PENDINGSTARTED合法返回 STARTED 新对象STARTEDCOMPLETED合法返回 COMPLETED 新对象任意状态PENDING非法抛InvalidNextStateExceptionSTARTED / COMPLETEDSTARTED非法抛InvalidNextStateExceptionCOMPLETED任意状态非法抛InvalidNextStateException异常类型定义在 InvalidNextStateException.java携带Transition: A - B的格式化信息便于定位是哪一步非法迁移。App 主程序完整演示幂等行为演示入口在 App.java。它通过 Spring Boot 的CommandLineRunner依次展示「重复创建不产生新记录」「重复启动被拒绝」「合法完成」三个关键行为SpringBootApplication Slf4j public class App { public static void main(String[] args) { var context SpringApplication.run(App.class, args); if (args.length 0 test.equals(args[0])) { // Close the context immediately during tests to prevent Tomcat/background threads from // hanging the JVM context.close(); } } Bean public CommandLineRunner run(RequestService requestService, RequestRepository requestRepository) { return args - { Request req requestService.create(UUID.randomUUID()); requestService.create(req.getUuid()); requestService.create(req.getUuid()); LOGGER.info( Nb of requests : {}, requestRepository.count()); // 1, processRequest is idempotent req requestService.start(req.getUuid()); try { req requestService.start(req.getUuid()); } catch (InvalidNextStateException ex) { LOGGER.error(Cannot start request twice!); } req requestService.complete(req.getUuid()); LOGGER.info(Request: {}, req); }; } }对应的程序输出来自原文档19:01:54.382 INFO [main] com.iluwatar.idempotentconsumer.App : Nb of requests : 1 19:01:54.395 ERROR [main] com.iluwatar.idempotentconsumer.App : Cannot start request twice! 19:01:54.399 INFO [main] com.iluwatar.idempotentconsumer.App : Request: Request(uuid2d5521ef-6b6b-4003-9ade-81e381fe9a63, statusCOMPLETED)三段输出分别验证了连续三次create后仓库记录数仍为 1幂等命中对已 STARTED 的记录再次start抛异常并被捕获重复被拒绝随后complete正常推进到 COMPLETED合法迁移放行。测试验证幂等语义的自动化保障仓库为上述行为提供了完整的测试覆盖可作为你接入该模式时的验收参照RequestServiceTests.java 使用 Mockito 模拟仓库逐一验证记录不存在时create恰好save一次记录已存在时create不再saveverify(requestRepository, times(0)).save(any())start/complete在记录缺失或状态非法时抛对应异常且不落库。RequestStateMachineTests.java 覆盖状态机全部路径PENDING→STARTED、STARTED→COMPLETED合法任何到PENDING的迁移、COMPLETED 出发的任意迁移均抛InvalidNextStateException。这些测试把「幂等」从口头承诺固化为可回归的契约是模式落地中最值得借鉴的部分。何时使用幂等消费者模式消息可能因网络抖动或生产者重试而多次到达时微服务必须在重复消息面前仍保证状态变更一致时容错的事件驱动通信对系统可靠性至关重要时水平扩展要求消费者具备无状态、可重复执行能力时。实际应用场景支付处理系统需要消化重复的扣款事件避免重复扣费电商订单服务需要处理重复的购买请求避免生成重复订单通知服务对失败的消息投递进行重试时避免重复发送通知分布式事务系统重复事件是常态需要靠幂等来兜底。收益与权衡收益Benefits阻止重复副作用防止重复扣款、重复建单在消息被重复或延迟投递时提升可靠性简化错误处理与重试逻辑——消费者可以放心重试因为重试本身是安全的。权衡Trade-offs需要精心设计消息处理记录的追踪方案维护幂等令牌idempotency token或幂等状态会带来额外开销可能需要额外的存储或数据库事务来支撑「先查后写」的判重逻辑。相关模式Outbox 模式发件箱模式使用专门的表或存储来可靠地发布事件并在源头做去重。幂等消费者模式从「消费端」兜底Outbox 从「生产端」兜底二者常组合使用以形成端到端的可靠投递。如何运行本模块该模块的构建配置见 microservices-idempotent-consumer/pom.xml它是 java-design-patterns 多模块仓库的子模块artifactId 为microservices-idempotent-consumer依赖 Spring Boot Data JPA、Hibernate并用 H2 内嵌数据库在运行时模拟持久层。通过 Maven Assembly 插件打包出的可执行 JAR 以com.iluwatar.idempotentconsumer.App为主类。在仓库根目录下可执行./mvnw -pl microservices-idempotent-consumer package java -jar microservices-idempotent-consumer/target/microservices-idempotent-consumer-*-jar-with-dependencies.jar运行后即可在控制台看到与上文一致的幂等演示输出。测试可执行./mvnw -pl microservices-idempotent-consumer test总结幂等消费者模式解决的不是「如何消灭重复消息」而是「如何让重复消息变得无害」。仓库示例给出的实现配方非常清晰UUID 唯一标识作为幂等键落到主键上→ 先查后写避免重复落库 → 状态机拒绝非法迁移。其中create用「查询命中即返回」实现幂等start/complete用「非法迁移即抛异常」拦截重复。这套思路可以平移到支付、电商订单、通知等任何面临重复投递的现实业务中是构建可靠微服务消息链路的必修课。赞分享示例工程教程【免费下载链接】java-design-patternsDesign patterns implemented in Java项目地址https://gitcode.com/GitHub_Trending/ja/java-design-patterns点击查看免费下载相关推荐微服务幂等性Idempotency实战指南从 HTTP 接口到消息消费的容错设计微服务幂等性Idempotency实战指南从 HTTP 接口到消息消费的容错设计 幂等性Idempotency是构建容错微服务系统的核心技术它确保同文档知识库Midscene.js 使用指南一句自然语言驱动的跨平台 UI 自动化测试Midscene.js 使用指南一句自然语言驱动的跨平台 UI 自动化测试 上周界面改版旧测试脚本里的选择器全失效了这次却一行代码都没改。Midscene消息队列流处理后端Apache RocketMQ消费者幂等处理分布式系统中的消息去重终极方案Apache RocketMQ消费者幂等处理分布式系统中的消息去重终极方案 引言为什么消息幂等至关重要 在分布式系统中消息队列Message Queu消息队列流处理后端上一篇Mac鼠标优化终极方案让你的普通鼠标秒变苹果级生产力工具下一篇Open3D点云配准初始化基于特征的粗配准方法创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表