Spring Integration DSL ScatterGather流程:多HTTP端点JSON聚合问题咨询
Spring Integration Scatter-Gather 实现多HTTP端点JSON结果聚合
我帮你把这段代码整理规范,并且补充完整的Scatter-Gather聚合逻辑,这样你就能顺利实现多HTTP端点返回JSON的聚合需求了:
@Bean public IntegrationFlow myFlow(IMyService myService, IMyOtherService myOtherService) { return f -> f .enrichHeaders(eh -> eh.headerExpression(Headers.PAYLOAD, "payload")) .handle(HeaderPrinter::headerPrinter) .enrichHeaders(httpRequestHeaderEnricher()) .scatterGather( // 配置Scatterer:分散调用多个HTTP端点 scatterer -> scatterer .applySequence(true) // 开启序列标识,帮助聚合器识别同一请求的消息组 // 第一个HTTP端点调用子流 .addFlow(subFlow -> subFlow .handle(Http.outboundGateway("http://first-endpoint/api/data") .httpMethod(HttpMethod.GET) .expectedResponseType(Map.class))) // 将JSON响应转为Map方便聚合 // 第二个HTTP端点调用子流,可按需添加更多端点 .addFlow(subFlow -> subFlow .handle(Http.outboundGateway("http://second-endpoint/api/info") .httpMethod(HttpMethod.GET) .expectedResponseType(Map.class))), // 配置Gatherer:聚合多个端点的返回结果 gatherer -> gatherer .outputProcessor(gatherResult -> { // 自定义聚合逻辑:把多个JSON结果合并为单个JSON对象 Map<String, Object> aggregatedJson = new HashMap<>(); gatherResult.getMessages().forEach(message -> { Map<String, Object> endpointPayload = (Map<String, Object>) message.getPayload(); // 直接合并所有字段(重复key会被后续端点结果覆盖) aggregatedJson.putAll(endpointPayload); // 可选:用自定义标识区分不同端点的结果 // String endpointKey = (String) message.getHeaders().get("endpoint-tag"); // aggregatedJson.put(endpointKey, endpointPayload); }); return aggregatedJson; }) .releaseLockBeforeSend(true) // 释放锁再发送聚合结果,避免资源占用 ) // 聚合完成后,传递给业务服务处理 .handle(myService::processAggregatedData); } // 辅助方法:配置HTTP请求头,确保端点接收JSON格式请求 private HeaderEnricherSpec httpRequestHeaderEnricher() { return HeaderEnricherSpec.headerEnricher() .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) .header(HttpHeaders.ACCEPT, MediaType.APPLICATION_JSON_VALUE); }
关键要点说明:
- Scatterer核心配置:
applySequence(true)是核心,它会给每个分散的消息添加序列标识,让聚合器能精准识别属于同一请求的消息组;每个addFlow对应一个HTTP outbound gateway,指定目标端点地址、请求方法,并把JSON响应转为Map类型,方便后续合并操作。 - Gatherer聚合逻辑:通过
outputProcessor自定义聚合规则,这里我们把每个端点返回的JSON(已转为Map)合并到一个大Map中,最终这个Map会被自动序列化为单个JSON对象返回。如果担心字段重复覆盖,可以给每个子流的消息添加自定义header(比如在子流中用enrichHeaders添加endpoint-tag),然后用这个header值作为聚合Map的key,把不同端点的结果分别存储。 - 后续处理:聚合完成后的JSON对象可以直接传递给你的业务服务(比如
myService)进行后续业务逻辑处理。
内容的提问来源于stack exchange,提问作者vdaras
相关产品推荐
相关产品推荐

