Spring Cloud Stream 3.2.4:批量发消息如何实现单条消费?
解决方案
你的问题根源在于当前生产者是把整个List<Message<String>>作为单个消息的payload发送到myOutDestination通道,因此消费者收到的是完整列表对象。要实现列表拆分为单条消息消费,有两种常用方案:
方式一:用Spring Integration Splitter组件拆分列表
通过Splitter可以自动将列表消息拆分为单个消息,再传递给消费方法,无需修改生产者逻辑。
调整消费者代码
修改MyTransformingProcessor,添加Splitter方法处理列表拆分:
public class MyTransformingProcessor { private final StreamBridge streamBridge; public MyTransformingProcessor(StreamBridge streamBridge) { this.streamBridge = streamBridge; } // 拆分消息列表,将列表中的每个Message单独发送到输出通道 @Splitter(inputChannel = "myOutDestination", outputChannel = "splitMessagesChannel") public List<Message<String>> splitMessageList(List<Message<String>> messageList) { return messageList; } // 处理拆分后的单条消息 @ServiceActivator(inputChannel = "splitMessagesChannel") public void accept(Message<String> source) { // 执行你的转换逻辑 streamBridge.send("mykafkatopic-out-0", source); } }
Spring Integration会自动识别@Splitter注解,将输入的列表拆分为单个元素,逐个发送到splitMessagesChannel,后续消费方法就能收到单条消息。
方式二:修改生产者发送逻辑(直接发送单条消息)
如果业务场景允许,也可以在生产者侧循环发送每条消息,避免发送列表:
@Scheduled(fixedDelay = 1000) public void get() { String payload = "my payload: " + random.nextInt(); String payload2 = "my payload 2: " + random.nextInt(); Message<String> msg1 = MessageBuilder.createMessage(payload, new MessageHeaders(Map.of())); Message<String> msg2 = MessageBuilder.createMessage(payload2, new MessageHeaders(Map.of())); // 分别发送每条消息 streamBridge.send("myOutDestination", msg1); streamBridge.send("myOutDestination", msg2); }
这种方式下,消费者无需额外配置,直接就能收到单条消息。
内容的提问来源于stack exchange,提问作者agienka
相关产品推荐
相关产品推荐

