【发布时间】:2017-02-06 01:33:08
【问题描述】:
我们有一个要求,客户端调用我们的 spring 集成 http 入站网关之一,给 API 的输入是 .csv 格式,一旦 请求被验证并发现正确的立即响应应该以状态 200 OK 发送。如果发生错误,则发送相应的错误消息。 我们使用直接和执行器通道的组合进行异步处理。这在使用 Spring Boot 父版本 1.2.5 时可以正常工作,但在升级到 1.4.0 版本时会失败。我们总是收到 500 Internal server error,原因是从日志中发现的原因是 MessageTimeoutException。
我们使用基于 java 的配置,配置如下。
pom.xml
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>1.4.0.RELEASE</version>
</parent>
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-http</artifactId>
<version>4.3.1.RELEASE</version>
</dependency>
@Configuration
public class ApplicationIntegrationConfig {
@Bean
public HttpRequestHandlingMessagingGateway httpMessageGateway(){
HttpRequestHandlingMessagingGateway gateway
= new HttpRequestHandlingMessagingGateway(Boolean.TRUE);
RequestMapping requestMapping = new RequestMapping();
requestMapping.setMethods(HttpMethod.POST);
requestMapping.setPathPatterns("/org/{orgId}/users");
requestMapping.setHeaders("Content-Type=text/csv");
gateway.setRequestMapping(requestMapping);
gateway.setRequestChannel(onBoardUserRequestChannel());
Map<String, Expression> customHeaderExpressions = new HashMap<>();
customHeaderExpressions.put("orgId", new SpelExpressionParser().
parseExpression("#pathVariables.orgId"));
gateway.setHeaderExpressions(customHeaderExpressions);
gateway.setErrorChannel(errorChannel());
gateway.setReplyTimeout(0);
return gateway;
}
@Bean
public MessageChannel processUserRequestChannel() {
DirectChannel channel =new DirectChannel();
channel.addInterceptor(new AuthenticationInterceptor());
return channel;
}
@Bean
public MessageChannel routeChannel() {
return new ExecutorChannel(Executors.newCachedThreadPool());
}
@Bean
public MessageChannel addUserChannel() {
return new ExecutorChannel(Executors.newCachedThreadPool());
}
@Bean
public MessageChannel removeUserChannel() {
return new ExecutorChannel(Executors.newCachedThreadPool());
}
@Bean
public MessageChannel errorChannel() {
return new DirectChannel();
}
}
分离器
@MessageEndpoint
public class PartnerUserOnBoardSplitter {
@Splitter(inputChannel= "processUserRequestChannel", outputChannel="routeChannel")
public List<UserDTO> split(Message message) throws ApplicationException {
List<UserDTO> userList = null;
try {
userList = validateAndCreateDTO(message);
}
} catch(Exception ex) {
throw new ApplicationException("<Message>");
}
return userList;
}
}
路由器
@MessageEndpoint
public class CustomRouter {
@Router(inputChannel="routeChannel")
public String resolveRoute(UserDTO dto) {
return (Operation.ADD.equals(dto.getOperation())) ? "addUserChannel" : "removeUserChannel";
}
}
public class ServiceActivator{
@ServiceActivator(inputChannel = "addUserChannel")
public addUser(UserDto dto){
//process add
}
@ServiceActivator(inputChannel = "removeUserChannel")
public removeUser(UserDto dto){
//process remove
}
}
【问题讨论】: