Spring Integration DSL中HTTP调用超时后如何返回自定义JSON错误信息?
解决方案
1. 捕获超时异常并转换为自定义JSON错误
推荐直接在HTTP出站调用的网关中添加异常处理Advice,将超时异常转换为自定义JSON并作为网关返回结果,便于后续流程判断终止。
代码实现
@Autowired private ObjectMapper objectMapper; private IntegrationFlow flow2() { return f -> f.handle( // 你的HTTP出站网关配置 Http.outboundGateway("http://your-target-url") .httpMethod(HttpMethod.POST) .expectedResponseType(String.class), // 配置超时异常处理Advice handler -> handler.advice(timeoutErrorAdvice()) ); } private Advice timeoutErrorAdvice() { ExpressionEvaluatingRequestHandlerAdvice advice = new ExpressionEvaluatingRequestHandlerAdvice(); // 开启异常捕获,避免异常向上抛出 advice.setTrapException(true); // 指定捕获超时异常类型 advice.setExceptionType(MessageTimeoutException.class); // 定义自定义错误响应内容 advice.setOnFailureExpressionString(""" T(java.util.Map).of( 'code', 'HTTP_CALL_TIMEOUT', 'message', 'HTTP出站调用超时,流程已终止', 'timeoutMs', 3000, 'timestamp', T(java.time.LocalDateTime).now().format(T(java.time.format.DateTimeFormatter).ISO_DATE_TIME) ) """); // 将错误对象转换为JSON字符串 advice.setFailureChannel(errorToJsonChannel()); return advice; } private MessageChannel errorToJsonChannel() { return MessageChannels.direct().get(); } @Bean public IntegrationFlow errorToJsonFlow() { return IntegrationFlows.from(errorToJsonChannel()) .handle((payload, headers) -> objectMapper.writeValueAsString(payload)) .get(); }
2. 终止主流程并返回错误
在scatter-gather步骤后添加路由逻辑,判断返回消息中是否包含错误信息。如果存在错误,则直接返回JSON错误,跳过后续聚合和业务处理:
@Bean public IntegrationFlow flow() { return flow -> flow.handle(validatorService, "validateRequest") .split() .channel(c -> c.executor(Executors.newCachedThreadPool())) .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(flow1()) .recipientFlow(flow2()) // flow2超时后会返回错误JSON .recipientFlow(flow3()), gatherer -> gatherer .releaseLockBeforeSend(true) .releaseStrategy(group -> group.size() == 2)) // 路由判断:是否为错误消息 .route( payload -> { if (payload instanceof MessageGroup messageGroup) { return messageGroup.getMessages().stream() .anyMatch(msg -> msg.getPayload().toString().contains("HTTP_CALL_TIMEOUT")); } return false; }, router -> router .channelMapping(true, "errorResponseChannel") .channelMapping(false, "normalProcessingChannel")) // 正常业务流程:聚合 + 后续处理 .from("normalProcessingChannel") .aggregate(someMethod1()) .to(someMethod2()) // 错误流程:直接返回JSON错误 .from("errorResponseChannel") .handle((errorPayload, headers) -> errorPayload) .get(); }
3. (可选)基于已配置errorChannel的处理方式
如果坚持使用现有errorChannel,可修改错误处理流程,将JSON错误发送回主流程的回复通道:
@Autowired private ObjectMapper objectMapper; @Autowired private MessageChannel mainReplyChannel; // 替换为主流程实际的回复通道 private IntegrationFlow returnError() { return IntegrationFlows.from("errorChannel") .handle((payload, headers) -> { ErrorMessage errorMessage = (ErrorMessage) payload; Throwable exception = errorMessage.getPayload(); Map<String, Object> errorMap = new HashMap<>(); errorMap.put("code", "HTTP_CALL_TIMEOUT"); errorMap.put("message", "HTTP调用超时,流程终止"); errorMap.put("timestamp", LocalDateTime.now().format(DateTimeFormatter.ISO_DATE_TIME)); if (exception instanceof MessageTimeoutException) { errorMap.put("timeoutMs", 3000); } return objectMapper.writeValueAsString(errorMap); }) .sendTo(mainReplyChannel) // 将错误发送回主流程回复通道 .get(); }
内容的提问来源于stack exchange,提问作者Somnath Mukherjee
相关产品推荐
相关产品推荐

