Inventory Cloud: базовое руководство Nacos

Java
Inventory Cloud: базовое руководство Nacos

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

Введение

Немного гидрологии, хахахахаха~~~~~~~~~~~~

Цель статьи:

  • Разберитесь с направлением отладки функции создания Nacos
  • Разберите модульную систему компонентов Nacos

План статьи: Этот документ в основном состоит из следующих основных частей.

  • Обнаружение службы Nacos
  • Загрузка конфигурации Nacos
  • Проверка здоровья Nacos
  • Стратегия маршрутизации Nacos

PS: Источник ссылки на документациюофициальная документацияРекомендуется прочитать документ, чтобы быстро научиться.

2. Компиляция исходного кода

2.1 Компиляция исходного кода Nacos

// Step 1 : 下载源码
https://github.com/alibaba/nacos.git

// Step 2 : 编译 Nacos
mvn -Prelease-nacos -Dmaven.test.skip=true clean install -U  
    
// Step 3 :  运行 Server 文件
nacos_code\distribution\target    

2.2 Работа с исходным кодом Nacos

// Step 1 : 下载 Nacos 源码
https://github.com/alibaba/nacos.git

// Step 2 : IDEA 导入 Nacos
此处添加 SpringBoot 启动 , 启动的类为   com.alibaba.nacos.Nacos


// Step 3 : IDEA 修改启动参数
-Dnacos.standalone=true -Dnacos.home=C:\\nacos
    
-nacos.standalone=true : 单机启动
-Dnacos.home=C:\\nacos : 日志路径
    
    
    
// PS : Nacos Application
@SpringBootApplication(scanBasePackages = "com.alibaba.nacos")
@ServletComponentScan
@EnableScheduling
public class Nacos {
    
    public static void main(String[] args) {
        SpringApplication.run(Nacos.class, args);
    }
}

3. Исходный код модуля

<modules>
	<!-- 配置管理-->
	<module>config</module>
        <!-- Nacos 内核 -->
	<module>core</module>
	<!-- 服务发现 -->
	<module>naming</module>
        <!-- 地址服务器--> 
	<module>address</module>
	<!-- 单元测试 -->
	<module>test</module>
        <!-- 接口抽象 -->
	<module>api</module>
        <!-- 客户端 -->
	<module>client</module>
	<!-- 案例 -->
	<module>example</module>
        <!-- 公共工具 -->
	<module>common</module>
        <!-- Server 构建发布 -->
	<module>distribution</module>
        <!-- 控制台,图形界面模块 -->
	<module>console</module>
         <!-- 元数据管理-->
	<module>cmdb</module>
        <!-- TODO : 猜测是集成 istio 完成流量控制-->
	<module>istio</module>
        <!-- 一致性管理 -->
	<module>consistency</module>
        <!-- 权限控制 -->
	<module>auth</module>
        <!-- 系统信息管理 Env 读取 , conf 读取 -->
	<module>sys</module>
</modules>

3.1 Обнаружение и управление услугами Nacos (c0-c20)

Управление услугами Nacos в основном сосредоточено наNamingВ модуле здесь объединено с клиентом, чтобы увидеть, какова логика обнаружения и управления услугами и как ими управлять.

3.1.1 Консоль для получения списка услуг

Внешний интерфейс:

  • C-CatalogController # listDetail : получить список сведений об услуге.
  • C- CatalogController # instanceList: Список экземпляров специальных сервисов
  • C- CatalogController # serviceDetail : сведения об услуге

Ядро проходит в основном черезServiceManagerДля обработки взгляните на внутреннюю связанную логику:

通过三个接口不难发现 ,其最终的调用核心都是 ServiceManager 类

// 内部类
C01- ServiceManager
    // 内部类
    PC- UpdatedServiceProcessor
    PC- ServiceUpdater
    PSC- ServiceChecksum
    PC- EmptyServiceAutoClean
    PC- ServiceReporter
    PSC- ServiceKey : 
    // 核心参数
    F- Map<String, Map<String, Service>> serviceMap = new ConcurrentHashMap<>();
    F- LinkedBlockingDeque<ServiceKey> toBeUpdatedServicesQueue = new LinkedBlockingDeque<>(1024 * 1024);
    F- ConsistencyService consistencyService;
    // 常用方法
    M- init : 初始化方法
    M- chooseServiceMap : 通过空间名获取 Server 集合
    M- addUpdatedServiceToQueue : 更新 Server
    M- onChange : server 改变
    M- onDelete : server 删除


// 另外 , ServiceManager 还存在一个依赖 : ConsistencyService
C02- ConsistencyService : 一致性服务接口  -> PS:C02_01
    M- put :向集群提交一个数据
    M- remove :从集群删除一个数据
    M- get :从集群获取数据
    M- listen :监听集群中某个key的变化
    M- unlisten :删除对某个key的监听
    M- isAvailable :返回当前一致性状态是否可用


//总结 : 此处的逻辑很简单 , 就是对集合的 CURD 操作 , 核心特点有以下几个 : 
1- Delete 时 , 会调用依赖对象 ConsistencyService (DelegateConsistencyServiceImpl) , 用于处理一致性需求

PS:C02_01 Архитектура службы согласованности

nacos-ConsistencyService.png

Когда есть несколько сервисов, как они хранятся?

Nacos-server-Map.jpg

// 如上述图所示 , 相关的对象存放在 clusterMap 中

// 注意 , Service 有2个
com.alibaba.nacos.api.naming.pojo.Service
com.alibaba.nacos.naming.core.Service


// Service 的属性
C- Service
    I- com.alibaba.nacos.api.naming.pojo.Service
    F- Selector selector
    F- Map<String, Cluster> clusterMap = new HashMap<>();
    F- Boolean enabled
    F- Boolean resetWeight
    F- String token
    F- List<String> owners

3.1.2 Служба Nacos и класс управления Config

