如何让自定义Spring Integration MessageSource批量生成单条Widget消息?
当前环境:Java 11 + Spring Integration 5.x
我创建了自定义的MessageSource<Widget>并集成到流中:
@Bean public IntegrationFlow myFlow(WidgetMessageSource widgetMessageSource) { return IntegrationFlows.from(widgetMessageSource) .get(); }
WidgetMessageSource的实现逻辑是批量查询远程系统获取Widget实例,但MessageSource的receive()方法要求返回单个Message<Widget>,而我需要将每个Widget都封装成独立的消息发送:
public class WidgetMessageSource implements MessageSource<Widget> { @Override public Message<Widget> receive() { List<Widget> widgets = fetchStagedWidgets(); // 如何生成List<Message<Widget>>而非单个Message<Widget>??? } private List<Widget> fetchStagedWidgets() { // 连接远程系统并拉取0个或多个widget return where_the_magic_happens(); // TODO } }
核心需求:
- 若
fetchStagedWidgets返回空列表,不发送任何消息 - 若返回多个
Widget,每个对应一条独立的Message<Widget>
不想修改fetchStagedWidgets的底层批量逻辑,询问Spring Integration的替代解决方案。
解决方案1:修改MessageSource返回列表,配合Splitter拆分
将WidgetMessageSource的泛型改为List<Widget>,返回批量结果后,用Spring Integration的split()组件自动拆分列表为单个消息。
修改MessageSource实现
public class WidgetMessageSource implements MessageSource<List<Widget>> { @Override public Message<List<Widget>> receive() { List<Widget> widgets = fetchStagedWidgets(); if (widgets.isEmpty()) { return null; // 无数据时返回null,Spring Integration会跳过本次轮询 } return MessageBuilder.withPayload(widgets).build(); } private List<Widget> fetchStagedWidgets() { // 原有批量逻辑保持不变 return where_the_magic_happens(); } }
更新集成流配置
@Bean public IntegrationFlow myFlow(WidgetMessageSource widgetMessageSource) { return IntegrationFlows.from(widgetMessageSource) .split() // 自动拆分List<Widget>为单个Widget消息 // 后续添加你的业务处理逻辑 .get(); }
说明:split()组件默认支持拆分集合、数组等可迭代对象,拆分后的每个元素会作为独立消息在流中传递;如果原消息的payload是空列表,split()不会产生任何输出,完全符合需求。
解决方案2:继承AbstractMessageSource,内部缓存批量结果
如果不想修改MessageSource的泛型,可以继承Spring Integration提供的AbstractMessageSource<Widget>,内部用队列缓存批量获取的Widget,每次receive()返回一个元素,直到队列空后再重新拉取。
public class WidgetMessageSource extends AbstractMessageSource<Widget> { private final Queue<Widget> widgetQueue = new ConcurrentLinkedQueue<>(); @Override protected Object doReceive() { if (widgetQueue.isEmpty()) { List<Widget> widgets = fetchStagedWidgets(); widgetQueue.addAll(widgets); } return widgetQueue.poll(); // 每次返回一个元素,队列空时返回null } private List<Widget> fetchStagedWidgets() { // 原有批量逻辑保持不变 return where_the_magic_happens(); } @Override public String getComponentType() { return "widget-message-source"; } }
说明:
AbstractMessageSource简化了MessageSource的实现,无需手动处理消息构建逻辑- 每次轮询时,先检查队列是否为空,空则批量拉取并填充队列;之后每次调用
receive()从队列取一个元素,直到队列耗尽 - 原有集成流配置无需修改,
myFlow代码可以保持原样
解决方案3:手动发送消息到通道
如果需要更灵活的发送控制,可以让WidgetMessageSource注入MessageChannel,批量获取后手动将每个Widget封装成消息发送到通道。
修改MessageSource实现
public class WidgetMessageSource implements MessageSource<Void> { private final MessageChannel outputChannel; // 通过构造器注入通道 public WidgetMessageSource(MessageChannel outputChannel) { this.outputChannel = outputChannel; } @Override public Message<Void> receive() { List<Widget> widgets = fetchStagedWidgets(); widgets.forEach(widget -> { Message<Widget> message = MessageBuilder.withPayload(widget).build(); outputChannel.send(message); }); return null; // 无需返回消息,已手动完成发送 } private List<Widget> fetchStagedWidgets() { // 原有批量逻辑保持不变 return where_the_magic_happens(); } }
更新集成流配置
@Bean public MessageChannel widgetOutputChannel() { return new DirectChannel(); } @Bean public IntegrationFlow myFlow(MessageChannel widgetOutputChannel) { WidgetMessageSource messageSource = new WidgetMessageSource(widgetOutputChannel); return IntegrationFlows.from(messageSource) .channel(widgetOutputChannel) // 后续添加你的业务处理逻辑 .get(); }
说明:这种方式适合需要自定义发送规则的场景,比如添加特定消息头、处理发送失败等,但需要注意线程安全和消息可靠性。
内容的提问来源于stack exchange,提问作者hotmeatballsoup

