消息队列(Kafka/RocketMQ)在削峰填谷与异步解耦中的深度实践

发布时间:2026/7/23 20:08:22
消息队列(Kafka/RocketMQ)在削峰填谷与异步解耦中的深度实践 消息队列Kafka/RocketMQ在削峰填谷与异步解耦中的深度实践又见面了我是高佣返利省赚客APP研发者微赚在省赚客APP的架构演进中消息队列MQ扮演着“中枢神经”的关键角色。面对双11、618等大促期间瞬间爆发的订单洪峰同步调用模式往往导致数据库连接池耗尽、接口响应超时甚至服务雪崩。我们引入RocketMQ与Kafka构建了一套高吞吐、低延迟的异步处理体系不仅实现了流量的削峰填谷更彻底解耦了订单创建、佣金计算、积分发放、风控审核等复杂业务链路。本文将深入代码层面解析如何利用MQ保障数据最终一致性与系统高可用。核心场景订单创建后的异步链路解耦在传统的单体或简单微服务架构中用户下单后系统需同步完成扣减库存、冻结资金、调用联盟API验单、计算返利、发送通知、记录日志。这一串行过程耗时可能高达2秒以上严重影响用户体验。我们将非核心逻辑剥离通过MQ异步执行。主流程仅负责落库和发消息响应时间压缩至200ms以内。packagejuwatech.cn.provinceearn.order.service.impl;importjuwatech.cn.provinceearn.order.entity.Order;importjuwatech.cn.provinceearn.order.repository.OrderRepository;importjuwatech.cn.provinceearn.mq.producer.OrderEventProducer;importjuwatech.cn.provinceearn.order.dto.OrderCreateRequest;importjuwatech.cn.provinceearn.order.enums.OrderStatus;importlombok.RequiredArgsConstructor;importlombok.extern.slf4j.Slf4j;importorg.springframework.stereotype.Service;importorg.springframework.transaction.annotation.Transactional;importjava.util.Date;/** * 订单创建服务 * 利用MQ实现核心链路与辅助链路的异步解耦 */Slf4jServiceRequiredArgsConstructorpublicclassOrderServiceImpl{privatefinalOrderRepositoryorderRepository;privatefinalOrderEventProducerorderEventProducer;/** * 创建订单 * 事务内仅做最核心的落库操作其余全部异步化 */Transactional(rollbackForException.class)publicOrdercreateOrder(OrderCreateRequestrequest){// 1. 构建订单实体OrderordernewOrder();order.setUserId(request.getUserId());order.setProductId(request.getProductId());order.setAmount(request.getAmount());order.setStatus(OrderStatus.CREATED);order.setCreateTime(newDate());// 2. 持久化订单 (本地事务)orderRepository.save(order);log.info(Order saved locally, orderId: {},order.getId());// 3. 发送消息到MQ// 注意此处利用Spring的事务监听器或手动确保消息发送在事务提交后// 为防止本地事务回滚但消息已发出的情况最佳实践是使用事务型MQ或本地消息表orderEventProducer.sendOrderCreatedEvent(order);returnorder;}}可靠投递本地消息表与事务型MQ实现分布式系统中最棘手的问题是“本地事务成功但消息发送失败”。我们采用了“本地消息表 定时任务补偿”的最终一致性方案确保消息必达。同时对于关键金融链路直接利用RocketMQ的事务消息机制。packagejuwatech.cn.provinceearn.mq.producer;importjuwatech.cn.provinceearn.mq.entity.LocalMessage;importjuwatech.cn.provinceearn.mq.repository.LocalMessageRepository;importjuwatech.cn.provinceearn.order.entity.Order;importcom.fasterxml.jackson.databind.ObjectMapper;importlombok.RequiredArgsConstructor;importlombok.extern.slf4j.Slf4j;importorg.apache.rocketmq.client.producer.SendResult;importorg.apache.rocketmq.spring.core.RocketMQTemplate;importorg.springframework.messaging.support.MessageBuilder;importorg.springframework.stereotype.Component;importorg.springframework.transaction.support.TransactionSynchronizationAdapter;importorg.springframework.transaction.support.TransactionSynchronizationManager;importjava.util.UUID;/** * 订单事件生产者 * 实现基于本地消息表的可靠投递机制 */Slf4jComponentRequiredArgsConstructorpublicclassOrderEventProducer{privatefinalRocketMQTemplaterocketMQTemplate;privatefinalLocalMessageRepositorymessageRepository;privatefinalObjectMapperobjectMapper;privatestaticfinalStringTOPIC_ORDER_CREATEDTOPIC_ORDER_CREATED;privatestaticfinalStringTAGS_DEFAULTTAG_CREATE;/** * 发送订单创建事件 * 策略先写本地消息表事务提交后由监听器异步发送 */publicvoidsendOrderCreatedEvent(Orderorder){try{StringmessageIdUUID.randomUUID().toString();StringbodyobjectMapper.writeValueAsString(order);// 1. 构建本地消息记录LocalMessagelocalMessagenewLocalMessage();localMessage.setMessageId(messageId);localMessage.setTopic(TOPIC_ORDER_CREATED);localMessage.setTags(TAGS_DEFAULT);localMessage.setBody(body);localMessage.setStatus(TO_SEND);// 待发送localMessage.setRetryCount(0);localMessage.setCreateTime(newjava.util.Date());// 2. 保存本地消息 (与订单保存在同一事务中)messageRepository.save(localMessage);// 3. 注册事务同步回调事务提交后才真正发送MQ消息TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronizationAdapter(){OverridepublicvoidafterCommit(){sendMessageAsync(localMessage);}});log.info(Local message recorded for orderId: {},order.getId());}catch(Exceptione){log.error(Failed to record local message,e);thrownewRuntimeException(Order creation failed due to message logging error,e);}}/** * 异步发送消息到RocketMQ */privatevoidsendMessageAsync(LocalMessagemsg){try{SendResultresultrocketMQTemplate.syncSend(msg.getTopic():msg.getTags(),MessageBuilder.withPayload(msg.getBody()).build());if(result.getSendStatus().name().equals(SEND_OK)){// 更新本地消息状态为已发送messageRepository.updateStatus(msg.getMessageId(),SENT);log.info(Message sent successfully: {},msg.getMessageId());}else{log.warn(Message send failed, status: {},result.getSendStatus());// 触发重试逻辑或报警}}catch(Exceptione){log.error(Exception occurred while sending message,e);// 依赖定时任务扫描TO_SEND状态的消息进行补偿重发}}}消费端幂等设计与流量削峰消息到达消费者后必须解决重复消费问题网络抖动导致ACK丢失。我们采用“数据库唯一键约束 Redis去重表”的双重幂等机制。同时通过设置消费者的并发度和拉取策略控制处理速率将突发流量平滑为数据库可承受的稳定负载。packagejuwatech.cn.provinceearn.mq.consumer;importjuwatech.cn.provinceearn.mq.service.CommissionCalculationService;importjuwatech.cn.provinceearn.common.util.RedisDistributedLock;importlombok.RequiredArgsConstructor;importlombok.extern.slf4j.Slf4j;importorg.apache.rocketmq.spring.annotation.RocketMQMessageListener;importorg.apache.rocketmq.spring.core.RocketMQListener;importorg.springframework.data.redis.core.StringRedisTemplate;importorg.springframework.stereotype.Component;importjava.util.concurrent.TimeUnit;/** * 订单创建消息消费者 * 负责异步计算佣金与发放积分 * 实现幂等性与限流保护 */Slf4jComponentRocketMQMessageListener(topicTOPIC_ORDER_CREATED,consumerGroupCG_COMMISSION_CALC,consumeThreadMax20,// 控制并发度实现削峰consumeTimeout180// 消费超时时间)RequiredArgsConstructorpublicclassOrderCreatedConsumerimplementsRocketMQListenerString{privatefinalCommissionCalculationServicecalculationService;privatefinalStringRedisTemplateredisTemplate;privatefinalRedisDistributedLockdistributedLock;privatestaticfinalStringIDEMPOTENT_KEY_PREFIXidempotent:order:;OverridepublicvoidonMessage(StringmessageBody){log.info(Received message: {},messageBody);// 1. 解析消息获取OrderId (假设JSON格式)StringorderIdextractOrderId(messageBody);StringidempotentKeyIDEMPOTENT_KEY_PREFIXorderId;// 2. 幂等性校验利用Redis SETNX实现快速去重BooleanisNewredisTemplate.opsForValue().setIfAbsent(idempotentKey,PROCESSING,24,TimeUnit.HOURS);if(Boolean.FALSE.equals(isNew)){log.warn(Duplicate message detected, ignoring. OrderId: {},orderId);return;// 直接ACK不再处理}try{// 3. 执行业务逻辑 (佣金计算、积分入账)// 此处可加分布式锁防止极端并发下的数据竞争视业务复杂度而定calculationService.processCommission(messageBody);log.info(Commission calculated successfully for Order: {},orderId);}catch(Exceptione){log.error(Failed to process commission for Order: {},orderId,e);// 抛出异常让RocketMQ触发重试机制// 注意需配合最大重试次数配置避免死循环thrownewRuntimeException(e);}}privateStringextractOrderId(Stringbody){// 简化解析逻辑returnbody.split(\orderId\:\)[1].split(\)[0];}}结语通过引入消息队列省赚客APP成功将同步阻塞的单体架构转型为高弹性的异步事件驱动架构。本地消息表保障了数据的绝对可靠幂等设计消除了重复消费隐患而灵活的并发控制则完美实现了削峰填谷。这套机制支撑了平台在亿级流量下的稳定运行是构建高并发分布式系统的基石。本文著作权归 省赚客app 研发团队转载请注明出处