欢迎光临
我们一直在努力

微服务组件源码1——服务调用概述

1、常用的方案概述

1.1、方案总结

在SpringCloud第一代中是通过Ribbon负载均衡,第二代是在Ribbon基础上封装了LoadBalancer。LoadBalancer和OpenFeign都属于SpringCloud第二代的产品。

在SpringBoot中可以通过一下方式达到项目通讯:

* RestTemplate结合@LoadBalanced(通信+负载均衡)

* LoadBalancerClient(通信+自定义负载均衡策略)

* Feign(通信+负载均衡)

* OkHttp(通信)

* HttpClient(通信)

1.2、客户端和服务端负载均衡

* 负载均衡分为客户端负载均衡和服务端负载均衡,它们之间的区别在于:服务清单所在的位置。

* 我们通常说的负载均衡都是服务端的负载均衡,其中可以分为硬件的负载均衡和软件负载均衡:硬件的负载均衡就是在服务器节点之间安装用于负载均衡的设备,

比如F5;软件负载均衡则是在服务器上安装一些具有负载均衡功能的模块或软件来完成请求分发的工作,

比如nginx。服务端的负载均衡会在服务端维护一个服务清单,然后通过心跳检测来剔除故障节点以保证服务清单中的节点都正常可用。

* 客户端负载均衡指客户端都维护着自己要访问的服务端实例清单,而这些服务端清单来自服务注册中心。

客户端负载均衡也需要心跳检测维护清单服务的健康性,只不过这个工作要和服务注册中心配合完成。

2、LoadBalancerClient

2.1、代码实例

2.1.1、服务提供方代码

//nacos注册中心的服务名:cloud-producer-server
//两数求和
@PostMapping ("getSum")
public String getSum(@RequestParam (value = "num1") Integer num1, @RequestParam (value = "num2") Integer num2) {
    return "两数求和结果=" + (num1 + num2);
}

2.1.2、服务消费方代码

* 指定服务,通过LoadBalancerClient自动获取某个服务实例与请求地址

  使用时必须指定一个服务名:loadBalancerClient.choose(serviceId);

@LoadBalancerClient(name = "SERVICE-A",configuration = xxx.class)

public class LoadBalancerTest {}

@Component
public class LoadBalancerUtil {
    // 注入LoadBalancerClient
    @Autowired
    LoadBalancerClient loadBalancerClient;

    //通过 LoadBalancer 获取提供服务的hostip
    public String getService(String serviceId) {
        //获取实例服务中的某一个服务
        ServiceInstance instance = loadBalancerClient.choose(serviceId);
        //获取服务的ip地址和端口号
        String host = instance.getHost();
        int port = instance.getPort();
        //格式化最终的访问地址
        return String.format("http://%s:%s", host, port);
    }
}

* 方式一:使用choose方法获取实例,使用restTemplate访问

通过RestTemplated请求远程服务地址url并接收返回值

多次访问服务消费方的 api/invoke/getByLoadBalancer 接口,并且通过打印出来的 hostAndIp 信息是不一样的,

可以看出 LoadBalancerClient 是轮询调用服务提供方的,这也是 LoadBalancerClient 的默认负载均衡策略

@RestController
@RequestMapping (value = "api/invoke")
public class InvokeController {
    @Autowired
    private LoadBalancerUtil loadBalancerUtil;

    /**
     * 使用 SpringCloud 的负载均衡策略组件 LoadBalancerClient 进行远程服务调用
     */
    @GetMapping ("getByLoadBalancer")
    public String getByLoadBalancer(Integer num1, Integer num2) {
        String hostAndIp = loadBalancerUtil.getService("cloud-producer-server");
        //打印服务的请求地址与端口,方便测试负载功能
        System.out.println(hostAndIp);

        String url = hostAndIp + "/cloud-producer-server/getSum";
        MultiValueMap<String, Object> params = new LinkedMultiValueMap<>();
        params.add("num1", num1);
        params.add("num2", num2);

        RestTemplate restTemplate = new RestTemplate();
        String result = restTemplate.postForObject(url, params, String.class);

        return result;
    }
}

* 方式二:使用execute方法,本质和方法一一样

  原理是一样的,LoadBalancerClient只提供ServiceInstance的发现,实际远处调用还的用户自己写

  execute在获取到一个ServiceInstance后回回调传递进来的LoadBalancerRequest的apply方法,这个方法需要用户自定义

public class MyTestController {

    @Autowired
    LoadBalancerClient loadBalancerClient;

    public void query() throws IOException {
        loadBalancerClient.execute("serviceA",new LoadBalancerRequest(){
            @Override
            public Object apply(ServiceInstance instance) throws Exception {
                //获取服务的ip地址和端口号
                String hostAndIp = instance.getHost();
                String url = hostAndIp + "/cloud-producer-server/getSum";
                MultiValueMap<String, Object> params = new LinkedMultiValueMap<>();
                //自己使用通讯工具去远处调用
                RestTemplate restTemplate = new RestTemplate();
                String result = restTemplate.postForObject(url, params, String.class);
                return result;
            }
        });

    }
}

2.2、LoadBalancerClient 原理

2.2.1、概述

* LoadBalancerClient是SpringCloud提供的一种负载均衡客户端

– LoadBalancerClient在初始化时会通过Eureka Client向 Eureka服务端获取所有服务实例的注册信息并缓存在本地,

