Spring Cloud Stream不使用@Output发送单条消息及关闭轮询的方法
解决方案
你当前遇到的重复发送问题,本质是使用了Supplier轮询生产模型:Spring Cloud Stream默认会每秒调用一次Supplier的返回值作为消息发送,适合定时生产场景,不符合你按需发送单条消息的需求。官方推荐使用StreamBridge实现自定义发送,完全符合你不使用废弃API、不依赖KafkaTemplate的要求。
1. 配置无需修改
你现有的yaml配置可以完全保留,不需要调整。
2. 调整代码实现
MessageService定位为普通业务服务类即可,不需要作为配置类,也不用定义Supplier类型的Bean:
import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Service; import lombok.RequiredArgsConstructor; @Service @RequiredArgsConstructor public class MessageService { // Spring Cloud Stream会自动装配StreamBridge实例 private final StreamBridge streamBridge; /** * 发送自定义消息 * @param payload 自定义消息负载 */ public void sendCustomMessage(Object payload) { // 第一个参数对应你配置中的绑定名:message-out-0 streamBridge.send("message-out-0", MessageBuilder.withPayload(payload) // 可按需添加消息头、Kafka消息key等配置 .setHeader("kafka_messageKey", "自定义消息key") .build()); } }
启动时发送单条消息的逻辑调整如下:
@Component @RequiredArgsConstructor public class ApplicationAnnouncer implements CommandLineRunner { private final MessageService messageService; @Override public void run(String... args) throws Exception { // 直接传入自定义负载即可,仅调用一次就只会发送一条 messageService.sendCustomMessage("Application started..."); } }
可选:如果坚持用Supplier实现单次发送
如果不想用StreamBridge,你可以通过修改轮询配置限制发送次数,在原有yaml中添加以下配置即可:
spring: cloud: stream: poller: max-messages-per-poll: 1 # 每次轮询只取1条 fixed-delay: -1 # 仅执行一次轮询,不再重复执行
该方案灵活性远低于StreamBridge,不推荐使用。
内容的提问来源于stack exchange,提问作者H3AR7B3A7
相关产品推荐
相关产品推荐

