Support LB stats. (#447)
This commit is contained in:
committed by
GitHub
parent
c060e07115
commit
ecfba4b6e9
@@ -93,15 +93,15 @@ public class FeignBlockingLoadBalancerClient implements Client {
|
||||
if (LOG.isWarnEnabled()) {
|
||||
LOG.warn(message);
|
||||
}
|
||||
supportedLifecycleProcessors
|
||||
.forEach(lifecycle -> lifecycle.onComplete(new CompletionContext<ResponseData, ServiceInstance>(
|
||||
CompletionContext.Status.DISCARD, lbResponse)));
|
||||
supportedLifecycleProcessors.forEach(lifecycle -> lifecycle
|
||||
.onComplete(new CompletionContext<ResponseData, ServiceInstance, RequestDataContext>(
|
||||
CompletionContext.Status.DISCARD, lbRequest, lbResponse)));
|
||||
return Response.builder().request(request).status(HttpStatus.SERVICE_UNAVAILABLE.value())
|
||||
.body(message, StandardCharsets.UTF_8).build();
|
||||
}
|
||||
String reconstructedUrl = loadBalancerClient.reconstructURI(instance, originalUri).toString();
|
||||
Request newRequest = buildRequest(request, reconstructedUrl);
|
||||
return executeWithLoadBalancerLifecycleProcessing(delegate, options, newRequest, lbResponse,
|
||||
return executeWithLoadBalancerLifecycleProcessing(delegate, options, newRequest, lbRequest, lbResponse,
|
||||
supportedLifecycleProcessors);
|
||||
}
|
||||
|
||||
|
||||
@@ -48,21 +48,23 @@ final class LoadBalancerUtils {
|
||||
}
|
||||
|
||||
static Response executeWithLoadBalancerLifecycleProcessing(Client feignClient, Request.Options options,
|
||||
Request feignRequest, org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse,
|
||||
Request feignRequest, org.springframework.cloud.client.loadbalancer.Request lbRequest,
|
||||
org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse,
|
||||
Set<LoadBalancerLifecycle> supportedLifecycleProcessors, boolean loadBalanced) throws IOException {
|
||||
supportedLifecycleProcessors.forEach(lifecycle -> lifecycle.onStartRequest(lbRequest, lbResponse));
|
||||
try {
|
||||
Response response = feignClient.execute(feignRequest, options);
|
||||
if (loadBalanced) {
|
||||
supportedLifecycleProcessors.forEach(
|
||||
lifecycle -> lifecycle.onComplete(new CompletionContext<>(CompletionContext.Status.SUCCESS,
|
||||
lbResponse, buildResponseData(response))));
|
||||
lbRequest, lbResponse, buildResponseData(response))));
|
||||
}
|
||||
return response;
|
||||
}
|
||||
catch (Exception exception) {
|
||||
if (loadBalanced) {
|
||||
supportedLifecycleProcessors.forEach(lifecycle -> lifecycle
|
||||
.onComplete(new CompletionContext<>(CompletionContext.Status.FAILED, exception, lbResponse)));
|
||||
supportedLifecycleProcessors.forEach(lifecycle -> lifecycle.onComplete(
|
||||
new CompletionContext<>(CompletionContext.Status.FAILED, exception, lbRequest, lbResponse)));
|
||||
}
|
||||
throw exception;
|
||||
}
|
||||
@@ -83,9 +85,10 @@ final class LoadBalancerUtils {
|
||||
}
|
||||
|
||||
static Response executeWithLoadBalancerLifecycleProcessing(Client feignClient, Request.Options options,
|
||||
Request feignRequest, org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse,
|
||||
Request feignRequest, org.springframework.cloud.client.loadbalancer.Request lbRequest,
|
||||
org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse,
|
||||
Set<LoadBalancerLifecycle> supportedLifecycleProcessors) throws IOException {
|
||||
return executeWithLoadBalancerLifecycleProcessing(feignClient, options, feignRequest, lbResponse,
|
||||
return executeWithLoadBalancerLifecycleProcessing(feignClient, options, feignRequest, lbRequest, lbResponse,
|
||||
supportedLifecycleProcessors, true);
|
||||
}
|
||||
|
||||
|
||||
@@ -44,7 +44,6 @@ import org.springframework.cloud.client.loadbalancer.LoadBalancerClient;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancerLifecycle;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancerLifecycleValidator;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancerProperties;
|
||||
import org.springframework.cloud.client.loadbalancer.RequestData;
|
||||
import org.springframework.cloud.client.loadbalancer.ResponseData;
|
||||
import org.springframework.cloud.client.loadbalancer.RetryableRequestContext;
|
||||
import org.springframework.cloud.client.loadbalancer.RetryableStatusCodeException;
|
||||
@@ -59,7 +58,7 @@ import org.springframework.retry.policy.NeverRetryPolicy;
|
||||
import org.springframework.retry.support.RetryTemplate;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
import static org.springframework.cloud.openfeign.loadbalancer.LoadBalancerUtils.executeWithLoadBalancerLifecycleProcessing;
|
||||
import static org.springframework.cloud.openfeign.loadbalancer.LoadBalancerUtils.buildRequestData;
|
||||
|
||||
/**
|
||||
* A {@link Client} implementation that provides Spring Retry support for requests
|
||||
@@ -107,7 +106,10 @@ public class RetryableFeignBlockingLoadBalancerClient implements Client {
|
||||
Set<LoadBalancerLifecycle> supportedLifecycleProcessors = LoadBalancerLifecycleValidator
|
||||
.getSupportedLifecycleProcessors(
|
||||
loadBalancerClientFactory.getInstances(serviceId, LoadBalancerLifecycle.class),
|
||||
RetryableRequestContext.class, RequestData.class, ServiceInstance.class);
|
||||
RetryableRequestContext.class, ResponseData.class, ServiceInstance.class);
|
||||
String hint = getHint(serviceId);
|
||||
DefaultRequest<RetryableRequestContext> lbRequest = new DefaultRequest<>(
|
||||
new RetryableRequestContext(null, buildRequestData(request), hint));
|
||||
// On retries the policy will choose the server and set it in the context
|
||||
// and extract the server and update the request being made
|
||||
if (context instanceof LoadBalancedRetryContext) {
|
||||
@@ -119,9 +121,7 @@ public class RetryableFeignBlockingLoadBalancerClient implements Client {
|
||||
+ "Reattempting service instance selection");
|
||||
}
|
||||
ServiceInstance previousServiceInstance = lbContext.getPreviousServiceInstance();
|
||||
String hint = getHint(serviceId);
|
||||
DefaultRequest<RetryableRequestContext> lbRequest = new DefaultRequest<>(
|
||||
new RetryableRequestContext(previousServiceInstance, request, hint));
|
||||
lbRequest.getContext().setPreviousServiceInstance(previousServiceInstance);
|
||||
supportedLifecycleProcessors.forEach(lifecycle -> lifecycle.onStart(lbRequest));
|
||||
retrievedServiceInstance = loadBalancerClient.choose(serviceId, lbRequest);
|
||||
if (LOG.isDebugEnabled()) {
|
||||
@@ -136,9 +136,9 @@ public class RetryableFeignBlockingLoadBalancerClient implements Client {
|
||||
}
|
||||
org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse = new DefaultResponse(
|
||||
retrievedServiceInstance);
|
||||
supportedLifecycleProcessors.forEach(
|
||||
lifecycle -> lifecycle.onComplete(new CompletionContext<ResponseData, ServiceInstance>(
|
||||
CompletionContext.Status.DISCARD, lbResponse)));
|
||||
supportedLifecycleProcessors.forEach(lifecycle -> lifecycle
|
||||
.onComplete(new CompletionContext<ResponseData, ServiceInstance, RetryableRequestContext>(
|
||||
CompletionContext.Status.DISCARD, lbRequest, lbResponse)));
|
||||
feignRequest = request;
|
||||
}
|
||||
else {
|
||||
@@ -153,8 +153,9 @@ public class RetryableFeignBlockingLoadBalancerClient implements Client {
|
||||
}
|
||||
org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse = new DefaultResponse(
|
||||
retrievedServiceInstance);
|
||||
Response response = executeWithLoadBalancerLifecycleProcessing(delegate, options, feignRequest, lbResponse,
|
||||
supportedLifecycleProcessors, retrievedServiceInstance != null);
|
||||
Response response = LoadBalancerUtils.executeWithLoadBalancerLifecycleProcessing(delegate, options,
|
||||
feignRequest, lbRequest, lbResponse, supportedLifecycleProcessors,
|
||||
retrievedServiceInstance != null);
|
||||
int responseStatus = response.status();
|
||||
if (retryPolicy != null && retryPolicy.retryableStatusCode(responseStatus)) {
|
||||
if (LOG.isDebugEnabled()) {
|
||||
|
||||
@@ -37,9 +37,9 @@ import org.mockito.junit.jupiter.MockitoExtension;
|
||||
import org.springframework.cloud.client.DefaultServiceInstance;
|
||||
import org.springframework.cloud.client.ServiceInstance;
|
||||
import org.springframework.cloud.client.loadbalancer.CompletionContext;
|
||||
import org.springframework.cloud.client.loadbalancer.DefaultRequestContext;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancerLifecycle;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancerProperties;
|
||||
import org.springframework.cloud.client.loadbalancer.RequestDataContext;
|
||||
import org.springframework.cloud.client.loadbalancer.ResponseData;
|
||||
import org.springframework.cloud.loadbalancer.blocking.client.BlockingLoadBalancerClient;
|
||||
import org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory;
|
||||
@@ -152,15 +152,18 @@ class FeignBlockingLoadBalancerClientTests {
|
||||
|
||||
feignBlockingLoadBalancerClient.execute(request, options);
|
||||
|
||||
Collection<org.springframework.cloud.client.loadbalancer.Request<Object>> lifecycleLogRequests = ((TestLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
Collection<org.springframework.cloud.client.loadbalancer.Request<RequestDataContext>> lifecycleLogRequests = ((TestLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
.get("loadBalancerLifecycle")).getStartLog().values();
|
||||
Collection<CompletionContext<Object, ServiceInstance>> anotherLifecycleLogRequests = ((AnotherLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
Collection<org.springframework.cloud.client.loadbalancer.Request<RequestDataContext>> lifecycleLogStartedRequests = ((TestLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
.get("loadBalancerLifecycle")).getStartRequestLog().values();
|
||||
Collection<CompletionContext<ResponseData, ServiceInstance, RequestDataContext>> anotherLifecycleLogRequests = ((AnotherLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
.get("anotherLoadBalancerLifecycle")).getCompleteLog().values();
|
||||
assertThat(lifecycleLogRequests)
|
||||
.extracting(lbRequest -> ((DefaultRequestContext) lbRequest.getContext()).getHint())
|
||||
assertThat(lifecycleLogRequests).extracting(lbRequest -> lbRequest.getContext().getHint())
|
||||
.contains(callbackTestHint);
|
||||
assertThat(lifecycleLogStartedRequests).extracting(lbRequest -> lbRequest.getContext().getHint())
|
||||
.contains(callbackTestHint);
|
||||
assertThat(anotherLifecycleLogRequests)
|
||||
.extracting(completionContext -> ((ResponseData) completionContext.getClientResponse()).getHttpStatus())
|
||||
.extracting(completionContext -> completionContext.getClientResponse().getHttpStatus())
|
||||
.contains(HttpStatus.OK);
|
||||
}
|
||||
|
||||
@@ -180,30 +183,43 @@ class FeignBlockingLoadBalancerClientTests {
|
||||
|
||||
}
|
||||
|
||||
protected static class TestLoadBalancerLifecycle implements LoadBalancerLifecycle<Object, Object, ServiceInstance> {
|
||||
protected static class TestLoadBalancerLifecycle
|
||||
implements LoadBalancerLifecycle<RequestDataContext, ResponseData, ServiceInstance> {
|
||||
|
||||
final ConcurrentHashMap<String, org.springframework.cloud.client.loadbalancer.Request<Object>> startLog = new ConcurrentHashMap<>();
|
||||
final Map<String, org.springframework.cloud.client.loadbalancer.Request<RequestDataContext>> startLog = new ConcurrentHashMap<>();
|
||||
|
||||
final ConcurrentHashMap<String, CompletionContext<Object, ServiceInstance>> completeLog = new ConcurrentHashMap<>();
|
||||
final Map<String, org.springframework.cloud.client.loadbalancer.Request<RequestDataContext>> startRequestLog = new ConcurrentHashMap<>();
|
||||
|
||||
final Map<String, CompletionContext<ResponseData, ServiceInstance, RequestDataContext>> completeLog = new ConcurrentHashMap<>();
|
||||
|
||||
@Override
|
||||
public void onStart(org.springframework.cloud.client.loadbalancer.Request<Object> request) {
|
||||
public void onStart(org.springframework.cloud.client.loadbalancer.Request<RequestDataContext> request) {
|
||||
startLog.put(getName() + UUID.randomUUID(), request);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onComplete(CompletionContext<Object, ServiceInstance> completionContext) {
|
||||
public void onStartRequest(org.springframework.cloud.client.loadbalancer.Request<RequestDataContext> request,
|
||||
org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse) {
|
||||
startRequestLog.put(getName() + UUID.randomUUID(), request);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onComplete(CompletionContext<ResponseData, ServiceInstance, RequestDataContext> completionContext) {
|
||||
completeLog.put(getName() + UUID.randomUUID(), completionContext);
|
||||
}
|
||||
|
||||
ConcurrentHashMap<String, org.springframework.cloud.client.loadbalancer.Request<Object>> getStartLog() {
|
||||
Map<String, org.springframework.cloud.client.loadbalancer.Request<RequestDataContext>> getStartLog() {
|
||||
return startLog;
|
||||
}
|
||||
|
||||
ConcurrentHashMap<String, CompletionContext<Object, ServiceInstance>> getCompleteLog() {
|
||||
Map<String, CompletionContext<ResponseData, ServiceInstance, RequestDataContext>> getCompleteLog() {
|
||||
return completeLog;
|
||||
}
|
||||
|
||||
Map<String, org.springframework.cloud.client.loadbalancer.Request<RequestDataContext>> getStartRequestLog() {
|
||||
return startRequestLog;
|
||||
}
|
||||
|
||||
protected String getName() {
|
||||
return this.getClass().getSimpleName();
|
||||
}
|
||||
|
||||
@@ -38,11 +38,11 @@ import org.mockito.junit.jupiter.MockitoExtension;
|
||||
import org.springframework.cloud.client.DefaultServiceInstance;
|
||||
import org.springframework.cloud.client.ServiceInstance;
|
||||
import org.springframework.cloud.client.loadbalancer.CompletionContext;
|
||||
import org.springframework.cloud.client.loadbalancer.DefaultRequestContext;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancedRetryFactory;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancerLifecycle;
|
||||
import org.springframework.cloud.client.loadbalancer.LoadBalancerProperties;
|
||||
import org.springframework.cloud.client.loadbalancer.ResponseData;
|
||||
import org.springframework.cloud.client.loadbalancer.RetryableRequestContext;
|
||||
import org.springframework.cloud.loadbalancer.blocking.client.BlockingLoadBalancerClient;
|
||||
import org.springframework.cloud.loadbalancer.blocking.retry.BlockingLoadBalancedRetryPolicy;
|
||||
import org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory;
|
||||
@@ -196,15 +196,18 @@ class RetryableFeignBlockingLoadBalancerClientTests {
|
||||
|
||||
feignBlockingLoadBalancerClient.execute(request, options);
|
||||
|
||||
Collection<org.springframework.cloud.client.loadbalancer.Request<Object>> lifecycleLogRequests = ((TestLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
Collection<org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext>> lifecycleLogRequests = ((TestLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
.get("loadBalancerLifecycle")).getStartLog().values();
|
||||
Collection<CompletionContext<Object, ServiceInstance>> anotherLifecycleLogRequests = ((AnotherLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
Collection<org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext>> lifecycleLogStartedRequests = ((TestLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
.get("loadBalancerLifecycle")).getStartRequestLog().values();
|
||||
Collection<CompletionContext<ResponseData, ServiceInstance, RetryableRequestContext>> anotherLifecycleLogRequests = ((AnotherLoadBalancerLifecycle) loadBalancerLifecycleBeans
|
||||
.get("anotherLoadBalancerLifecycle")).getCompleteLog().values();
|
||||
assertThat(lifecycleLogRequests)
|
||||
.extracting(lbRequest -> ((DefaultRequestContext) lbRequest.getContext()).getHint())
|
||||
assertThat(lifecycleLogRequests).extracting(lbRequest -> lbRequest.getContext().getHint())
|
||||
.contains(callbackTestHint);
|
||||
assertThat(lifecycleLogStartedRequests).extracting(lbRequest -> lbRequest.getContext().getHint())
|
||||
.contains(callbackTestHint);
|
||||
assertThat(anotherLifecycleLogRequests)
|
||||
.extracting(completionContext -> ((ResponseData) completionContext.getClientResponse()).getHttpStatus())
|
||||
.extracting(completionContext -> completionContext.getClientResponse().getHttpStatus())
|
||||
.contains(HttpStatus.OK);
|
||||
}
|
||||
|
||||
@@ -224,30 +227,45 @@ class RetryableFeignBlockingLoadBalancerClientTests {
|
||||
|
||||
}
|
||||
|
||||
protected static class TestLoadBalancerLifecycle implements LoadBalancerLifecycle<Object, Object, ServiceInstance> {
|
||||
protected static class TestLoadBalancerLifecycle
|
||||
implements LoadBalancerLifecycle<RetryableRequestContext, ResponseData, ServiceInstance> {
|
||||
|
||||
final ConcurrentHashMap<String, org.springframework.cloud.client.loadbalancer.Request<Object>> startLog = new ConcurrentHashMap<>();
|
||||
final Map<String, org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext>> startLog = new ConcurrentHashMap<>();
|
||||
|
||||
final ConcurrentHashMap<String, CompletionContext<Object, ServiceInstance>> completeLog = new ConcurrentHashMap<>();
|
||||
final Map<String, org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext>> startRequestLog = new ConcurrentHashMap<>();
|
||||
|
||||
final Map<String, CompletionContext<ResponseData, ServiceInstance, RetryableRequestContext>> completeLog = new ConcurrentHashMap<>();
|
||||
|
||||
@Override
|
||||
public void onStart(org.springframework.cloud.client.loadbalancer.Request<Object> request) {
|
||||
public void onStart(org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext> request) {
|
||||
startLog.put(getName() + UUID.randomUUID(), request);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onComplete(CompletionContext<Object, ServiceInstance> completionContext) {
|
||||
public void onStartRequest(
|
||||
org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext> request,
|
||||
org.springframework.cloud.client.loadbalancer.Response<ServiceInstance> lbResponse) {
|
||||
startRequestLog.put(getName() + UUID.randomUUID(), request);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onComplete(
|
||||
CompletionContext<ResponseData, ServiceInstance, RetryableRequestContext> completionContext) {
|
||||
completeLog.put(getName() + UUID.randomUUID(), completionContext);
|
||||
}
|
||||
|
||||
ConcurrentHashMap<String, org.springframework.cloud.client.loadbalancer.Request<Object>> getStartLog() {
|
||||
Map<String, org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext>> getStartLog() {
|
||||
return startLog;
|
||||
}
|
||||
|
||||
ConcurrentHashMap<String, CompletionContext<Object, ServiceInstance>> getCompleteLog() {
|
||||
Map<String, CompletionContext<ResponseData, ServiceInstance, RetryableRequestContext>> getCompleteLog() {
|
||||
return completeLog;
|
||||
}
|
||||
|
||||
Map<String, org.springframework.cloud.client.loadbalancer.Request<RetryableRequestContext>> getStartRequestLog() {
|
||||
return startRequestLog;
|
||||
}
|
||||
|
||||
protected String getName() {
|
||||
return this.getClass().getSimpleName();
|
||||
}
|
||||
@@ -255,7 +273,7 @@ class RetryableFeignBlockingLoadBalancerClientTests {
|
||||
}
|
||||
|
||||
protected static class AnotherLoadBalancerLifecycle
|
||||
extends FeignBlockingLoadBalancerClientTests.TestLoadBalancerLifecycle {
|
||||
extends RetryableFeignBlockingLoadBalancerClientTests.TestLoadBalancerLifecycle {
|
||||
|
||||
@Override
|
||||
protected String getName() {
|
||||
|
||||
Reference in New Issue
Block a user