– 并且每10秒向 EurekaClient 发送“ ping ”,来判断服务的可用性。如果服务的可用性发生了改变或者服务数量和之前的不一致,则更新或者重新拉取。

– 最后,在得到服务注册列表信息后,ILoadBalancer 根据 IRule 的策略进行负载均衡(默认策略为轮询)。

* 当使用 LoadBalancerClient 进行远程调用时

– LoadBalancerClient 先通过目标服务名在本地服务注册清单中获取服务提供方的某一个实例,比如订单服务需要访问商品服务,商品服务有3个节点

– LoadBalancerClient 会通过 choose() 方法获取到3个节点中的一个服务,拿到服务的信息之后取出服务IP信息,就可以得到完整的想要访问的IP地址和接口

– 最后通过 RestTempate 访问商品服务。

2.2.2、LoadBalancerClient接口

* 核心方法就是ServiceInstance choose(String serviceId);

  – LoadBalancerClient的execute相当于对choose的扩展,会获取到service实例后回调LoadBalancerRequest

  – 而LoadBalancerRequest参数需要用户自定义对此serviceInstance的处理,LoadBalancerClient不去关心与如何处理serviceInstance

* 两个实现类

  当RibbonLoadBalancerClient不存在时,才会使用BlockingLoadBalancerClient

public interface ServiceInstanceChooser {
    //根据服务的名称 serviceId 来选择其中一个服务实例,即根据 serviceId 获取ServiceInstance
    ServiceInstance choose(String serviceId);

}

public interface LoadBalancerClient extends ServiceInstanceChooser {
    //执行请求
    <T> T execute(String serviceId, LoadBalancerRequest<T> request) throws IOException;
    <T> T execute(String serviceId, ServiceInstance serviceInstance,LoadBalancerRequest<T> request) throws IOException;
    //重构URL
    URI reconstructURI(ServiceInstance instance, URI original);

}

2.2.3、BlockingLoadBalancerClient实现类

2.2.3.1、自动配置内容

* 其中@LoadBalancerClients用于注入指定的配置类,这里没有指定.class。那就没有了

* @LoadBalancerClients导入的LoadBalancerClientConfigurationRegistrar会从容器中收集延迟对象,类型为ObjectProvider<List<LoadBalancerClientSpecification>> configurations,

会封装到LoadBalancerClientFactory对象中

– 其中LoadBalancerClientSpecification类型的bean是通过所有@LoadBalancerClient和@LoadBalancerClients注解中指定的

@LoadBalancerClient(name = "SERVICE-A",configuration = xxx.class)

   设置到LoadBalancerClientFactory中,执行clientFactory.setConfigurations( Specification configurations),最终会加载到clientFactory中创建的name(对应@LoadBalancerClient的name属性)

容器中

  – @LoadBalancerClients设置的defaultConfiguration也会翻转为LoadBalancerClientSpecification,容器name为default.类名

* 此LoadBalancerClientFactory最终会作为BlockingLoadBalancerClient的构造参数

@Configuration(proxyBeanMethods = false)
@ConditionalOnClass(RestTemplate.class)
@Conditional(OnNoRibbonDefaultCondition.class)
protected static class BlockingLoadbalancerClientConfig {

   @Bean
   @ConditionalOnBean(LoadBalancerClientFactory.class)
   @Primary
   public BlockingLoadBalancerClient blockingLoadBalancerClient(
         LoadBalancerClientFactory loadBalancerClientFactory) {
      return new BlockingLoadBalancerClient(loadBalancerClientFactory);
   }

}

@Configuration(proxyBeanMethods = false)
@LoadBalancerClients
@AutoConfigureBefore({ ReactorLoadBalancerClientAutoConfiguration.class,
      LoadBalancerBeanPostProcessorAutoConfiguration.class,
      ReactiveLoadBalancerAutoConfiguration.class })
public class LoadBalancerAutoConfiguration {

   private final ObjectProvider<List<LoadBalancerClientSpecification>> configurations;

   public LoadBalancerAutoConfiguration(
         ObjectProvider<List<LoadBalancerClientSpecification>> configurations) {
      this.configurations = configurations;
   }

   @Bean
   public LoadBalancerClientFactory loadBalancerClientFactory() {
      LoadBalancerClientFactory clientFactory = new LoadBalancerClientFactory();
      clientFactory.setConfigurations(
            this.configurations.getIfAvailable(Collections::emptyList));
      return clientFactory;
   }

}

@Configuration(proxyBeanMethods = false)
@Retention(RetentionPolicy.RUNTIME)
@Target({ ElementType.TYPE })
@Documented
@Import(LoadBalancerClientConfigurationRegistrar.class)
public @interface LoadBalancerClients {

   LoadBalancerClient[] value() default {};
   Class<?>[] defaultConfiguration() default {};

}

