Move to reactor netty httpclient from WebClient

This commit is contained in:
Spencer Gibb
2017-01-16 22:05:23 -07:00
parent fd6a038bd8
commit ba68a74dfa
5 changed files with 100 additions and 69 deletions

View File

@@ -25,7 +25,7 @@ import org.springframework.cloud.gateway.filter.route.SetResponseHeaderRouteFilt
import org.springframework.cloud.gateway.filter.route.SetStatusRouteFilter;
import org.springframework.cloud.gateway.handler.FilteringWebHandler;
import org.springframework.cloud.gateway.handler.RoutePredicateHandlerMapping;
import org.springframework.cloud.gateway.handler.WebClientRoutingWebHandler;
import org.springframework.cloud.gateway.handler.RoutingWebHandler;
import org.springframework.cloud.gateway.handler.predicate.CookieRoutePredicate;
import org.springframework.cloud.gateway.handler.predicate.HeaderRoutePredicate;
import org.springframework.cloud.gateway.handler.predicate.HostRoutePredicate;
@@ -35,8 +35,8 @@ import org.springframework.cloud.gateway.handler.predicate.RoutePredicate;
import org.springframework.cloud.gateway.handler.predicate.UrlRoutePredicate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.client.reactive.ReactorClientHttpConnector;
import org.springframework.web.reactive.function.client.WebClient;
import reactor.ipc.netty.http.client.HttpClient;
/**
* @author Spencer Gibb
@@ -47,8 +47,11 @@ public class GatewayAutoConfiguration {
@Bean
@ConditionalOnMissingBean
public WebClient webClient() {
return WebClient.builder(new ReactorClientHttpConnector(opts -> opts.disablePool())).build();
public HttpClient httpClient() {
return HttpClient.create(opts -> {
//opts.poolResources(PoolResources.elastic("proxy"));
opts.disablePool();
});
}
@Bean
@@ -63,21 +66,21 @@ public class GatewayAutoConfiguration {
}
@Bean
public WebClientRoutingWebHandler webClientRoutingWebHandler(WebClient webClient) {
return new WebClientRoutingWebHandler(webClient);
public RoutingWebHandler routingWebHandler(HttpClient httpClient) {
return new RoutingWebHandler(httpClient);
}
@Bean
public FilteringWebHandler filteringWebHandler(WebClientRoutingWebHandler webHandler,
List<GlobalFilter> globalFilters,
Map<String, RouteFilter> routeFilters) {
public FilteringWebHandler filteringWebHandler(RoutingWebHandler webHandler,
List<GlobalFilter> globalFilters,
Map<String, RouteFilter> routeFilters) {
return new FilteringWebHandler(webHandler, globalFilters, routeFilters);
}
@Bean
public RoutePredicateHandlerMapping routePredicateHandlerMapping(FilteringWebHandler webHandler,
Map<String, RoutePredicate> predicates,
RouteReader routeReader) {
Map<String, RoutePredicate> predicates,
RouteReader routeReader) {
return new RoutePredicateHandlerMapping(webHandler, predicates, routeReader);
}
@@ -188,3 +191,4 @@ public class GatewayAutoConfiguration {
}
}

View File

@@ -11,6 +11,8 @@ import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebFilterChain;
import org.springframework.web.util.UriComponentsBuilder;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_ROUTE_ATTR;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.getAttribute;
import reactor.core.publisher.Mono;
@@ -30,7 +32,7 @@ public class RouteToRequestUrlFilter implements GlobalFilter, Ordered {
@Override
public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
Route route = getAttribute(exchange, ServerWebExchangeUtils.GATEWAY_ROUTE_ATTR, Route.class);
Route route = getAttribute(exchange, GATEWAY_ROUTE_ATTR, Route.class);
if (route == null) {
return chain.filter(exchange);
}
@@ -39,7 +41,7 @@ public class RouteToRequestUrlFilter implements GlobalFilter, Ordered {
.uri(route.getUri())
.build(true)
.toUri();
exchange.getAttributes().put(ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR, requestUrl);
exchange.getAttributes().put(GATEWAY_REQUEST_URL_ATTR, requestUrl);
return chain.filter(exchange);
}

View File

@@ -3,9 +3,9 @@ package org.springframework.cloud.gateway.filter;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.core.Ordered;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.NettyDataBuffer;
import org.springframework.core.io.buffer.NettyDataBufferFactory;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.web.reactive.function.client.ClientResponse;
import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebFilterChain;
@@ -15,6 +15,7 @@ import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.g
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.ipc.netty.http.client.HttpClientResponse;
/**
* @author Spencer Gibb
@@ -35,7 +36,7 @@ public class WriteResponseFilter implements GlobalFilter, Ordered {
// NOTICE: nothing in "pre" filter stage as CLIENT_RESPONSE_ATTR is not added
// until the WebHandler is run
return chain.filter(exchange).then(() -> {
ClientResponse clientResponse = getAttribute(exchange, CLIENT_RESPONSE_ATTR, ClientResponse.class);
HttpClientResponse clientResponse = getAttribute(exchange, CLIENT_RESPONSE_ATTR, HttpClientResponse.class);
if (clientResponse == null) {
return Mono.empty();
}
@@ -48,7 +49,13 @@ public class WriteResponseFilter implements GlobalFilter, Ordered {
return Mono.empty();
});
Flux<DataBuffer> body = clientResponse.body((inputMessage, context) -> inputMessage.getBody());
NettyDataBufferFactory factory = (NettyDataBufferFactory) response.bufferFactory();
//TODO: what if it's not netty
final Flux<NettyDataBuffer> body = clientResponse.receive()
.retain() //TODO: needed?
.map(factory::wrap);
return response.writeWith(body);
});
}

View File

@@ -0,0 +1,69 @@
package org.springframework.cloud.gateway.handler;
import java.net.URI;
import java.util.Optional;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpStatus;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebHandler;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.CLIENT_RESPONSE_ATTR;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR;
import io.netty.buffer.Unpooled;
import io.netty.handler.codec.http.DefaultHttpHeaders;
import io.netty.handler.codec.http.HttpMethod;
import reactor.core.publisher.Mono;
import reactor.ipc.netty.NettyPipeline;
import reactor.ipc.netty.http.client.HttpClient;
/**
* @author Spencer Gibb
*/
public class RoutingWebHandler implements WebHandler {
private final HttpClient httpClient;
public RoutingWebHandler(HttpClient httpClient) {
this.httpClient = httpClient;
}
@Override
public Mono<Void> handle(ServerWebExchange exchange) {
Optional<URI> requestUrl = exchange.getAttribute(GATEWAY_REQUEST_URL_ATTR);
ServerHttpRequest request = exchange.getRequest();
final HttpMethod method = HttpMethod.valueOf(request.getMethod().toString());
final String url = requestUrl.get().toString();
final DefaultHttpHeaders httpHeaders = new DefaultHttpHeaders();
request.getHeaders().forEach(httpHeaders::set);
return this.httpClient.request(method, url, req ->
req.options(NettyPipeline.SendOptions::flushOnEach)
.headers(httpHeaders)
//.sendHeaders()
.send(request.getBody()
.map(DataBuffer::asByteBuffer)
.map(Unpooled::wrappedBuffer)))
.then(res -> {
// Defer committing the response until all route filters have run
// Put client response as ServerWebExchange attribute and write response later WriteResponseFilter
exchange.getAttributes().put(CLIENT_RESPONSE_ATTR, res);
ServerHttpResponse response = exchange.getResponse();
// put headers and status so filters can modify the response
final HttpHeaders headers = new HttpHeaders();
res.responseHeaders().forEach(entry -> headers.add(entry.getKey(), entry.getValue()));
response.getHeaders().putAll(headers);
response.setStatusCode(HttpStatus.valueOf(res.status().code()));
return Mono.empty();
});
}
}