Nacos 中主要通过 NamingService 和 ConfigService 对服务和配置进行控制 , 其底层原理仍然为 :   , 这里来简单看一下 

ConfigService -> NacosConfigService
NamingService -> NacosNamingService


这2个类归属于 com.alibaba.nacos.client.naming 包

// PS : 注意 ,要使用 nacos-config-spring-boot-starter 和 nacos-discovery-spring-boot-starter 包

3.1.3 Выход и уничтожение Nacos

Когда зарегистрированная нами служба будет закрыта, Nacos выполнит выход приложения с сервера в конце жизненного цикла.

Шаг 1: Closeable останавливает связанные классы и останавливает проект.Когда мы нажимаем «Стоп», мы видим следующую серию журналов.

com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.beat.BeatReactor do shutdown stop
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.core.EventDispatcher do shutdown begin
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.core.EventDispatcher do shutdown stop
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.core.HostReactor do shutdown begin
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.core.PushReceiver do shutdown begin
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.core.PushReceiver do shutdown stop
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.backups.FailoverReactor do shutdown begin
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.backups.FailoverReactor do shutdown stop
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.core.HostReactor do shutdown stop
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.net.NamingProxy do shutdown begin
com.alibaba.nacos.client.naming          : [NamingHttpClientManager] Start destroying NacosRestTemplate
com.alibaba.nacos.client.naming          : [NamingHttpClientManager] Destruction of the end
com.alibaba.nacos.client.naming          : com.alibaba.nacos.client.naming.net.NamingProxy do shutdown stop


这里可以看到 , 这里进行了关闭操作 , 其中主要的 Close 操作都是基于 com.alibaba.nacos.common.lifecycle.Closeable 进行实现

 // PS : 此处的调用逻辑为 TODO 

Шаг 2. Выход из системы NacosServiceRegistry.

В дополнение к этому Closeable будет закрыт, он также выйдет из службы, здесь в основном NacosServiceRegistry

public void deregister(Registration registration) {
    if (StringUtils.isEmpty(registration.getServiceId())) {
        return;
    }

    NamingService namingService = namingService();
    String serviceId = registration.getServiceId();
    String group = nacosDiscoveryProperties.getGroup();

    try {
        namingService.deregisterInstance(serviceId, group, registration.getHost(),
					registration.getPort(), nacosDiscoveryProperties.getClusterName());
    } catch (Exception e) {
        // 省略 log
    }

}
    
// PS : 此处的原理为 继承了 ServiceRegistry , 实现销毁逻辑
C- AbstractAutoServiceRegistration

Система Nacos_Closeable

Nacos_Closeable.png

3.1.4 Процесс проверки работоспособности

Nacos обеспечивает проверку работоспособности сервисов в режиме реального времени, предотвращая запросы к неработоспособным хостам или экземплярам сервисов.Nacos поддерживает транспортный уровень (PING или TCP) и прикладной уровень (например, HTTP, MySQL, определяемый пользователем) проверки работоспособности.

Интерфейсы, связанные с проверкой работоспособности

  • Отправить пульс экземпляра (InstanceController): /nacos/v1/ns/instance/beat
  • Обновить состояние работоспособности экземпляра (HealthController): /nacos/v1/ns/health/instance

Клиент инициирует сердцебиение

Сервер будет периодически инициировать операцию пульса для вызова 2 интерфейсов:

C- BeatReactor
PC- BeatTask : 内部类
      
// 其中会有2个步骤 :

// Step 1 : BeatReactor # addBeatInfo 中添加定时任务
 public void addBeatInfo(String serviceName, BeatInfo beatInfo) {
        NAMING_LOGGER.info("[BEAT] adding beat: {} to beat map.", beatInfo);
        String key = buildKey(serviceName, beatInfo.getIp(), beatInfo.getPort());
        BeatInfo existBeat = null;
        //fix #1733
        if ((existBeat = dom2Beat.remove(key)) != null) {
            existBeat.setStopped(true);
        }
        dom2Beat.put(key, beatInfo);
        // 添加心跳信息
        executorService.schedule(new BeatTask(beatInfo), beatInfo.getPeriod(), TimeUnit.MILLISECONDS);
        MetricsMonitor.getDom2BeatSizeMonitor().set(dom2Beat.size());
}

// Step 2 :  BeatTask 中调用 Server 接口进行心跳操作
JsonNode result = serverProxy.sendBeat(beatInfo, BeatReactor.this.lightBeatEnabled);


// PS : 心跳的间隔默认是 5秒 ( com.alibaba.nacos.api.common.Constants)
public static final long DEFAULT_HEART_BEAT_TIMEOUT = TimeUnit.SECONDS.toMillis(15);
public static final long DEFAULT_IP_DELETE_TIMEOUT = TimeUnit.SECONDS.toMillis(30);
public static final long DEFAULT_HEART_BEAT_INTERVAL = TimeUnit.SECONDS.toMillis(5);

Сервер обнаруживает сердцебиение

// 服务端同样会对实例进行检测 , 核心类为 ClientBeatCheckTask

// Step 1 : ClientBeatCheckTask 的创建
C- Service
    ?- 在 Service 初始化时 ,即开始了 Task 任务
    
    
public void init() {
    HealthCheckReactor.scheduleCheck(clientBeatCheckTask);
    //............
}

// 这里也可以看到默认时间的设置
public static void scheduleCheck(ClientBeatCheckTask task) {
        futureMap.computeIfAbsent(task.taskKey(),
                k -> GlobalExecutor.scheduleNamingHealth(task, 5000, 5000, TimeUnit.MILLISECONDS));
}



// Step 2 : ClientBeatCheckTask 的运行
C- ClientBeatCheckTask
    ?- 检查并更新临时实例的状态,如果它们已经过期则删除它们。
    

