使用Supplier模式向RabbitMQ生产消息的问题及优化咨询
问题分析与解决方案
需求说明
- 轮询消息并发布至RabbitMQ,涉及100个门店
- 每个门店每次轮询可获取5条以上消息
- 轮询周期为每分钟
现有代码
public class FetchMessages{ @Scheduled(fixedRateString = "60000") private void sendToClxRmq() { //Code to fetch the messages; //loop below line for all the 100 stores and messages polled by each of the store publishMessage.sendMessages(publishMessageDto); } } public class PublishMessage { PublishMessageDto publishMessageDto; @Bean public Supplier<Message<PublishMessageDto>> routeMessage() { return () -> { if (ObjectUtils.isNotEmpty(publishMessageDto)) { return MessageBuilder.withPayload(publishMessageDto) .setHeader("store", publishMessageDto.getStore()) .build(); } else { return null; } }; } public void sendMessages(final PublishMessageDto publishMessageDto) { this.publishMessageDto = publishMessageDto; } }
问题描述
Supplier Bean默认每秒触发一次get()方法,但sendToClxRmq()可能每秒为单个门店生成超过5条消息,导致单例的publishMessageDto被频繁覆盖,最终仅能发送最新的一条消息。
咨询问题
如何通过Supplier模式解决该问题?是否应使用StreamBridge或有更优方案?
解决方案
1. 基于Supplier模式的修复
核心思路是用线程安全的队列缓存待发送消息,避免单例变量被覆盖,让Supplier每次从队列中取出一条消息发送,保证消息不丢失、不覆盖。
修改后的PublishMessage类:
import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; public class PublishMessage { // 线程安全的队列,用于缓存待发送消息 private final BlockingQueue<PublishMessageDto> messageQueue = new LinkedBlockingQueue<>(); @Bean public Supplier<Message<PublishMessageDto>> routeMessage() { return () -> { // 从队列取出消息,队列为空时返回null,Supplier会等待下一次触发 PublishMessageDto dto = messageQueue.poll(); if (dto != null) { return MessageBuilder.withPayload(dto) .setHeader("store", dto.getStore()) .build(); } return null; }; } public void sendMessages(final PublishMessageDto publishMessageDto) { try { // 将消息加入队列,队列满时会阻塞(可根据业务需求设置超时时间) messageQueue.put(publishMessageDto); } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 此处可添加日志记录异常 } } }
说明:通过LinkedBlockingQueue缓存所有待发送消息,sendMessages()负责将消息入队,Supplier的get()方法每次从队列出队一条消息发送,彻底解决消息覆盖问题,同时支持流量削峰。
2. StreamBridge方案(更优选择)
如果使用Spring Cloud Stream,StreamBridge是更适配该场景的方案——它支持动态发送消息,无需依赖Supplier的定时触发逻辑,获取到消息后可直接发送,代码更简洁高效。
修改后的FetchMessages类:
import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.messaging.support.MessageBuilder; import org.springframework.scheduling.annotation.Scheduled; public class FetchMessages{ private final StreamBridge streamBridge; // 构造注入StreamBridge public FetchMessages(StreamBridge streamBridge) { this.streamBridge = streamBridge; } @Scheduled(fixedRateString = "60000") private void sendToClxRmq() { //Code to fetch the messages; //loop below line for all the 100 stores and messages polled by each of the store // 直接发送消息到RabbitMQ的指定绑定目标 streamBridge.send("rabbit-out-0", MessageBuilder.withPayload(publishMessageDto) .setHeader("store", publishMessageDto.getStore()) .build()); } }
配套配置文件示例(application.yaml):
spring: cloud: stream: rabbit: bindings: rabbit-out-0: producer: exchange: your-exchange-name routing-key-expression: headers['store'] bindings: rabbit-out-0: destination: your-exchange-name
说明:StreamBridge跳过了Supplier的定时触发环节,消息生成后立即发送,完全避免了消息覆盖问题,代码逻辑更直观,耦合度更低。
3. 方案选择建议
- 如果需要流量削峰或控制发送速率,选择Supplier+队列的模式,可通过队列容量和Supplier触发频率调节发送节奏;
- 如果是批量轮询后立即发送的场景,StreamBridge是更优选择,代码简洁、时效性更强,无需额外缓存逻辑。
内容的提问来源于stack exchange,提问作者Joe Lat
相关产品推荐
相关产品推荐

