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

如何实现Spring Integration通用流程处理多类型请求

Spring Integration 通用集成流程改造方案

改造核心思路是抽象公共层、剥离具体类型依赖,从接口定义、服务层、流程配置三层逐层消除硬编码类型强转,实现流程复用。


步骤1:定义公共请求顶层接口

把所有请求类型(LionRequest/CatRequest等)共有的属性访问、业务方法抽成顶层接口,所有具体请求类实现该接口,从数据模型层统一规范:

// 公共请求顶层接口
public interface BaseRequest {
    // 统一声明原有代码中被流程调用的公共方法
    RequestDetail getDetails();
}

// 具体请求类实现接口,原有业务逻辑保留
public class LionRequest implements BaseRequest {
    private RequestDetail details;
    @Override
    public RequestDetail getDetails() { return details; }
    // 其余LionRequest自有属性、方法
}

public class CatRequest implements BaseRequest {
    private RequestDetail details;
    @Override
    public RequestDetail getDetails() { return details; }
    // 其余CatRequest自有属性、方法
}

步骤2:服务层抽象剥离具体类型依赖

将原有LionsServiceImpl中硬编码依赖LionRequest的方法,全部调整为依赖BaseRequest接口,不同请求的差异化逻辑通过策略模式路由,不要在流程层做类型判断:

// 抽象通用服务接口
public interface CommonRequestService {
    void validateRequest(BaseRequest request);
    Object saveRequest(BaseRequest request, String dbId, Integer bizId, String sourceSystemCode) throws JsonProcessingException;
    Object getData(BaseRequest request, SourceSystem sourceSystem);
    Object prepareCDRequest(BaseRequest request);
}

如果LionRequest和CatRequest的处理逻辑存在差异,在服务实现类内部维护「请求类型-处理实现」的映射即可,流程层无需感知具体类型。

步骤3:改造集成流程配置,消除强转

修复原有代码的线程安全问题

原配置类中把dbId定义为类成员变量,会导致所有请求共用同一个ID,存在严重线程安全问题,改造时将dbId的生成逻辑放到单条请求的处理链路中,通过消息头传递给下游。

改造后的通用流程配置如下:

@Configuration
public class CommonIntegrationConfiguration {
    @Autowired
    private CommonRequestService commonRequestService;
    private final SequenceGenerator sequenceGenerator = new SequenceGenerator();

    // 通用主流程
    @Bean
    public IntegrationFlow commonProcessFlow() {
        return flow -> flow
            .<BaseRequest, MessageHeaders>handle((payload, headers) -> {
                commonRequestService.validateRequest(payload);
                // 生成当前请求专属dbId,放入消息头供下游使用
                headers.put("reqDbId", String.valueOf(sequenceGenerator.nextId()));
                return payload;
            })
            .split()
            .channel(c -> c.executor(Executors.newCachedThreadPool()))
            // 转换为公共接口类型,不再绑定具体请求类
            .convert(BaseRequest.class)
            .scatterGather(
                scatterer -> scatterer
                    .applySequence(true)
                    .recipientFlow(saveFlow())
                    .recipientFlow(queryFlow())
                    .recipientFlow(prepareFlow()),
                gatherer -> gatherer.releaseLockBeforeSend(true)
            )
            .log()
            .aggregate(a -> a.outputProcessor(MessageGroup::getMessages))
            .channel("output-flow");
    }

    // 原flow1:持久化请求子流程
    @Bean
    public IntegrationFlow saveFlow() {
        return flowDef -> flowDef
            .channel(c -> c.executor(Executors.newCachedThreadPool()))
            .<BaseRequest, MessageHeaders>handle((payload, headers) -> {
                try {
                    String dbId = headers.get("reqDbId", String.class);
                    SourceSystem source = headers.get("sourceSystem", SourceSystem.class);
                    Integer bizId = payload.getDetails().getId();
                    return commonRequestService.saveRequest(payload, dbId, bizId, source.getSourceSystemCode());
                } catch (JsonProcessingException e) {
                    return e.getMessage();
                }
            })
            .nullChannel();
    }

    // 原flow2:查询数据子流程
    @Bean
    public IntegrationFlow queryFlow() {
        return flowDef -> flowDef
            .channel(c -> c.executor(Executors.newCachedThreadPool()))
            .<Message>handle(msg -> {
                BaseRequest payload = (BaseRequest) msg.getPayload();
                SourceSystem source = msg.getHeaders().get("sourceSystem", SourceSystem.class);
                return commonRequestService.getData(payload, source);
            })
            .log();
    }

    // 原flow3:生成CD请求子流程
    @Bean
    public IntegrationFlow prepareFlow() {
        return flowDef -> flowDef
            .channel(c -> c.executor(Executors.newCachedThreadPool()))
            .<BaseRequest>handle(commonRequestService::prepareCDRequest);
    }

    // 输出通道配置
    @Bean
    public MessageChannel replyChannel() {
        return MessageChannels.executor("output-flow", outputExecutor()).get();
    }

    @Bean
    public ThreadPoolTaskExecutor outputExecutor() {
        ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
        pool.setCorePoolSize(4);
        pool.setMaxPoolSize(4);
        pool.initialize();
        return pool;
    }
}

步骤4:改造网关层支持通用请求

将原有绑定LionRequest的网关方法,改为支持所有BaseRequest实现类的通用方法:

@MessagingGateway
public interface CommonRequestGateway {
    @Gateway(requestChannel = "commonProcessFlow.input")
    <T extends BaseRequest> void processRequest(
        @Payload T request,
        @Header("sourceSystem") SourceSystem sourceSystem
    );
}

扩展优化点

  • 如果不同类型请求需要走的子流程组合不同,不要硬编码scatterGather的收件人列表,替换为路由器route(),根据请求类型、请求头动态决定要投递的子流程。
  • 不要在流程的handle方法中写instanceof类型判断,所有差异化业务逻辑下沉到服务层通过策略模式实现,保证流程层的通用性。
  • 线程池不要每次创建子通道都用Executors.newCachedThreadPool(),统一用Spring管理的线程池Bean,避免无限制创建线程导致资源耗尽。

内容的提问来源于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.26 18:57:23