public void run() {
    //...... 省略
    
    // Step 1 : 获取所有实例
    List<Instance> instances = service.allIPs(true);
    
    for (Instance instance : instances) {
        // 如果时间大于心跳超时时间 , 则修改健康状态
        if (System.currentTimeMillis() - instance.getLastBeat() > instance.getInstanceHeartBeatTimeOut()) {
            if (!instance.isMarked()) {
                if (instance.isHealthy()) {
                    instance.setHealthy(false);
                    getPushService().serviceChanged(service);
                    ApplicationUtils.publishEvent(new InstanceHeartbeatTimeoutEvent(this, instance));
                }
            }
        }
    }
            
    if (!getGlobalConfig().isExpireInstance()) {
        return;
    }

    for (Instance instance : instances) { 
        if (instance.isMarked()) {
            continue;
        }
        // 如果心跳大于删除时间 , 则删除实例
        if (System.currentTimeMillis() - instance.getLastBeat() > instance.getIpDeleteTimeout()) {
            deleteIp(instance);
        }
    }      
}
    
// PS : 第一次修改的是健康状态 , 后面才修改的实例数
    

Основная логика проверки работоспособности ACK и одновременная отправка информации

Непреднамеренно обнаружен механизм ACK, это делается через инцидент.Этот режим в основном предназначен для отправки обновлений через порт UDP, когда сервер меняется на сервере.

PS: когда клиент запрашивает экземпляр службы, если указан порт udp, сервер создаст udpClient.

for (Instance instance : service.allIPs(Lists.newArrayList(clusterName))) {
    if (instance.getIp().equals(ip) && instance.getPort() == port) {
        instance.setHealthy(valid);
        // 发布事件 ServiceChangeEvent
        pushService.serviceChanged(service);
        break;
    }
} 
                                                                     
// ServiceChangeEvent 事件的处理
C- PushService                                                             
public void onApplicationEvent(ServiceChangeEvent event) {
        Service service = event.getService();
        String serviceName = service.getName();
        String namespaceId = service.getNamespaceId();
        
        Future future = GlobalExecutor.scheduleUdpSender(() -> {
            try {
                // 通过 ServerName 和 命名空间 获取PushClient 集合
                ConcurrentMap<String, PushClient> clients = clientMap
                        .get(UtilsAndCommons.assembleFullServiceName(namespaceId, serviceName));
                if (MapUtils.isEmpty(clients)) {
                    return;
                }
                
                Map<String, Object> cache = new HashMap<>(16);
                long lastRefTime = System.nanoTime();
                
                // 循环所有的 PushClient
                for (PushClient client : clients.values()) {
                    if (client.zombie()) {
                        clients.remove(client.toString());
                        continue;
                    }
                    
                    Receiver.AckEntry ackEntry;
                    // 获取缓存 key ,并且从缓存中获取实体数据
                    String key = getPushCacheKey(serviceName, client.getIp(), client.getAgent());
                    byte[] compressData = null;
                    Map<String, Object> data = null;
                    if (switchDomain.getDefaultPushCacheMillis() >= 20000 && cache.containsKey(key)) {
                        org.javatuples.Pair pair = (org.javatuples.Pair) cache.get(key);
                        compressData = (byte[]) (pair.getValue0());
                        data = (Map<String, Object>) pair.getValue1();
                        
                    }
                    
                    // 构建 ACK 实体类
                    if (compressData != null) {
                        ackEntry = prepareAckEntry(client, compressData, data, lastRefTime);
                    } else {
                        ackEntry = prepareAckEntry(client, prepareHostsData(client), lastRefTime);
                        if (ackEntry != null) {
                            cache.put(key, new org.javatuples.Pair<>(ackEntry.origin.getData(), ackEntry.data));
                        }
                    }
                    
                    // UDP ACK 校验 , 同时推送 ACK 实体
                    udpPush(ackEntry);
                }
            } catch (Exception e) {
                Loggers.PUSH.error("[NACOS-PUSH] failed to push serviceName: {} to client, error: {}", serviceName, e);
                
            } finally {
                futureMap.remove(UtilsAndCommons.assembleFullServiceName(namespaceId, serviceName));
            }
            
        }, 1000, TimeUnit.MILLISECONDS);
        
        futureMap.put(UtilsAndCommons.assembleFullServiceName(namespaceId, serviceName), future);
        
    } 
    
private static Receiver.AckEntry udpPush(Receiver.AckEntry ackEntry) {
        if (ackEntry == null) {
            return null;
        }
        
        if (ackEntry.getRetryTimes() > MAX_RETRY_TIMES) {
            ackMap.remove(ackEntry.key);
            udpSendTimeMap.remove(ackEntry.key);
            failedPush += 1;
            return ackEntry;
        }
        
        try {
            if (!ackMap.containsKey(ackEntry.key)) {
                totalPush++;
            }
            ackMap.put(ackEntry.key, ackEntry);
            udpSendTimeMap.put(ackEntry.key, System.currentTimeMillis());
            
            //  Socket Send 
            udpSocket.send(ackEntry.origin);
            
            ackEntry.increaseRetryTime();
            
            GlobalExecutor.scheduleRetransmitter(new Retransmitter(ackEntry),
                    TimeUnit.NANOSECONDS.toMillis(ACK_TIMEOUT_NANOS), TimeUnit.MILLISECONDS);
            
            return ackEntry;
        } catch (Exception e) {
            ackMap.remove(ackEntry.key);
            udpSendTimeMap.remove(ackEntry.key);
            failedPush += 1;
            
            return null;
        }
}     




Использование порогов проверки работоспособности

При настройке услугиЧисло с плавающей запятой 0-1 может быть настроено, определите порог проверки работоспособности, класс, соответствующий порогу, — com.alibaba.nacos.api.naming.pojo.Service

C- Service
    F- name : 服务名
    F- protectThreshold : 健康阈值
    F- appName : 应用名 
    F- groupName : 组名
    F- metadata : 元数据  
    
// 阈值的使用                                                                  
C- InstanceController
	M- doSrvIpxt :

// 核心逻辑              
                                                                     
