You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 14:01:35