Inventory Cloud: Стратегия балансировки нагрузки

задняя часть
Inventory Cloud: Стратегия балансировки нагрузки

Это второй день моего участия в Gengwen Challenge.Подробности мероприятия смотрите:Обновить вызов

Общая документация:Каталог статей
Github : github.com/black-ant

Введение

Это одна из серии статей Фейна, в которой в основном объясняются связанные с этим знания Фейна о балансировке нагрузки.Статья включает следующее содержание:

  • Обработка балансировки нагрузки, когда Feign использует ленту

2. Стратегия голосования Файна

Сам Feign имеет набор логики опроса для обработки вызовов сервера.Основной процесс выглядит следующим образом:

  • Шаг 1: Вызовите submit, чтобы получить информацию о сервере
  • Шаг 2: selectServer Выберите сервер
  • Шаг 3: Обратный вызов для выполнения запроса в ServerOperation

Логика раннего вызова

Чтобы облегчить понимание, давайте взглянем на его предыдущую логику вызова Feign в основном инициирует прокси метода через Invoke:

  • Шаг 1: Запустите InvoKe (FeignInvocationHandler)
  • Шаг 2: Вызов SynchronousMethodHandler, обработка логики метода
  • Шаг 3: LoadBalancerFeignClient вызывает основную логику
  • Шаг 4: AbstractLoadBalancerAwareClient в середине перехода, обработка URL-адреса, создание реального URL-адреса
  • Шаг 5: Последний вызов — это внутренний класс клиента (по умолчанию).

2.1 Получить запись сервера

Видно, что на четвертом шаге выше запускается соответствующая обработка балансировки нагрузки:

executeWithLoadBalancer — это метод балансировки нагрузки в AbstractLoadBalancerAwareClient, который находится вcom.netflix.loadbalancerв сумке

PS: сам OpenFeign содержит зависимости Ribbon!

Во-первых, давайте посмотримexecuteWithLoadBalancerОсновная логика, где Submit — это вызов синтаксического сахара Function:

C- AbstractLoadBalancerAwareClient
public T executeWithLoadBalancer(final S request, final IClientConfig requestConfig) throws ClientException {
        LoadBalancerCommand<T> command = buildLoadBalancerCommand(request, requestConfig);

    try {
        // Step 1 : submit 中获取 Server
        // Step 2 : 通过 ServerOperation 语法糖 , 执行具体的操作
        return command.submit(// 暂时省略 Function:003 )
            .toBlocking()
            .single();
    } catch (Exception e) {
        Throwable t = e.getCause();
        if (t instanceof ClientException) {
            throw (ClientException) t;
        } else {
            throw new ClientException(e);
        }
    }
}

Давайте посмотрим, что делается в методе Submit:

  • Шаг 1: Подготовка контейнера
  • Шаг 2. Включен прослушиватель балансировки нагрузки.
  • Шаг 3: Установите количество попыток
    • maxRetrysSame: максимальное количество попыток выполнения на сервере
    • maxRetrysNext: максимальное количество отдельных серверов для повторной попытки
  • ШАГ 4: Получите сервер с помощью балансировщика нагрузки]
  • Шаг 5: Обработка политики повторных попыток
  • Шаг 6: обработка исключений процесса
