Это второй день моего участия в 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,"....");
}
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
}
}
Система правил обеспечивает идеальную систему, посмотрите на структурную целостность системы Правило:
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, возможно ли настроить набор или расширить внешний алгоритм
Эти вещи будут проанализированы в последующих статьях, а мы подождем и увидим.