如何实现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
相关产品推荐
相关产品推荐