double threshold = service.getProtectThreshold();
// IPMap 中可用的健康实例数/服务总数的比例 如果小于阈值 , 则达到保护阈值                                                          
if ((float) ipMap.get(Boolean.TRUE).size() / srvedIPs.size() <= threshold) {
      
	if (isCheck) {
		result.put("reachProtectThreshold", true);
	}
	ipMap.get(Boolean.TRUE).addAll(ipMap.get(Boolean.FALSE));
	ipMap.get(Boolean.FALSE).clear();
}
                                                                     
                                                                     
PS : 这里联想后面 , Client 做 Balancer 时 , 获取的 Server 实际上就是全部健康的实例了                        
    

Но какова цель порога?

Если предположить, что в инстансе большое количество исключений, то финальное давление будет на эти исправные инстансы.В это время может произойти цепная реакция.

Чтобы избежать таких ситуаций, все экземпляры возвращаются при достижении порога работоспособности.

Это цель предложения **ipMap.get(Boolean.TRUE).addAll(ipMap.get(Boolean.FALSE));**.

Хотя клиентская сторона может столкнуться с аномальными экземплярами, это может предотвратить сбой всей системы.

/nacos/v1/ns/health/параметр экземпляра

имя Типы Являются обязательными описывать
namespaceId нить нет Идентификатор пространства имен
serviceName нить да наименование услуги
groupName нить нет имя группы
clusterName нить нет имя кластера
ip нить да IP-адрес экземпляра службы
port int да порт экземпляра службы
healthy boolean да Это здорово

Работа с весами

Вес в основном настраивается в объекте экземпляра.

C- Instance
    M- instanceId
    M- ip
    M- port
    M- weight : 权重
    M- healthy : 健康情况
    M- enabled
    M- ephemeral
    M- clusterName
    M- serviceName
    M- metadata

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

Ссылка на оригинал @blog.CSDN.net/Крис прерывает/искусство…

public class NacosWeightLoadBalanceRule extends AbstractLoadBalancerRule {
 
  @Override
  public void initWithNiwsConfig(IClientConfig clientConfig) {}
 
  @Resource private NacosDiscoveryProperties nacosDiscoveryProperties;
 
  @Override
  public Server choose(Object key) {
    // 1.获取服务的名称
    BaseLoadBalancer loadBalancer = (BaseLoadBalancer) this.getLoadBalancer();
    String serverName = loadBalancer.getName();
    // 2.此时Nacos Client会自动实现基于权重的负载均衡算法
    NamingService namingService = nacosDiscoveryProperties.namingServiceInstance();
    try {
      Instance instance = namingService.selectOneHealthyInstance(serverName);
      return new NacosServer(instance);
    } catch (NacosException e) {
      e.printStackTrace();
    }
    return null;
  }

@Bean
public IRule getLoadBalancerRule(){
	return new NacosWeightLoadBalancerRule();
}


// PS : 个人以为 , 权重是给 Client 端自行处理的

Другие пункты

// Nacos 的默认值

@Value("${nacos.naming.empty-service.auto-clean:false}")
private boolean emptyServiceAutoClean;

@Value("${nacos.naming.empty-service.clean.initial-delay-ms:60000}")
private int cleanEmptyServiceDelay;

@Value("${nacos.naming.empty-service.clean.period-time-ms:20000}")
private int cleanEmptyServicePeriod;

3.2 Процесс настройки Nacos C30-C60

3.2.1 Управление конфигурацией Nacos

// 获取配置 , 这里主要有几个步骤 :
C30- NacosConfigService
	M30_01- getConfigInner(String tenant, String dataId, String group, long timeoutMs)
        - 构建 ConfigResponse , 为其设置 dataId , tenant , group
        1- 调用 LocalConfigInfoProcessor.getFailover 优先使用本地配置
        2- 调用 ClientWorker , 获取远程配置 -> PS:M30_01_01
        3- 仍然没有 , LocalConfigInfoProcesso.getSnapshot 获取快照   
        End- configFilterChainManager 进行 Filter 链处理
            
// PS:M30_01_01 ClientWorker 的处理
ClientWorker 中进行了远程服务的请求 , 核心代码 : 
agent.httpGet(Constants.CONFIG_CONTROLLER_PATH, null, params, agent.getEncode(), readTimeout);

// 可以看到 , 这里并没有太负载的逻辑 , 仍然是 Rest 请求 : PS 看官网的消息 , 2.0 会采用长连接 , 这里应该会有变动
- Constants.CONFIG_CONTROLLER_PATH : /v1/cs/configs


Основная логика LocalConfigInfoProcessor

// 这里本地配置是指本地 File 文件 , 这里通过源码推断一下使用的方式 : 
C31- LocalConfigInfoProcessor
    M31_01- getFailover
    	- 获取 localPath -> PS:M31_01_01
	M32_02- saveSnapshot : 获取成功后 , 会保存快照
        ?- 保存路径 : 省略\nacos\config\fixed-127.0.0.1_8848_nacos\snapshot\one1\test1
        
            
// PS:M31_01_01 localPath 参数
C:\Users\10169\nacos\config\fixed-127.0.0.1_8848_nacos\data\config-data\one1\test1    

Pro 1: Какие знания можно увидеть в этом исходном коде?

Поскольку существует локальный файл, означает ли это, что я могу сначала использовать локальную конфигурацию, изменив этот путь? Тест проходит успешно после изменения следующего пути здесьОпустить \nacos\config\fixed-127.0.0.1_8848_nacos\data\config-data\one1\test1

PS : 除了这个路径 , SpringBoot 运行在配置文件中直接配置路径 -> 
spring:
  cloud:
    config:
      # 相同配置,本地优先
      override-none: true

Pro 2: Использование фильтров

Как вы можете видеть выше, при обработке конфигурации есть обработка фильтра по умолчанию.
configFilterChainManager.doFilter(null, cr);

