You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 17:25:27