public class LoadBalancerClientConfigurationRegistrar
      implements ImportBeanDefinitionRegistrar {

   private static String getClientName(Map<String, Object> client) {
      if (client == null) {
         return null;
      }
      String value = (String) client.get("value");
      if (!StringUtils.hasText(value)) {
         value = (String) client.get("name");
      }
      if (StringUtils.hasText(value)) {
         return value;
      }
      throw new IllegalStateException(
            "Either 'name' or 'value' must be provided in @LoadBalancerClient");
   }

   private static void registerClientConfiguration(BeanDefinitionRegistry registry,
         Object name, Object configuration) {
      BeanDefinitionBuilder builder = BeanDefinitionBuilder
            .genericBeanDefinition(LoadBalancerClientSpecification.class);
      builder.addConstructorArgValue(name);
      builder.addConstructorArgValue(configuration);
      registry.registerBeanDefinition(name + ".LoadBalancerClientSpecification",
            builder.getBeanDefinition());
   }

   @Override
   public void registerBeanDefinitions(AnnotationMetadata metadata,
         BeanDefinitionRegistry registry) {
      Map<String, Object> attrs = metadata
            .getAnnotationAttributes(LoadBalancerClients.class.getName(), true);
      if (attrs != null && attrs.containsKey("value")) {
         AnnotationAttributes[] clients = (AnnotationAttributes[]) attrs.get("value");
         for (AnnotationAttributes client : clients) {
            registerClientConfiguration(registry, getClientName(client),
                  client.get("configuration"));
         }
      }
      if (attrs != null && attrs.containsKey("defaultConfiguration")) {
         String name;
         if (metadata.hasEnclosingClass()) {
            name = "default." + metadata.getEnclosingClassName();
         }
         else {
            name = "default." + metadata.getClassName();
         }
         registerClientConfiguration(registry, name,
               attrs.get("defaultConfiguration"));
      }
      Map<String, Object> client = metadata
            .getAnnotationAttributes(LoadBalancerClient.class.getName(), true);
      String name = getClientName(client);
      if (name != null) {
         registerClientConfiguration(registry, name, client.get("configuration"));
      }
   }

}

2.3.3.2、基于LoadBalancerClientFactory实现

* LoadBalancerClientFactory是一个容器集合,继承NamedContextFactory

– 其内包含一个容器表:Map<String, ApplicationContextInitializer<GenericApplicationContext>> applicationContextInitializers;

  包含一个收集到设置进来的LoadBalancerClientSpecification配置表:Map<String, C> configurations = new ConcurrentHashMap<>();

– 以name为一个容器维度,没有容器都会存在一个独立的环境,里面存在一个配置:

loadbalancer.client.name = name

  – 配置了父容器为当前Spring容器,在NamedContextFactory中实现了ApplicationContextAware接口的setApplicationContext方法设置了this.parent = parent;

* LoadBalancerClientFactory提供唯一的方法就是getInstance(String serviceId)

  会在以serviceId为name的容器里,获取ReactorServiceInstanceLoadBalancer类型的bean

* 每一个容器都有一个默认的配置:LoadBalancerClientConfiguration,有构造方法中穿件

  – 会为本容器创建一个RoundRobinLoadBalancer和ServiceInstanceListSupplier

* getContext(String name)创建流程

  – 创建一个GenericApplicationContext,设置父容器为当前Spring容器

  – 加入环境信息:使用构造参数传递进来的propertyName作为key,value为方法参数name

  – 从成员变量configurations表中获取此name对应的ClientSpecification和全局默认的ClientSpecification,作为此容器的自动配置类

  – context.refresh()开始初始化此容器

public class LoadBalancerClientFactory
      extends NamedContextFactory<LoadBalancerClientSpecification>
      implements ReactiveLoadBalancer.Factory<ServiceInstance> {

   public static final String NAMESPACE = "loadbalancer";
   public static final String PROPERTY_NAME = NAMESPACE + ".client.name";

   public LoadBalancerClientFactory() {
      super(LoadBalancerClientConfiguration.class, NAMESPACE, PROPERTY_NAME);
   }

   public String getName(Environment environment) {
      return environment.getProperty(PROPERTY_NAME);
   }

   @Override
   public ReactiveLoadBalancer<ServiceInstance> getInstance(String serviceId) {
      return getInstance(serviceId, ReactorServiceInstanceLoadBalancer.class);
   }

}

 */