View File

@@ -1,51 +0,0 @@
package org.springframework.cloud.gateway.handler;
import java.net.URI;
import java.util.Optional;
import org.springframework.http.server.reactive.ServerHttpRequest;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.web.reactive.function.client.ClientRequest;
import org.springframework.web.reactive.function.client.WebClient;
import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebHandler;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.CLIENT_RESPONSE_ATTR;
import static org.springframework.cloud.gateway.support.ServerWebExchangeUtils.GATEWAY_REQUEST_URL_ATTR;
import reactor.core.publisher.Mono;
/**
* @author Spencer Gibb
*/
public class WebClientRoutingWebHandler implements WebHandler {
private final WebClient webClient;
public WebClientRoutingWebHandler(WebClient webClient) {
this.webClient = webClient;
}
@Override
public Mono<Void> handle(ServerWebExchange exchange) {
Optional<URI> requestUrl = exchange.getAttribute(GATEWAY_REQUEST_URL_ATTR);
ServerHttpRequest request = exchange.getRequest();
ClientRequest<Void> clientRequest = ClientRequest
.method(request.getMethod(), requestUrl.get())
.headers(request.getHeaders())
.body((r, context) -> r.writeWith(request.getBody()));
return this.webClient.exchange(clientRequest).then(clientResponse -> {
// Defer committing the response until all route filters have run
// Put client response as ServerWebExchange attribute and write response later WriteResponseFilter
exchange.getAttributes().put(CLIENT_RESPONSE_ATTR, clientResponse);
ServerHttpResponse response = exchange.getResponse();
// put headers and status so filters can modify the response
response.getHeaders().putAll(clientResponse.headers().asHttpHeaders());
response.setStatusCode(clientResponse.statusCode());
return Mono.<Void>empty();
});
}
}