    // 依照这个逻辑 , 是可以进行更多配置的
C32- ConfigFilterChainManager
    ?- 其中允许自定义添加 Filter
    M32_01- addFilter
    M32_02- doFilter
    
    


TODO : Как вживить фильтр здесь нужно доработать, я не нашел интерфейс для добавления фильтра, странно....

3.2.2 Аварийное восстановление конфигурации Nacos

Nacos LocalConfigInfoProcessor предоставляет функции аварийного восстановления двумя способами:Локальная конфигурация и обработка моментальных снимков

локальная конфигурация

Как упоминалось выше, измените указанный путь для достижения

Снимок конфигурации

Клиентский SDK Nacos создаст моментальный снимок конфигурации локально. Когда клиент не может подключиться к серверу Nacos, снимок конфигурации можно использовать для отображения общих возможностей аварийного восстановления системы. Снимки конфигурации аналогичны локальным коммитам в Git, а также аналогичны кешам.Они будут обновляться в соответствующее время, но не существует концепции срока действия кеша.

3.2.3 Обработка динамической конфигурации Nacos

Динамическая конфигурация в основном относится к мониторингу изменений конфигурации:

перевал НакосДолгий опрос для обнаружения изменений конфигурации, соответствующий базовый классLongPollingRunnable # checkUpdateDataIds, взгляните на отладку здесь >>>>


    
 class LongPollingRunnable implements Runnable {
        
        private final int taskId;
        
        public LongPollingRunnable(int taskId) {
            this.taskId = taskId;
        }
        
