Spring Cloud Stream 3.1+函数式生产者如何接收REST控制器传入的参数
问题原因
- 你选择的
Supplier函数式接口是Spring Cloud Stream面向主动定时生成消息场景设计的,框架默认会每秒轮询调用一次Supplier的get()方法生成并发送消息,这就是你应用启动后无限自动发消息的根本原因,该类型本身不适合你这种收到REST请求才按需发消息的场景。 Supplier的逻辑是封装内部自发生成消息,设计上就不支持接收外部传入的动态参数,所以你没法把控制器收到的请求数据传入使用,属于选型错误。
解决方案
3.1+版本官方专门提供了StreamBridge组件用于按需发送消息,完全可以替代原来的MessageChannel发送方式,改动成本极低:
- 移除你之前写的
Supplier类型Bean,修改Producer类实现如下:
import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.messaging.Message; import org.springframework.messaging.MessageBuilder; import org.springframework.messaging.support.MessageHeaders; import org.springframework.util.MimeTypeUtils; import org.springframework.stereotype.Component; @Component public class Producer { private final StreamBridge streamBridge; // 构造注入StreamBridge,Spring会自动装配该实例 public Producer(StreamBridge streamBridge) { this.streamBridge = streamBridge; } public void produce(int messageId, Object message) { Message<Object> msg = MessageBuilder .withPayload(message) .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON) .setHeader("partitionKey", messageId) .build(); // 第一个参数为绑定名,和你yaml中配置的produce-out-0保持一致即可 streamBridge.send("produce-out-0", msg); } }
- 你现有的
application.yaml配置、REST控制器代码完全不需要修改,原有逻辑可以直接复用,只有接收到接口请求时才会触发消息发送,不会出现自动无限发消息的问题。
内容的提问来源于stack exchange,提问作者Vin
相关产品推荐
相关产品推荐