public abstract class NamedContextFactory<C extends NamedContextFactory.Specification>
      implements DisposableBean, ApplicationContextAware {

   private final Map<String, ApplicationContextInitializer<GenericApplicationContext>> applicationContextInitializers;
   private final String propertySourceName;
   private final String propertyName;
   private final Map<String, GenericApplicationContext> contexts = new ConcurrentHashMap<>();
   private Map<String, C> configurations = new ConcurrentHashMap<>();
   private ApplicationContext parent;
   private Class<?> defaultConfigType;

   public NamedContextFactory(Class<?> defaultConfigType, String propertySourceName, String propertyName) {
      this(defaultConfigType, propertySourceName, propertyName, new HashMap<>());
   }

   public NamedContextFactory(Class<?> defaultConfigType, String propertySourceName, String propertyName,
         Map<String, ApplicationContextInitializer<GenericApplicationContext>> applicationContextInitializers) {
      this.defaultConfigType = defaultConfigType;
      this.propertySourceName = propertySourceName;
      this.propertyName = propertyName;
      this.applicationContextInitializers = applicationContextInitializers;
   }

   
   public void setConfigurations(List<C> configurations) {
      for (C client : configurations) {
         this.configurations.put(client.getName(), client);
      }
   }

   protected GenericApplicationContext getContext(String name) {
      if (!this.contexts.containsKey(name)) {
         synchronized (this.contexts) {
            if (!this.contexts.containsKey(name)) {
               this.contexts.put(name, createContext(name));
            }
         }
      }
      return this.contexts.get(name);
   }

   public GenericApplicationContext createContext(String name) {
      GenericApplicationContext context = buildContext(name);
      if (applicationContextInitializers.get(name) != null) {
         applicationContextInitializers.get(name).initialize(context);
         context.refresh();
         return context;
      }
      registerBeans(name, context);
      context.refresh();
      return context;
   }

   public void registerBeans(String name, GenericApplicationContext context) {
      Assert.isInstanceOf(AnnotationConfigRegistry.class, context);
      AnnotationConfigRegistry registry = (AnnotationConfigRegistry) context;
      if (this.configurations.containsKey(name)) {
         for (Class<?> configuration : this.configurations.get(name).getConfiguration()) {
            registry.register(configuration);
         }
      }
      for (Map.Entry<String, C> entry : this.configurations.entrySet()) {
         if (entry.getKey().startsWith("default.")) {
            for (Class<?> configuration : entry.getValue().getConfiguration()) {
               registry.register(configuration);
            }
         }
      }
      registry.register(PropertyPlaceholderAutoConfiguration.class, this.defaultConfigType);
   }

   public GenericApplicationContext buildContext(String name) {
      GenericApplicationContext context;
      if (this.parent != null) {
        //…       

}
      else {
         context = AotDetector.useGeneratedArtifacts() ? new GenericApplicationContext()
               : new AnnotationConfigApplicationContext();
      }
      context.getEnvironment().getPropertySources().addFirst(
            new MapPropertySource(this.propertySourceName, Collections.singletonMap(this.propertyName, name)));
      if (this.parent != null) {
         context.setParent(this.parent);
      }
      context.setDisplayName(generateDisplayName(name));
      return context;
   }

  public <T> T getInstance(String name, Class<T> type) {
      GenericApplicationContext context = getContext(name);
      try {
         return context.getBean(type);
      }
      catch (NoSuchBeanDefinitionException e) {
         // ignore
      }
      return null;
   }
   //封装为FactoryObjectProvider或ObjectProvider等延迟加载的类型
   public <T> ObjectProvider<T> getLazyProvider(String name, Class<T> type) {
      return new ClientFactoryObjectProvider<>(this, name, type);
   }
   public <T> ObjectProvider<T> getProvider(String name, Class<T> type) {
      GenericApplicationContext context = getContext(name);
      return context.getBeanProvider(type);
   }

   public <T> T getInstance(String name, Class<?> clazz, Class<?>… generics) {
      ResolvableType type = ResolvableType.forClassWithGenerics(clazz, generics);
      return getInstance(name, type);
   }

   public <T> Map<String, T> getInstances(String name, Class<T> type) {
      GenericApplicationContext context = getContext(name);

      return BeanFactoryUtils.beansOfTypeIncludingAncestors(context, type);
   }

   public interface Specification {

      String getName();
      Class<?>[] getConfiguration();

   }

//
}

@Configuration(proxyBeanMethods = false)
@EnableConfigurationProperties(LoadBalancerProperties.class)
@ConditionalOnDiscoveryEnabled
public class LoadBalancerClientConfiguration {

   private static final int REACTIVE_SERVICE_INSTANCE_SUPPLIER_ORDER = 193827465;

   @Bean
   @ConditionalOnMissingBean
   LoadBalancerProperties loadBalancerProperties() {
      return new LoadBalancerProperties();
   }

   @Bean
   @ConditionalOnMissingBean
   public ReactorLoadBalancer<ServiceInstance> reactorServiceInstanceLoadBalancer(
         Environment environment,
         LoadBalancerClientFactory loadBalancerClientFactory) {
      String name = environment.getProperty(LoadBalancerClientFactory.PROPERTY_NAME);
      return new RoundRobinLoadBalancer(loadBalancerClientFactory.getLazyProvider(name,
            ServiceInstanceListSupplier.class), name);
   }

 
   @Configuration(proxyBeanMethods = false)
   @ConditionalOnBlockingDiscoveryEnabled
   @Order(REACTIVE_SERVICE_INSTANCE_SUPPLIER_ORDER + 1)
   public static class BlockingSupportConfiguration {

      @Bean
      @ConditionalOnBean(DiscoveryClient.class)
      @ConditionalOnMissingBean
      public ServiceInstanceListSupplier discoveryClientServiceInstanceListSupplier(
            DiscoveryClient discoveryClient, Environment env,
            ApplicationContext context) {
         DiscoveryClientServiceInstanceListSupplier delegate = new DiscoveryClientServiceInstanceListSupplier(
               discoveryClient, env);
         ObjectProvider<LoadBalancerCacheManager> cacheManagerProvider = context
               .getBeanProvider(LoadBalancerCacheManager.class);
         if (cacheManagerProvider.getIfAvailable() != null) {
            return new CachingServiceInstanceListSupplier(delegate,
                  cacheManagerProvider.getIfAvailable());
         }
         return delegate;
      }

      @Bean
      @ConditionalOnBean(DiscoveryClient.class)
      @ConditionalOnMissingBean
      public ServiceInstanceSupplier discoveryClientServiceInstanceSupplier(
            DiscoveryClient discoveryClient, Environment env,
            ApplicationContext context) {
         DiscoveryClientServiceInstanceSupplier delegate = new DiscoveryClientServiceInstanceSupplier(
               discoveryClient, env);
         ObjectProvider<LoadBalancerCacheManager> cacheManagerProvider = context
               .getBeanProvider(LoadBalancerCacheManager.class);
         if (cacheManagerProvider.getIfAvailable() != null) {
            return new CachingServiceInstanceSupplier(delegate,
                  cacheManagerProvider.getIfAvailable());
         }
         return delegate;
      }

   }

//
}