        @Override
        public void run() {
            
             // .... 核心语句
             List<String> changedGroupKeys = checkUpdateDataIds(cacheDatas, inInitializingCacheList);
        }
     
 }     
            
List<String> checkUpdateDataIds(List<CacheData> cacheDatas, List<String> inInitializingCacheList) throws Exception {
        StringBuilder sb = new StringBuilder();
        for (CacheData cacheData : cacheDatas) {
            if (!cacheData.isUseLocalConfigInfo()) {
                sb.append(cacheData.dataId).append(WORD_SEPARATOR);
                sb.append(cacheData.group).append(WORD_SEPARATOR);
                if (StringUtils.isBlank(cacheData.tenant)) {
                    sb.append(cacheData.getMd5()).append(LINE_SEPARATOR);
                } else {
                    sb.append(cacheData.getMd5()).append(WORD_SEPARATOR);
                    sb.append(cacheData.getTenant()).append(LINE_SEPARATOR);
                }
                if (cacheData.isInitializing()) {
                    // It updates when cacheData occours in cacheMap by first time.
                    inInitializingCacheList
                            .add(GroupKey.getKeyTenant(cacheData.dataId, cacheData.group, cacheData.tenant));
                }
            }
        }
        boolean isInitializingCacheList = !inInitializingCacheList.isEmpty();
        return checkUpdateConfigStr(sb.toString(), isInitializingCacheList);
} 


List<String> checkUpdateConfigStr(String probeUpdateString, boolean isInitializingCacheList) throws Exception {
        
        Map<String, String> params = new HashMap<String, String>(2);
        params.put(Constants.PROBE_MODIFY_REQUEST, probeUpdateString);
        Map<String, String> headers = new HashMap<String, String>(2);
    	// 长轮询方式
        headers.put("Long-Pulling-Timeout", "" + timeout);
        
        // told server do not hang me up if new initializing cacheData added in
        if (isInitializingCacheList) {
            headers.put("Long-Pulling-Timeout-No-Hangup", "true");
        }
        
        if (StringUtils.isBlank(probeUpdateString)) {
            return Collections.emptyList();
        }
        
        try {
            // In order to prevent the server from handling the delay of the client's long task,
            // increase the client's read timeout to avoid this problem.
            
            long readTimeoutMs = timeout + (long) Math.round(timeout >> 1);
            // /v1/cs/configs/listener
            HttpRestResult<String> result = agent
                    .httpPost(Constants.CONFIG_CONTROLLER_PATH + "/listener", headers, params, agent.getEncode(),
                            readTimeoutMs);
            
            if (result.ok()) {
                setHealthServer(true);
                return parseUpdateDataIdResponse(result.getData());
            } else {
                setHealthServer(false);
            }
        } catch (Exception e) {
            setHealthServer(false);
            throw e;
        }
        return Collections.emptyList();
}


// 对应的 Controller 为 ConfigController
    public void listener(HttpServletRequest request, HttpServletResponse response)
            throws ServletException, IOException {
       
        // .............
        
        Map<String, String> clientMd5Map;
        try {
            clientMd5Map = MD5Util.getClientMd5Map(probeModify);
        } catch (Throwable e) {
            throw new IllegalArgumentException("invalid probeModify");
        }
        
        // do long-polling
        inner.doPollingConfig(request, response, clientMd5Map, probeModify.length());
    }

// 对应轮询的接口
public String doPollingConfig(HttpServletRequest request, HttpServletResponse response,
            Map<String, String> clientMd5Map, int probeRequestSize) throws IOException {
        
        // Long polling.
        if (LongPollingService.isSupportLongPolling(request)) {
            longPollingService.addLongPollingClient(request, response, clientMd5Map, probeRequestSize);
            return HttpServletResponse.SC_OK + "";
        }
        
        // Compatible with short polling logic.
        List<String> changedGroups = MD5Util.compareMd5(request, response, clientMd5Map);
        
        // Compatible with short polling result.
        String oldResult = MD5Util.compareMd5OldResult(changedGroups);
        String newResult = MD5Util.compareMd5ResultString(changedGroups);
        
        String version = request.getHeader(Constants.CLIENT_VERSION_HEADER);
        if (version == null) {
            version = "2.0.0";
        }
        int versionNum = Protocol.getVersionNumber(version);
        
        // Before 2.0.4 version, return value is put into header.
        if (versionNum < START_LONG_POLLING_VERSION_NUM) {
            response.addHeader(Constants.PROBE_MODIFY_RESPONSE, oldResult);
            response.addHeader(Constants.PROBE_MODIFY_RESPONSE_NEW, newResult);
        } else {
            request.setAttribute("content", newResult);
        }
        
        Loggers.AUTH.info("new content:" + newResult);
        
        // Disable cache.
        response.setHeader("Pragma", "no-cache");
        response.setDateHeader("Expires", 0);
        response.setHeader("Cache-Control", "no-cache,no-store");
        response.setStatus(HttpServletResponse.SC_OK);
        return HttpServletResponse.SC_OK + "";
}


// LongPollingService

这里想看详情推荐看这一篇 https://www.jianshu.com/p/acb9b1093a54

3.2.4 Метаданные Nacos

Давайте посмотрим, что такое метаданные Nacos?

Данные Nacos (такие как конфигурация и обслуживание) информация описания,Например, версии службы, веса, стратегии ликвидации последствий стихийных бедствий, стратегии балансировки нагрузки, конфигурация аутентификации, различные настраиваемые метки (Label)., по сфере действия он делится на метаинформацию об уровне обслуживания, метаинформацию кластера и метаинформацию экземпляра.

// 方式一 : 配置服务的时候配置
spring:
  application:
    name: nacos-config-server
  cloud:
    nacos:
      discovery:
        server-addr: 127.0.0.1:8848
        metadata:
          version: v1
              
// 方式二 : 页面直接配置
              
// 方式三 :  API 调用 , 通过 Delete , Update 等请求方式决定类型
批量更新实例元数据 (InstanceController) : /nacos/v1/ns/instance/metadata/batch  
批量删除实例元数据 (InstanceController) : /nacos/v1/ns/instance/metadata/batch 

image-20210527135529407.png

3.3 Балансировка нагрузки Nacos

Балансировка нагрузки Nacos принадлежитСлужба динамического DNS

Nacos_service_list.jpg

Служба динамической DNS поддерживает взвешенную маршрутизацию, упрощая реализацию балансировки нагрузки среднего уровня, более гибких политик маршрутизации, управления трафиком и простых служб разрешения DNS для внутренней сети центра обработки данных. Динамические службы DNS также упрощают реализацию обнаружения служб на основе протокола DNS, чтобы исключить риск привязки к частным API обнаружения служб поставщика.

Основываясь на знаниях, связанных с Feign, мы знаем, что,Обработка правила балансировки выполняется совместно с BaseLoadBalancer., давайте проанализируем, какой структурой они связаны

Шаг 1: Фейн звонит Нако.

Начинаем с BaseLoadBalancer, выполняем Debug-обработку и доходим до PredicateBasedRule

C- PredicateBasedRule
    M- choose(Object key)
    	-  Optional<Server> server = getPredicate().chooseRoundRobinAfterFiltering(lb.getAllServers(), key);
            ?- 此处可以看到 , 其中有一个 lb.getAllServers 的操作  , 此处的 lb 为 DynamicServerListLoadBalancer
            
// PS : getAllServers , 此处可以看到 ,其中的 Servers 已经全部放在 List 中了         
public List<Server> getAllServers() {
	return Collections.unmodifiableList(allServerList);
}


// 跟踪一下放入的逻辑 , 放入逻辑的起点是ILoadBalancer Bean 的加载 , 其调用链为 : 
C- RibbonClientConfiguration # ribbonLoadBalancer : 构建一个 ILoadBalancer
C- ZoneAwareLoadBalancer : 进入 ZoneAwareLoadBalancer 构造函数
C- DynamicServerListLoadBalancer : 进入 构造函数
C- DynamicServerListLoadBalancer # restOfInit : init 操作
C- DynamicServerListLoadBalancer # updateListOfServers : 更新 Server 列表主流程 , 此处第一次获取相关的 Server List , 后续Debug 第一节点
C- DynamicServerListLoadBalancer # updateAllServerList : 设置 ServerList 
    
    
public void updateListOfServers() {
    List<T> servers = new ArrayList<T>();
    if (serverListImpl != null) {
        servers = serverListImpl.getUpdatedListOfServers();
        if (filter != null) {
            servers = filter.getFilteredListOfServers(servers);
        }
    }
    updateAllServerList(servers);
}  


可以看到 , 其中有2个获取 Server 的逻辑方法 , 在这里看一下家族体系 , 就清楚了
serverListImpl.getUpdatedListOfServers();
filter.getFilteredListOfServers(servers);
    

// 下述图片中就很清楚了 , 存在一个实现类 NacosServerList 实现类 , 从Nacos 中获取服务列表 
private List<NacosServer> getServers() {
    try {
        String group = discoveryProperties.getGroup();
        List<Instance> instances = discoveryProperties.namingServiceInstance()
					.selectInstances(serviceId, group, true);
        return instancesToServerList(instances);
    } catch (Exception e) {
        throw new IllegalStateException(....);
    }
}


// 负载均衡策略 
负载均衡策略是基于 Balance  

Nacos_ServerList.png

3.4 Обработка кластеров

Использование кластера Nacos

Кластер Nacos относительно прост в использовании, вам нужно только настроить соответствующую служебную информацию в /conf/cluster.conf >>>>


#it is ip
#example
127.0.0.1:8848
127.0.0.1:8849
127.0.0.1:8850

Отслеживание источника кластера NACOS

Давайте посмотрим на уровень исходного кода, как обрабатывается эта логика?

Основные классы обработки находятся по адресу com.alibaba.nacos.core.cluster.


// Step 1 : 获取配置的方式
C- EnvUtil 
public static String getClusterConfFilePath() {
	return Paths.get(getNacosHome(), "conf", "cluster.conf").toString();
}    

// 读取 Cluster 配置
public static List<String> readClusterConf() throws IOException {
    try (Reader reader = new InputStreamReader(new FileInputStream(new File(getClusterConfFilePath())),
                StandardCharsets.UTF_8)) {
        return analyzeClusterConf(reader);
    } catch (FileNotFoundException ignore) {
        List<String> tmp = new ArrayList<>();
        String clusters = EnvUtil.getMemberList();
        if (StringUtils.isNotBlank(clusters)) {
            String[] details = clusters.split(",");
            for (String item : details) {
                tmp.add(item.trim());
            }
        }
        return tmp;
    }
}


// Step 2 : Cluster 的使用 
AbstractMemberLookup
    
// 主要使用集中在 ServerManager 中
C- ServerManager 
	F- ServerMemberManager memberManager;  

C- ServerMemberManager : Nacos中的集群节点管理
	M- init : 集群节点管理器初始化
	M- getSelf : 获取本地节点信息
	M- getmemberaddressinfo : 获取正常成员节点的地址信息
	M- allMembers : 获取集群成员节点的列表
	M- update : 更新目标节点信息
	M- isUnHealth : 目标节点是否健康
	M- initAndStartLookup : 初始化寻址模式

// TODO : 其他方法就省略了 , 后期准备进行相关的性能分析 , 集群的源码梳理预计放在那一部分分析

Суммировать

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

Кроме того, выходит Nacos 2.0. См. документ с использованием метода длинного соединения Socket. Если будет возможность в будущем, сравните разницу между ними.

приложение

# Приложение 1: Сервер вызова вручную

package com.alibaba.nacos.discovery.service;

import com.alibaba.nacos.api.NacosFactory;
import com.alibaba.nacos.api.config.ConfigService;
import com.alibaba.nacos.api.exception.NacosException;
import com.alibaba.nacos.api.naming.NamingService;
import com.alibaba.nacos.api.naming.pojo.Cluster;
import com.alibaba.nacos.api.naming.pojo.Instance;
import com.alibaba.nacos.api.naming.pojo.Service;
import com.alibaba.nacos.api.naming.pojo.healthcheck.AbstractHealthChecker;
import netscape.javascript.JSObject;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.stereotype.Component;

import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Properties;

/**
 * @Classname NacosClientService
 * @Description TODO
 * @Date 2021/5/26
 * @Created by zengzg
 */
@Component
public class NacosClientNodesService implements ApplicationRunner {