// C- LoadBalancerCommand : 做过魔改 , 仅展示核心的代码
public Observable<T> submit(final ServerOperation<T> operation) {

    // Step 1 : 容器准备
    final ExecutionInfoContext context = new ExecutionInfoContext();
    
    // Step 2 : ExecutionContextListenerInvoker 负载均衡器在执行的不同阶段调用的侦听器
    if (listenerInvoker != null) {
        try {
            listenerInvoker.onExecutionStart();
        } catch (AbortExecutionException e) {
            return Observable.error(e);
        }
    }

    // Step 3 : 重试次数设置
    final int maxRetrysSame = retryHandler.getMaxRetriesOnSameServer();
    final int maxRetrysNext = retryHandler.getMaxRetriesOnNextServer();

    // Step 4 : 使用负载均衡器获得 Server 
    // 选择 Server , 此处已经通过 Rule 选择完成
    Observable<Server> servers = server == null ? selectServer() : Observable.just(server);
    Observable<T> o = servers.concatMap(// PS : Function:001 详见);
        
    // Step 5 : 重试策略的处理
    if (maxRetrysNext > 0 && server == null){
        o = o.retry(retryPolicy(maxRetrysNext, false));
    }
    
    // Step 6 : 流程异常处理
    return o.onErrorResumeNext( // PS : Function:002 详见);
}

Фокус:На шаге 4 selectServer завершил процесс выбора балансировки нагрузки, а затем в concatMap выполняются дальнейшие операции (функция 1):

[Pro]: Что делается в функции:001 ?

Суммировать:Проще говоря, это аудит данных, запуск прослушивателя и расширение метода для последующих обратных вызовов.

C- LoadBalancerCommand

// 观察者对象编号 : Observable:001
new Func1<Server, Observable<T>>() {
    @Override
    // Called for each server being selected
    public Observable<T> call(Server server) {
        // 容器设置 Server
        context.setServer(server);
    
        // 捕获LoadBalancer中的每个服务器(节点)的各种统计数据
        final ServerStats stats = loadBalancerContext.getServerStats(server);
                    
        //  尝试调用和重试操作 (Observable:002)
        Observable<T> o = Observable.just(server).concatMap( new Func1<Server, Observable<T>>(){
            @Override
            public Observable<T> call(final Server server) {
                context.incAttemptCount();
                loadBalancerContext.noteOpenConnection(stats);
                                    
                if (listenerInvoker != null) {
                    try {
                        // 当选择服务器并且请求将在服务器上执行时调用
                        listenerInvoker.onStartWithServer(context.toExecutionInfo());
                    } catch (AbortExecutionException e) {
                        return Observable.error(e);
                    }
                }
            
                // 计时器开始
                final Stopwatch tracer = loadBalancerContext.getExecuteTracer().start();

                // 提供接收基于推送的通知的机制 -> UNDO:001 / Observable:003
                return operation.call(server).doOnEach(new Observer<T>() {

                    private T entity;

                    // 省略其中代码 , 其中实现了 四个方法 :

                    // - onCompleted : 完成后调用
                    recordStats(tracer, stats, entity, null)

                    // - onError : 异常时调用 , 区别是传入了 exception
                    recordStats(tracer, stats, null, e)
                    listenerInvoker.onExceptionWithServer(e, context.toExecutionInfo())

                    // - onNext : 为观察者提供一个要观察的新项目 , 设置 entity + 触发监听器
                    this.entity = entity
                    listenerInvoker.onExecutionSuccess(entity, context.toExecutionInfo())


                    // - recordStats : 停止计时 + 更新统计数据
                    tracer.stop()
                    oadBalancerContext.noteRequestCompletion(
                        stats, entity, exception, tracer.getDuration(TimeUnit.MILLISECONDS), retryHandler)
                });
        }
}
 


[Pro]: Что делается в функции:002 ?

Функция 2 — это процесс обработки исключений, который в основном делает три вещи:

  • Создавайте разные ClientExceptions с разными конфигурациями повторных попыток.
  • listenerInvoker существует, запускает событие
  • Возвращает объект наблюдателя, который будет вызываться, когда наблюдатель подпишется на этот Observable.ObserverизObserver#onErrorметод
C- LoadBalancerCommand
new Func1<Throwable, Observable<T>>() {
        @Override
        public Observable<T> call(Throwable e) {
            // 建不同的 ClientException
            if (context.getAttemptCount() > 0) {
                if (maxRetrysNext > 0 && context.getServerAttemptCount() == (maxRetrysNext + 1)) {
                    e = new ClientException(......);
                }
                else if (maxRetrysSame > 0 && context.getAttemptCount() == (maxRetrysSame + 1)) {
                    e = new ClientException(......);
                }
            }
            
            if (listenerInvoker != null) {
                listenerInvoker.onExecutionFailed(e, context.toFinalExecutionInfo());
            }
            return Observable.error(e);
        }
    }

2.2 selectServer Выберите сервер

Как упоминалось ранее, ядром балансировки нагрузки является SelectServer:

C- LoadBalancerCommand
private Observable<Server> selectServer() {
    return Observable.create(new OnSubscribe<Server>() {
        @Override
        public void call(Subscriber<? super Server> next) {
            // 通过 loadBalancerContext 进行负载均衡处理
            Server server = loadBalancerContext.getServerFromLoadBalancer(loadBalancerURI, loadBalancerKey);   
            
            // Observable 模式的其他调用
            next.onNext(server);
            next.onCompleted();
        }
    });
}

2.3 Основная логика getServerFromLoadBalancer

Этот раздел логики относительно длинный, но основного кода не так много, мы оставляем только основные части.

C- LoadBalancerContext
public Server getServerFromLoadBalancer(@Nullable URI original, @Nullable Object loadBalancerKey) throws ClientException {
    
    // 属性准备
    String host = null;
    int port = -1;
    
    // 省略点一 : 会尝试从 URI original (http:///template/get?desc=order-server 中获取 Host 和 Post
    // 当然 ,正常模式下是获取不到的 , 所以这里省略 host != null 的情况 >>>
        
    // 获取 ILoadBalancer 对象 , 此处主要为 ZoneAwareLoadBalancer , 会调用父类 BaseLoadBalancer
    ILoadBalancer lb = getLoadBalancer();
    if (host == null) {
        if (lb != null){
            // 核心逻辑 : 选择 Server 对象 --> 详见 2.4 
            Server svc = lb.chooseServer(loadBalancerKey);
            host = svc.getHost();
            return svc;
        } else {
            // 省略负载均衡器不存在的逻辑
        }
    } else {
        // 省略 host 不为 null
    }
    
    // 最终构建一个 新 Server 对象返回
    return new Server(host, port);
}



// PS : 其中省略了很多判空抛出异常的逻辑 , 通常结构如下所示
if (host == null){
    throw new ClientException(ClientException.ErrorType.GENERAL,"....");
}


Feign_ILoadBalancer.png

2.4 Выберите сервер с помощью правила в ChooseServer

Вот точка использования окончательного алгоритма, который обрабатывается вызовом окончательной стратегии через Rule:

C- BaseLoadBalancer
public Server chooseServer(Object key) {
    if (counter == null) {
        // 用于跟踪某个事件发生频率的监视器类型
        counter = createCounter();
    }
    counter.increment();
    
    if (rule == null) {
        return null;
    } else {
        // 调用 Rule 返回具体的 Server
        return rule.choose(key);
        
        // 忽略 try-catch
    }
}

Система правил обеспечивает идеальную систему, посмотрите на структурную целостность системы Правило:

Feign_IRule.png

PS: здесь в основном используется PredicateBasedRule

2.5 Анализ случая использования PredicateBasedRule

В ней есть несколько основных логик:

  • lb.getAllServers(): получить все серверы
  • getEligibleServers(servers, loadBalancerKey) :
  • incrementAndGetModulo(eligible.size()) :

Step 1: выберите Select Server, чтобы увидеть, вот конкретный выбор через Predicate

C- PredicateBasedRule 
@Override
public Server choose(Object key) {
    ILoadBalancer lb = getLoadBalancer();
    // 注意此处的 lb.getAllServers() , 已经获取了所有的 Server 列表
    Optional<Server> server = getPredicate().chooseRoundRobinAfterFiltering(lb.getAllServers(), key);
    if (server.isPresent()) {
        return server.get();
    } else {
        return null;
    }       
}

Step 2: сервер выбирает основной метод

В этом методе сначала получите правильный список серверов (Step 3), после использования алгоритма выбора конкретного Сервера (Step 4)

C- PredicateBasedRule 
// 选择合适的 Service
public Optional<Server> chooseRoundRobinAfterFiltering(List<Server> servers, Object loadBalancerKey) {
	List<Server> eligible = getEligibleServers(servers, loadBalancerKey);
	if (eligible.size() == 0) {
            return Optional.absent();
	}
        // 核心语句 : 此处从 Server 中获取 Server
	return Optional.of(eligible.get(incrementAndGetModulo(eligible.size())));
} 

Step 3: получить правильный список серверов

C- PredicateBasedRule 
public List<Server> getEligibleServers(List<Server> servers, Object loadBalancerKey) {
    if (loadBalancerKey == null) {
        return ImmutableList.copyOf(Iterables.filter(servers, this.getServerOnlyPredicate())); 
    } else {
        List<Server> results = Lists.newArrayList();
        for (Server server: servers) {
            // 
            if (this.apply(new PredicateKey(loadBalancerKey, server))) {
                results.add(server);
            }
        }
        return results;            
    }
}

Step 4: Используйте алгоритм для выбора определенного сервера

C- PredicateBasedRule 
private int incrementAndGetModulo(int modulo) {
    for (;;) {
        int current = nextIndex.get();
        int next = (current + 1) % modulo;
        if (nextIndex.compareAndSet(current, next) && current < modulo){
            return current;
        }
    }
}

2.6 Выполнение ServerOperation (Функция:003)

Когда Submit выполняет серверную обработку (подробности см. -> UNDO:001), объект наблюдателя, подготовленный ранее, будет вызван обратно, и URL-адрес, который вам нужен, генерируется здесь через Сервер для завершения процесса балансировки нагрузки.

Здесь также использование наблюдателя, после завершения всей обработки последующий код логики выполняется по подписке >>

C- AbstractLoadBalancerAwareClient
new ServerOperation<T>() {
    @Override
    public Observable<T> call(Server server) {
        // 获取实际的 URL , 指向具体的 Server
        URI finalUri = reconstructURIWithServer(server, request.getUri());
        S requestForServer = (S) request.replaceUri(finalUri);
        try {
            return Observable.just(AbstractLoadBalancerAwareClient.this.execute(requestForServer, requestConfig));
        } catch (Exception e) {
            return Observable.error(e);
        }
    }
}

разное

Вопрос 1: Логика вызова наблюдателя и функции

Большой проблемой в балансировке нагрузки Feign является использование наблюдателей и синтаксического сахара Function.Этот режим обратного вызова может легко запутать всю логику:

  • Step 1: получить объект сервера в LoadBalancerCommand.
  • Step 2: Вторичная обработка сервера в функции: 002
  • Step 3: Выполнение прослушивателей и таймингов в Observable:002
  • Step 4: Вернитесь к функции: 003 в AbstractLoadBalancerAwareClient, чтобы сгенерировать URL-адрес и последующую логику.
  • Step 5: в Observable:003 onNext запускает прослушиватель и устанавливает сущность объекта
  • Step 6: Observable:003 в onCompleted подписывается на завершение запроса

Вопрос 2: Использование ExecutionListener

  • Функция:Слушатели, вызываемые балансировщиком нагрузки на разных этапах исполнения
  • использовать :Предусмотрены следующие методы
    • onExecutionStart : запускать
    • onStartWithServer : Вызывается, когда сервер выбран и запрос будет выполняться на сервере
    • onExceptionWithServer : На сервере есть исключение
    • onExecutionSuccess : выполнение успешно
    • onExecutionFailed : Не удалось выполнить
public void onExecutionStart(ExecutionContext<I> context) {
    for (ExecutionListener<I, O> listener : listeners) {
        if (!isListenerDisabled(listener)) {
            listener.onExecutionStart(context.getChildContext(listener));
        }
    }
}

Суммировать

В этой статье я стремлюсь прояснить логику балансировки нагрузки.Теперь она в основном соответствует требованиям.На самом деле есть еще над чем подумать во всем процессе.

Например:

  • Использование объектов multi-observer в Feign, хотя читабельность стала хуже, но бизнес-возможности сильно улучшились, как лучше спроектировать эту логику
  • Использование Feign Ribbon Rule, возможно ли настроить набор или расширить внешний алгоритм

Эти вещи будут проанализированы в последующих статьях, а мы подождем и увидим.