2.3.3.3、BlockingLoadBalancerClient实现

* 每一个容器都有一个默认的配置:LoadBalancerClientConfiguration

  会为本容器创建一个RoundRobinLoadBalancer和DiscoveryClientServiceInstanceListSupplier,其中name为serverId

new RoundRobinLoadBalancer(loadBalancerClientFactory.getLazyProvider(name,
      ServiceInstanceListSupplier.class), name);

* 核心方法choose方法时基于ReactiveLoadBalancer对象的choose()实现的

  – 而ReactiveLoadBalancer的choose()有借助于(ServiceInstanceListSupplier)DiscoveryClientServiceInstanceListSupplier实现,获取到指定serverId的实例列表

  – 在基于此列表进行轮询,返回一个serverIntance

  – 最终要的是,列表的查询是借助如传递进来的DiscoveryClient实例。例如EurekaDiscoveryClient

   注:这个实例需要注入到每一个name容器中才能使用,可以通过@LoadBalancerClient注解配置一全局默认的配置类型,在里面指定DiscoveryClient实例,这个每一个LoadBalancerClient都共享这个DiscoveryClient实例

public class BlockingLoadBalancerClient implements LoadBalancerClient {

   private final LoadBalancerClientFactory loadBalancerClientFactory;

   public BlockingLoadBalancerClient(
         LoadBalancerClientFactory loadBalancerClientFactory) {
      this.loadBalancerClientFactory = loadBalancerClientFactory;
   }

   @Override
   public <T> T execute(String serviceId, LoadBalancerRequest<T> request)
         throws IOException {
      ServiceInstance serviceInstance = choose(serviceId);
      if (serviceInstance == null) {
         throw new IllegalStateException("No instances available for " + serviceId);
      }
      return execute(serviceId, serviceInstance, request);
   }

   @Override
   public <T> T execute(String serviceId, ServiceInstance serviceInstance,
         LoadBalancerRequest<T> request) throws IOException {
      try {
         return request.apply(serviceInstance);
      }
      catch (IOException iOException) {
         throw iOException;
      }
      catch (Exception exception) {
         ReflectionUtils.rethrowRuntimeException(exception);
      }
      return null;
   }

   @Override
   public URI reconstructURI(ServiceInstance serviceInstance, URI original) {
      return LoadBalancerUriTools.reconstructURI(serviceInstance, original);
   }

   @Override
   public ServiceInstance choose(String serviceId) {
      ReactiveLoadBalancer<ServiceInstance> loadBalancer = loadBalancerClientFactory
            .getInstance(serviceId);
      if (loadBalancer == null) {
         return null;
      }
      Response<ServiceInstance> loadBalancerResponse = Mono.from(loadBalancer.choose())
            .block();
      if (loadBalancerResponse == null) {
         return null;
      }
      return loadBalancerResponse.getServer();
   }

}

public class RoundRobinLoadBalancer implements ReactorServiceInstanceLoadBalancer {

   private final AtomicInteger position;
   private ObjectProvider<ServiceInstanceListSupplier> serviceInstanceListSupplierProvider;
   private final String serviceId;

   public RoundRobinLoadBalancer(
         ObjectProvider<ServiceInstanceListSupplier> serviceInstanceListSupplierProvider,
         String serviceId, int seedPosition) {
      this.serviceId = serviceId;
      this.serviceInstanceListSupplierProvider = serviceInstanceListSupplierProvider;
      this.position = new AtomicInteger(seedPosition);
   }

   
   @Override
   public Mono<Response<ServiceInstance>> choose(Request request) {
      if (serviceInstanceListSupplierProvider != null) {
         ServiceInstanceListSupplier supplier = serviceInstanceListSupplierProvider
               .getIfAvailable(NoopServiceInstanceListSupplier::new);
         return supplier.get().next().map(this::getInstanceResponse);
      }
      ServiceInstanceSupplier supplier = this.serviceInstanceSupplier
            .getIfAvailable(NoopServiceInstanceSupplier::new);
      return supplier.get().collectList().map(this::getInstanceResponse);
   }

   private Response<ServiceInstance> getInstanceResponse(
         List<ServiceInstance> instances) {
      if (instances.isEmpty()) {
         return new EmptyResponse();
      }
      //默认为轮询
      int pos = Math.abs(this.position.incrementAndGet());
      ServiceInstance instance = instances.get(pos % instances.size());
      return new DefaultResponse(instance);
   }

}

public class DiscoveryClientServiceInstanceListSupplier
      implements ServiceInstanceListSupplier {

   private final String serviceId;

   private final Flux<ServiceInstance> serviceInstances;

   public DiscoveryClientServiceInstanceListSupplier(DiscoveryClient delegate,
         Environment environment) {
      this.serviceId = environment.getProperty(PROPERTY_NAME);
      this.serviceInstances = Flux
            .defer(() -> Flux.fromIterable(delegate.getInstances(serviceId)))
            .subscribeOn(Schedulers.boundedElastic());
   }

   public DiscoveryClientServiceInstanceListSupplier(ReactiveDiscoveryClient delegate,
         Environment environment) {
      this.serviceId = environment.getProperty(PROPERTY_NAME);
      this.serviceInstances = delegate.getInstances(serviceId);
   }

   @Override
   public String getServiceId() {
      return serviceId;
   }

   @Override
   public Flux<List<ServiceInstance>> get() {
      return serviceInstances.collectList().flux();
   }

}