    private Logger logger = LoggerFactory.getLogger(this.getClass());

    private NamingService namingService;

    @Value("${spring.cloud.nacos.config.server-addr}")
    private String serverAddr;

    @Override
    public void run(ApplicationArguments args) throws Exception {
        Properties properties = new Properties();
        properties.put("serverAddr", serverAddr);
        namingService = NacosFactory.createNamingService(properties);
    }

    /**
     * 获取 Nacos Config
     * 参数格式
     * {
     *     "instanceId": "192.168.0.97#9083#DEFAULT#DEFAULT_GROUP@@nacos-user-server",
     *      "ip": "192.168.0.97",
     *      "port": 9083,
     *      "weight": 1, 可以通过权重决定使用的 Server
     *      "healthy": true,
     *      "enabled": true,
     *      "ephemeral": true,
     *      "clusterName": "DEFAULT",
     *      "serviceName": "DEFAULT_GROUP@@nacos-user-server",
     *      "metadata": {
     *           "preserved.register.source": "SPRING_CLOUD"
     *      },
     *      "ipDeleteTimeout": 30000,
     *      "instanceHeartBeatInterval": 5000,
     *      "instanceHeartBeatTimeOut": 15000
     * }
     *
     * @param serviceName
     * @return
     */
    public List<Instance> get(String serviceName) {
        List<Instance> content = new LinkedList<Instance>();
        try {
            content = namingService.getAllInstances(serviceName);
            logger.info("------> 获取 Config serviceName [{}]  <-------", serviceName);
        } catch (NacosException e) {
            logger.error("E----> error :{} -- content :{}", e.getClass(), e.getMessage());
            e.printStackTrace();
        }

        return content;


    }

    /**
     * 创建 Nacos Config
     *
     * @param serviceName
     * @param ip
     * @param port
     */
    public void createOrUpdate(String serviceName, String ip, Integer port) {
        try {
            logger.info("------> 创建 Config GroupID [{}] -- DataID [{}] Success ,The value :[{}] <-------", serviceName, ip, port);
            namingService.registerInstance(serviceName, ip, port, "TEST1");
        } catch (NacosException e) {
            logger.error("E----> error :{} -- content :{}", e.getClass(), e.getMessage());
            e.printStackTrace();
        }
    }


    /**
     * 移除 Nacos Config
     *
     * @param serviceName
     * @param ip
     */
    public void delete(String serviceName, String ip, Integer port) {
        try {
            namingService.deregisterInstance(serviceName, ip, port, "DEFAULT");
            logger.info("------> 删除 Config GroupID [{}] -- DataID [{}] Success  <-------", serviceName, ip);
        } catch (NacosException e) {
            logger.error("E----> error :{} -- content :{}", e.getClass(), e.getMessage());
            e.printStackTrace();
        }
    }
}



# Приложение II: Ручной вызов Config

public class NacosClientConfigService implements ApplicationRunner {

    private Logger logger = LoggerFactory.getLogger(this.getClass());

    private ConfigService configService;

    @Value("${spring.cloud.nacos.config.server-addr}")
    private String serverAddr;

    @Override
    public void run(ApplicationArguments args) throws Exception {
        Properties properties = new Properties();
        properties.put("serverAddr", serverAddr);
        configService = NacosFactory.createConfigService(properties);
    }


    /**
     * 获取 Nacos Config
     *
     * @param dataId
     * @param groupId
     * @return
     */
    public String get(String dataId, String groupId) {
        String content = "";
        try {
            content = configService.getConfig(dataId, groupId, 5000);
            logger.info("------> 获取 Config GroupID [{}] -- DataID [{}] Success ,The value :[{}] <-------", dataId, groupId, content);

            configService.addListener(dataId, groupId, new ConfigListener());

        } catch (NacosException e) {
            logger.error("E----> error :{} -- content :{}", e.getClass(), e.getMessage());
            e.printStackTrace();
        }

        return content;


    }

