Spring Integration中如何在流程间传递对象且不修改原负载?
Spring Integration多流程场景下的对象传递方案
一、复合负载载体方案(推荐用于大对象)
核心思路是创建一个包含Object1和Object2的上下文DTO,将其作为消息负载传递,既保留两个对象的关联关系,又避免大对象存入消息头的风险。
1. 定义上下文DTO
public class ProcessContext { private Object1 obj1; private Object2 obj2; // 可添加聚合结果字段用于后续流程 private Object aggregatedResult; // 构造器、getter、setter public ProcessContext(Object1 obj1, Object2 obj2) { this.obj1 = obj1; this.obj2 = obj2; } }
2. 修改validateRequest方法
public ProcessContext validateRequest(LionRequest lionRequest) { lionValidationHelper.validateRequestAttributes(lionRequest); Object1 obj1 = someLogicTocreateTheObject1(); Object2 obj2 = someLogicTocreateTheObject2(); return new ProcessContext(obj1, obj2); }
3. 调整主流程与子流程
@Bean public IntegrationFlow flow() { return flow -> flow.handle(validatorService, "validateRequest") // 暂存上下文到消息头,提取Object1作为Scatter-Gather的处理负载 .enrichHeaders(h -> h.header("processContext", Message::getPayload)) .transform(ProcessContext::getObj1) .split() .channel(c -> c.executor(Executors.newCachedThreadPool())) .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(flow1()) .recipientFlow(flow2()) .recipientFlow(flow3()), gatherer -> gatherer .releaseLockBeforeSend(true) .releaseStrategy(group -> group.size() == 2)) .aggregate(lionService.someMethod()) // 聚合后将结果合并回上下文,恢复Object2的关联 .enrichPayload((payload, headers) -> { ProcessContext context = headers.get("processContext", ProcessContext.class); context.setAggregatedResult(payload); return context; }) // 后续流程可直接从ProcessContext中获取Object2操作 .gateway(someFlow()) .to(someFlow2()); } // 子流程无需修改,直接接收Object1处理 @Bean public IntegrationFlow flow1() { return flow -> flow.channel(c -> c.executor(Executors.newCachedThreadPool())) .enrichHeaders(h -> h.errorChannel("flow1ErrorChannel", true)) .handle(cdRequestService, "prepareCDRequestFromLionRequest"); }
二、消息存储方案(适合超大型对象)
如果Object2体积过大,不适合存入负载或消息头,可使用Spring Integration的MessageStore存储对象,通过消息ID作为标识传递。
1. 配置MessageStore
@Bean public MessageStore messageStore() { // 单机场景用SimpleMessageStore,分布式场景替换为RedisMessageStore等 return new SimpleMessageStore(); }
2. 修改validateRequest与流程
// 注入MessageStore,通过消息ID存储Object2 public Object1 validateRequest(LionRequest lionRequest, MessageHeaders headers) { lionValidationHelper.validateRequestAttributes(lionRequest); Object1 obj1 = someLogicTocreateTheObject1(); Object2 obj2 = someLogicTocreateTheObject2(); String messageId = headers.getId().toString(); messageStore.addMessage(MessageBuilder.withPayload(obj2).build()); return obj1; } @Bean public IntegrationFlow flow() { return flow -> flow.handle(validatorService, "validateRequest") .split() .channel(c -> c.executor(Executors.newCachedThreadPool())) .scatterGather(...) // 原有Scatter-Gather逻辑不变 .aggregate(lionService.someMethod()) // 从MessageStore中取出Object2,组合后传递给后续流程 .enrichPayload((payload, headers) -> { String messageId = headers.getId().toString(); Message<Object2> obj2Msg = messageStore.getMessage(messageId); Object2 obj2 = obj2Msg.getPayload(); return new ProcessContext(payload, obj2); }) .gateway(someFlow()) .to(someFlow2()); }
三、后续流程创建对象不影响原负载的方案
若需在validateRequest之后创建新对象并传递,同时保留原负载,可采用以下两种方式:
1. 使用Enricher添加关联对象
- 小对象存入消息头:
.flow(...) .handle(validatorService, "validateRequest") .enrichHeaders(h -> h.header("newObject", () -> createNewObject())) .scatterGather(...)
- 大对象存入复合负载:
.enrichPayload(payload -> new ProcessContext(payload, createNewObject()))
2. 自定义处理器包装消息
.handle((payload, headers) -> { Object1 obj1 = (Object1) payload; SomeNewObject newObj = createNewObject(); // 保留原负载,新对象存入消息头 return MessageBuilder.withPayload(obj1) .setHeader("newObject", newObj) .copyHeaders(headers) .build(); })
内容的提问来源于stack exchange,提问作者Somnath Mukherjee
相关产品推荐
相关产品推荐