2.3.4、RibbonLoadBalancerClient实现类

* 原理和BlockingLoadBalancerClient类似

  – 自动配置相关的内容,见下章节

  – 同样使用一个Name容器工厂SpringClientFactory来作为各个服务id的隔离容器,里面存储了需要的bean

* 核心方法execute(String serviceId, LoadBalancerRequest<T> request)的实现

  – 获取ILoadBalancer实现类(默认ZoneAwareLoadBalancer),使用loadBalancer.chooseServer()获取指定serviceId的一个服务实例

  – 调用入参LoadBalancerRequest,进行处理serviceInstance,区别在于

    这里会使用RibbonLoadBalancerContext记录此服务实例的统计信息

      

public class RibbonLoadBalancerClient implements LoadBalancerClient {

   private SpringClientFactory clientFactory;

   public RibbonLoadBalancerClient(SpringClientFactory clientFactory) {
      this.clientFactory = clientFactory;
   }

   @Override
   public URI reconstructURI(ServiceInstance instance, URI original) {
      //..
   }

   @Override
   public ServiceInstance choose(String serviceId) {
      return choose(serviceId, null);
   }
   public ServiceInstance choose(String serviceId, Object hint) {
      Server server = getServer(getLoadBalancer(serviceId), hint);
      if (server == null) {
         return null;
      }
      return new RibbonServer(serviceId, server, isSecure(server, serviceId),
            serverIntrospector(serviceId).getMetadata(server));
   }
   @Override
   public <T> T execute(String serviceId, LoadBalancerRequest<T> request)
         throws IOException {
      return execute(serviceId, request, null);
   }
   public <T> T execute(String serviceId, LoadBalancerRequest<T> request, Object hint)
         throws IOException {
      //获取ILoadBalancer实现类(默认ZoneAwareLoadBalancer
      ILoadBalancer loadBalancer = getLoadBalancer(serviceId);
      //loadBalancer.chooseServer();
      Server server = getServer(loadBalancer, hint);
      if (server == null) {
         throw new IllegalStateException("No instances available for " + serviceId);
      }
      //server的简单封装
      RibbonServer ribbonServer = new RibbonServer(serviceId, server,
            isSecure(server, serviceId),
            serverIntrospector(serviceId).getMetadata(server));

      return execute(serviceId, ribbonServer, request);
   }
   protected ILoadBalancer getLoadBalancer(String serviceId) {
      return this.clientFactory.getLoadBalancer(serviceId);
   }

   @Override
   public <T> T execute(String serviceId, ServiceInstance serviceInstance,
         LoadBalancerRequest<T> request) throws IOException {
      Server server = null;
      if (serviceInstance instanceof RibbonServer) {
         server = ((RibbonServer) serviceInstance).getServer();
      }
      if (server == null) {
         throw new IllegalStateException("No instances available for " + serviceId);
      }
        //用于记录服务实例的统计信息
      RibbonLoadBalancerContext context = this.clientFactory
            .getLoadBalancerContext(serviceId);
      RibbonStatsRecorder statsRecorder = new RibbonStatsRecorder(context, server);

      try {
         //调用入参LoadBalancerRequest,进行处理serviceInstance
         T returnVal = request.apply(serviceInstance);
         //对应的RibbonLoadBalancerContextnoteRequestCompletion
         statsRecorder.recordStats(returnVal);
         return returnVal;
      }
      catch (IOException ex) {
         statsRecorder.recordStats(ex);
         throw ex;
      }
      catch (Exception ex) {
         statsRecorder.recordStats(ex);
         ReflectionUtils.rethrowRuntimeException(ex);
      }
      return null;
   }

   private ServerIntrospector serverIntrospector(String serviceId) {
      ServerIntrospector serverIntrospector = this.clientFactory.getInstance(serviceId,
            ServerIntrospector.class);
      if (serverIntrospector == null) {
         serverIntrospector = new DefaultServerIntrospector();
      }
      return serverIntrospector;
   }

   private boolean isSecure(Server server, String serviceId) {
      IClientConfig config = this.clientFactory.getClientConfig(serviceId);
      ServerIntrospector serverIntrospector = serverIntrospector(serviceId);
      return RibbonUtils.isSecure(config, serverIntrospector, server);
   }

   // Note: This method could be removed?
   protected Server getServer(String serviceId) {
      return getServer(getLoadBalancer(serviceId), null);
   }
   protected Server getServer(ILoadBalancer loadBalancer) {
      return getServer(loadBalancer, null);
   }
   protected Server getServer(ILoadBalancer loadBalancer, Object hint) {
      if (loadBalancer == null) {
         return null;
      }
      // Use 'default' on a null hint, or just pass it on?
      return loadBalancer.chooseServer(hint != null ? hint : "default");
   }

   
   public static class RibbonServer implements ServiceInstance {

      private final String serviceId;
      private final Server server;
      private final boolean secure;
      private Map<String, String> metadata;
      public RibbonServer(String serviceId, Server server) {
         this(serviceId, server, false, Collections.emptyMap());
      }
      
      public RibbonServer(String serviceId, Server server, boolean secure,
            Map<String, String> metadata) {
         this.serviceId = serviceId;
         this.server = server;
         this.secure = secure;
         this.metadata = metadata;
      }