    /**
     * 创建 Nacos Config
     *
     * @param dataId
     * @param groupId
     * @param content
     */
    public void createOrUpdate(String dataId, String groupId, String content) {
        try {
            logger.info("------> 创建 Config GroupID [{}] -- DataID [{}] Success ,The value :[{}] <-------", dataId, groupId, content);
            configService.publishConfig(dataId, groupId, content);
        } catch (NacosException e) {
            logger.error("E----> error :{} -- content :{}", e.getClass(), e.getMessage());
            e.printStackTrace();
        }
    }


    /**
     * 移除 Nacos Config
     *
     * @param dataId
     * @param groupId
     */
    public void delete(String dataId, String groupId) {
        try {
            configService.removeConfig(dataId, groupId);
            logger.info("------> 删除 Config GroupID [{}] -- DataID [{}] Success  <-------", dataId, groupId);

            configService.removeListener(dataId, groupId, null);

        } catch (NacosException e) {
            logger.error("E----> error :{} -- content :{}", e.getClass(), e.getMessage());
            e.printStackTrace();
        }
    }


}


# Приложение 3: Использование NacosInjected

может пройти

<!-- 使用 Nacos Inject -->
<dependency>
    <groupId>com.alibaba.boot</groupId>
    <artifactId>nacos-config-spring-boot-starter</artifactId>
    <version>0.2.7</version>
</dependency>
<dependency>
    <groupId>com.alibaba.boot</groupId>
    <artifactId>nacos-discovery-spring-boot-starter</artifactId>
    <version>0.2.7</version>
</dependency>


@NacosInjected
private ConfigService configService;

    @NacosInjected
    private NamingService namingService;

# Приложение 4: Официальная схема архитектуры Naocs (транспорт)

Вот чистое обращение, можно посмотреть официальную документацию @ что cos.IO/this-capable/docs/…

Функциональная схема: image.png

  • Управление сервисом: реализовать сервис CRUD, доменное имя CRUD, проверку состояния сервиса, управление весом сервисов и другие функции.
  • Управление конфигурацией: реализовать управление конфигурацией CRUD, управление версиями, управление оттенками серого, управление мониторингом, push-трек, агрегирование данных и другие функции.
  • Управление метаданными: предоставление метаданных CURD и возможности маркировки
  • Механизм подключаемых модулей: реализовать возможность разделения и объединения трех модулей, а также реализовать механизм точки расширения SPI.
  • Механизм событий: реализовать асинхронное уведомление о событии, асинхронное уведомление об изменении данных SDK и другую логику.
  • Модуль журнала: управление классификацией журнала, уровнем журнала, переносимостью журнала (особенно во избежание конфликтов), форматом журнала, кодом исключения + справочной документацией.
  • Механизм обратного вызова: SDK уведомляет данные и вызывает обработку пользователя в унифицированном режиме. Интерфейсы и структуры данных должны быть расширяемыми
  • Режим адресации: решить различные режимы адресации, такие как IP, доменное имя, сервер имен, широковещательная рассылка и т. д., которые должны быть расширяемыми.
  • Push-канал: решить проблемы с производительностью push-уведомлений между сервером и хранилищем, между сервером, сервером и SDK.
  • Управление емкостью: управляйте емкостью каждого арендатора и группы, чтобы предотвратить перезапись хранилища и повлиять на доступность службы.
  • Управление трафиком: управление частотой запросов, количеством длинных ссылок, размером пакета и управлением потоком запросов в соответствии с арендаторами, группами и другими параметрами.
  • Механизм кэширования: каталог аварийного восстановления, локальный кэш, механизм кэширования сервера. Необходимы инструменты для использования каталога аварийного восстановления
  • Режим запуска: в соответствии с автономным режимом, режимом конфигурации, сервисным режимом, режимом DNS или всеми режимами, запуск различных программ + пользовательский интерфейс
  • Протокол согласованности: обращение к разным данным, к разным требованиям согласованности, к разным механизмам согласованности.
  • Модуль хранения: решение проблемы сохранения и непостоянства данных, а также решение проблемы фрагментации данных.
  • Сервер имен: решить проблему маршрутизации из пространства имен в идентификатор кластера и решить проблему сопоставления между пользовательской средой и физической средой nacos.
  • CMDB: Решите проблему хранения метаданных, подключитесь к сторонней системе cmdb и устраните взаимосвязь между приложениями, людьми и ресурсами.
  • Метрики: предоставление стандартных данных метрик для облегчения подключения к сторонним системам мониторинга.
  • Трассировка: предоставляет стандартные трассировки, которые удобны для подключения к системе SLA, отбеливанию журналов, трассировкам push-уведомлений и другим возможностям, а также могут подключаться к системе учета и выставления счетов.
  • Управление доступом: эквивалентно процессу предоставления услуг Alibaba Cloud, процессу распределения удостоверений, емкости и разрешений.
  • Управление пользователями: решите проблемы с управлением пользователями, входом в систему, sso и другими проблемами.
  • Управление правами: решение проблем распознавания личности, контроля доступа, управления ролями и т. д.
  • Система аудита: интерфейс расширения удобен для подключения к системам аудита разных компаний
  • Система уведомлений: основные изменения данных или операции, которые легко пройти через систему SMS и уведомить соответствующее лицо об изменениях данных
  • OpenAPI: предоставляет стандартный HTTP-интерфейс в стиле Rest, прост в использовании, удобен для многоязычной интеграции.
  • Консоль: простая в использовании консоль для управления службами, конфигурацией и т. д.
  • SDK: многоязычный SDK
  • Агент: аналогичный режим dns-f или интегрированный с такими схемами, как mesh
  • CLI: легкое управление продуктами из командной строки, такое же простое в использовании, как git

Модель домена:

image.png

Диаграмма классов SDK:

image.png