      //..

   }

}

3、Spring Ribbon (@LoadBalanced + RestTemplate)

3.1、代码实例

通过 Spring Cloud Ribbon 的封装,我们在微服务架构中使用客户端负载均衡非常简单,只需要两步:

    ① 服务提供者启动服务实例并注册到服务注册中心

    ② 服务消费者直接使用被 @LoadBalanced 注解修饰的 RestTemplate 来实现面向服务的接口调用

3.1.1、服务提供方代码

//nacos注册中心的服务名:cloud-producer-server
//两数求和
@PostMapping ("getSum")
public String getSum(@RequestParam (value = "num1") Integer num1, @RequestParam (value = "num2") Integer num2) {
    return "两数求和结果=" + (num1 + num2);
}

3.1.2、服务消费方代码

(1)使用 @LoadBalanced 注解修饰的 RestTemplate:

@LoadBalanced 注解用于开启负载均衡,标记 RestTemplate 使用 LoadBalancerClient 配置

@Configuration
public class RestConfig {
    /**
     * 创建restTemplate对象。
     * LoadBalanced注解表示赋予restTemplate使用Ribbon的负载均衡的能力(一定要加上注解,否则无法远程调用)
     */
    @Bean
    @LoadBalanced
    public RestTemplate restTemplate(){
        return new RestTemplate();
    }
}

(2)通过 RestTemplate 请求远程服务地址并接收返回值

@RestController
@RequestMapping (value = "api/invoke")
public class InvokeController{
    @Autowired
    private RestTemplate restTemplate;

    /**
     * 使用 RestTemplate 进行远程服务调用,并且使用 Ribbon 进行负载均衡
     */
    @ApiOperation (value = "RestTemplate", notes = "使用RestTemplate进行远程服务调用,并使用Ribbon进行负载均衡")
    @GetMapping ("getByRestTemplate")
    public String getByRestTemplate(Integer num1, Integer num2){
        //第一个cloud-producer-server代表在nacos注册中心中的服务名,第二个cloud-producer-server代表contextPath配置的项目路径
        String url = "http://cloud-producer-server/cloud-producer-server/getSum";
        MultiValueMap<String, Object> params = new LinkedMultiValueMap<>();
        params.add("num1", num1);
        params.add("num2", num2);

        //通过服务名的方式调用远程服务(非ip端口)
        return restTemplate.postForObject(url, params, String.class);
    }
}

默认情况下,Ribbon 也是使用轮询作为负载均衡策略,那么处理轮询策略,Ribbon 还有哪些负载均衡策略呢?

3.2、参数配置

3.2.1、配置方式

对于Ribbon的配置一般有两种配置方式:全局配置和指定客户端配置

  • 全局配置
  • * 依赖Spring的自动扫描(不推荐)

      Ribbon的配置类一定不能Spring扫描到。因为Ribbon有自己的子上下文,Spring的父上下文如果和Ribbon的子上下文重叠,会有各种各样的问题

    @Configuration
         public class MyRibbonConfiguration {
            @Bean
            public IPing ribbonPing(){
                return new PingUrl();
            }
        }

    * 使用@RibbonClients注解的defaultConfiguration进行全局默认配置

    @Configuration
    @RibbonClients(defaultConfiguration = MyRibbonConfiguration.class)
    public class GreetingServiceRibbonConf {}

    * 配置文件形式

    使用 ribbon.<key>=<value> 的形式配置,例如全局配置连接超时时间:

    ribbon:

       ConnectTimeout: 250

    (2)局部配置

      * 使用@RibbonClients注解,指定服务名name 和对应的配置文件configuration

    @Configuration
         @RibbonClient(name = "orderService",configuration = HelloRibbonConfiguration.class)
         public class RibbonConfiguration {}

    @Configuration
         public class HelloRibbonConfiguration {
            @Bean
            public IPing ribbonPing(){
                return new PingUrl();
           }
        }

    * 文件配置

        – 如果同时配置了全局配置和指定客户端配置,那么以指定客户端的配置为准。

    – prperties格式: <client>.ribbon.<key>=<value>

    – yaml格式如下

    my-service:

         ribbon:

         listOfServers: localhost:8080,localhost:8081 #指定服务列表

         NFLoadBalancerRuleClassName: com.netflix.loadbalancer.RandomRule #设置负载均衡策略

     注:全量的配置项可以参考com.netflix.client.config.CommonClientConfigKey中的配置,由于配置项太多,就不一一说明。

    (3)两种方式对比:

         * 代码配置:基于代码、更加灵活;但是线上修改得重新打包、发布,并且还有小坑(父子上下文问题)

         * 文件配置: 配置更加直观、优先级更高(相对代码配置)、线上修改无需重新打包、发布;但是极端场景下没有代码配置方式灵活。

    注意:如果代码配置和文件配置两种方式混用,文件配置优先级更高。

       

    尽量使用文件配置,属性方式实现不了的情况下再考虑代码配置

    同一个微服务内尽量保持单一性,使用同样的配置方式,避免两种方式混用,增加定位代码的复杂性。

    3.2.2、Ribbon的七种负载均衡策略

    * 配置方式

    my-service:

         ribbon:

         NFLoadBalancerRuleClassName: com.netflix.loadbalancer.RandomRule #设置负载均衡策略

    * com.netflix.loadbalancer.IRule接口,该接口的实现类主要用于定义负载均衡策略,我们找到它所有的实现类,如下:

    客户端本地会记录每个服务的状态ServerStats,会记录每个url的连接性、并发度、平均响应时间等信息,基于这些信息可进行不同的负载均衡策略

    * 随机策略 RandomRule

    随机数选择服务列表中的服务节点Server,如果当前节点不可用,则进入下一轮随机策略,直到选到可用服务节点为止

    * 轮询策略 RoundRobinRule

    按照接收的请求顺序,逐一分配到不同的后端服务器

    * 重试策略 RetryRule

    在选定的负载均衡策略机上重试机制,在一个配置时间段内当选择Server不成功,则一直尝试使用 subRule 的方式选择一个可用的server;

    * 可用过滤策略 PredicateBaseRule

    过滤掉连接失败 和 高并发连接 的服务节点,然后从健康的服务节点中以线性轮询的方式选出一个节点返回

    * 响应时间权重策略 WeightedRespinseTimeRule

    根据服务器的响应时间分配一个权重weight,响应时间越长,weight越小,被选中的可能性越低。

    主要通过后台线程定期地从 status 里面读取平均响应时间,为每个 server 计算一个 weight

    * 并发量最小可用策略 BestAvailableRule

    选择一个并发量最小的服务节点 server。ServerStats 的 activeRequestCount 属性记录了 server 的并发量,轮询所有的server,选择其中 activeRequestCount 最小的那个server,

    就是并发量最小的服务节点。该策略的优点是可以充分考虑每台服务节点的负载,把请求打到负载压力最小的服务节点上。

    但是缺点是需要轮询所有的服务节点,如果集群数量太大,那么就会比较耗时。

    * 区域权重策略 ZoneAvoidanceRule

    综合判断 server 所在区域的性能 和 server 的可用性,使用 ZoneAvoidancePredicate 和 AvailabilityPredicate 来判断是否选择某个server,

    前一个判断判定一个zone的运行性能是否可用,剔除不可用的zone(的所有server),AvailabilityPredicate 用于过滤掉连接数过多的Server。

    3.2.3、重试机制

    众所周知,分布式服务治理有CAP原则(一致性、可用性、可靠性),其中代表性的Eureka保证了AP(可用性和可靠性),而Zookeeper保证了CP(一致性、可靠性)。

    由于Eureka保证了可用性而舍去了一致性,因此无论是触发了保护机制还是服务剔除延迟,最终导致调用到故障的实例的时候,我们还是希望增强对这类问题的容错,因此Ribbon就提供了充实策略。

    对于重试机制的参数配置如下代码所示,具体含义描述已在代码中注释

    # 开启重试机制,默认为关闭spring:

      cloud:

        loadbalancer:

          retry:

            enable: true

    #断路器的超时时间需要大于Ribbon的超时时间,不然不会触发重试

    hystrix:

      command:

        default:

          execution:

            isolation:

              thread:

                timeoutInMilliseconds: 10000

    EUREKA-CLIENT:

      ribbon:

        #请求链接的超时时间

        ConnectTimeOut: 250

        #请求处理的超时时间

        ReadTimeOut: 1000

        #对所有操作都重试

        OkToRetyrOnAllOperations: true

        #切换实例的重试次数

        MaxAutoRetyiesNextServer: 2

        #对当前实例的重试次数

    4、openFeign

      微服务架构中,由于对服务粒度的拆分致使服务数量变多,而作为 Web 服务的调用端方,除了需要熟悉各种 Http 客户端,比如 okHttp、HttpClient 组件的使用,而且还要显式地序列化和反序列化请求和响应内容,从而导致出现很多样板代码,开发起来很痛苦。为了解决这个问题,Feign 诞生了,那么 Feign 是什么呢?

    Feign 就是一个 Http 客户端的模板,目标是减少 HTTP API 的复杂性,希望能将 HTTP 远程服务调用做到像 RPC 一样易用。Feign 集成 RestTemplate、Ribbon 实现了客户端的负载均衡的 Http 调用,并对原调用方式进行了封装,使得开发者不必手动使用 RestTemplate 调用服务,而是声明一个接口,并在这个接口中标注一个注解即可完成服务调用,这样更加符合面向接口编程的宗旨,客户端在调用服务端时也不需要再关注请求的方式、地址以及restTemplate是 forObject 还是 forEntity,结构更加明了,耦合也更低,简化了开发。但 Feign 已经停止迭代了,所以本篇文章我们也不过多的介绍,而在 Feign 的基础上,又衍生出了 openFeign,那么 openFeign 又是什么呢?

    openFeign 在 Feign 的基础上支持了 SpringMVC 的注解,如 @RequestMapping 等。OpenFeign 的 @FeignClient 可以解析 SpringMVC 的 @RequestMapping 注解下的接口,并通过动态代理的方式产生实现类,实现类中做负载均衡并调用其他服务。

        总的就是,openFeign 作为微服务架构下服务间调用的解决方案,是一种声明式、模板化的 HTTP 的模板,使 HTTP 请求就像调用本地方法一样,通过 openFeign 可以替代基于 RestTemplate 的远程服务调用,并且默认集成了 Ribbon 进行负载均衡。RibbonLoadBalancerClient

    赞(0)
    未经允许不得转载:171主机测评 » 微服务组件源码1——服务调用概述
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址