主頁 >  其他 > Nacos 2.0原始碼分析-Distro協議概覽

Nacos 2.0原始碼分析-Distro協議概覽

2021-08-31 20:07:52 其他

溫馨提示:
本文內容基于個人學習Nacos 2.0.1版本代碼總結而來,因個人理解差異,不保證完全正確,如有理解錯誤之處歡迎各位拍磚指正,相互學習;轉載請注明出處,

什么是Distro協議

今天來分析Nacos中使用的一種叫作Distro的協議,Distro是阿里巴巴內部使用的一種協議,用于實作分布式環境下的資料一致性,協議約定了節點之間通信的資料格式,資料決議規則,資料傳輸載體,它是一種臨時資料一致性協議,所管理的資料僅保留在記憶體中,

Distro協議用來做什么

Nacos作為一個分布式服務管理平臺(其最主要的功能之一),在分布式環境下每個節點上面的服務資訊都會有不同的狀態,當服務的可用狀態變更等一系列的問題都需要通知其他節點,每個節點上的服務串列也需要進行同步,Distro協議就是用于在不同節點之間同步節點之間的服務,除了字面上的同步之外,它還負責向其他節點報告自身的服務狀態,事實上也可以看做是一種同步,

本篇內容不設計該協議的具體操作,具體的協議實作可參考《Distro協議詳解》一文,本文僅從Nacos中所有關于Distro的類中來看看它能做什么,通過在Idea內搜索Distro開頭的類可以發現它有30個類,分別分布在nacos-corenacos-namingnacos-config模塊中,本篇只分析nacos-core模塊下的內容,因為它已經覆寫了Distro協議的完整流程,

提示:
這里可以先記住一個關鍵詞同步,所謂的同步無非就是從遠端獲取資料到本地,或者是從本地發送資料到遠端,同步的資料在這里肯定就是 服務相關的了,畢竟在官方檔案中都是這樣寫的:”服務(Service)是 Nacos 世界的一等公民”,
本篇介紹的所有內容均是為了服務于同步這個概念的,

Distro協議的核心組件

nacos-core模塊下,定義了Distro協議的所有組件,

distro
	component						Distro的一些組件,例如資料存盤物件、資料處理器、資料傳輸代理等
	entity							物體物件
	exception						例外處理
	task							任務相關
		delay						延遲任務相關組件
		execute						任務執行器相關組件
		load						加載任務相關組件
		verify						驗證任務相關組件
	DistroConfig.java				Distro配置資訊
	DistroConstants.java			Distro常量
	DistroProtocol.java 			Distro協議入口
	

com.alibaba.nacos.core.distributed.distro.component

DistroCallback

Distro回呼介面,用于異步處理之后需要回呼的場景,

public interface DistroCallback {
    
    /**
     * Callback when distro task execute successfully.
     */
    void onSuccess();
    
    /**
     * Callback when distro task execute failed.
     *
     * @param throwable throwable if execute failed caused by exception
     */
    void onFailed(Throwable throwable);
}

DistroComponentHolder

Distro組件持有者,它內部定義了一些容器(HashMap)來存盤Distro協議需要用到的資料,相當于一個大管家,

@Component
public class DistroComponentHolder {
	
    // 存盤不同型別的DistroData傳輸物件
    private final Map<String, DistroTransportAgent> transportAgentMap = new HashMap<>();
    // 存盤不同型別的DistroData裝載容器
    private final Map<String, DistroDataStorage> dataStorageMap = new HashMap<>();
    // 存盤不同型別的Distro失敗任務處理器
    private final Map<String, DistroFailedTaskHandler> failedTaskHandlerMap = new HashMap<>();
    // 存盤不同型別的DistroData資料處理器
    private final Map<String, DistroDataProcessor> dataProcessorMap = new HashMap<>();
    
    public DistroTransportAgent findTransportAgent(String type) {
        return transportAgentMap.get(type);
    }
    
    public void registerTransportAgent(String type, DistroTransportAgent transportAgent) {
        transportAgentMap.put(type, transportAgent);
    }
    
    public DistroDataStorage findDataStorage(String type) {
        return dataStorageMap.get(type);
    }
    
    public void registerDataStorage(String type, DistroDataStorage dataStorage) {
        dataStorageMap.put(type, dataStorage);
    }
    
    public Set<String> getDataStorageTypes() {
        return dataStorageMap.keySet();
    }
    
    public DistroFailedTaskHandler findFailedTaskHandler(String type) {
        return failedTaskHandlerMap.get(type);
    }
    
    public void registerFailedTaskHandler(String type, DistroFailedTaskHandler failedTaskHandler) {
        failedTaskHandlerMap.put(type, failedTaskHandler);
    }
    
    public void registerDataProcessor(DistroDataProcessor dataProcessor) {
        dataProcessorMap.putIfAbsent(dataProcessor.processType(), dataProcessor);
    }
    
    public DistroDataProcessor findDataProcessor(String processType) {
        return dataProcessorMap.get(processType);
    }
}

DistroDataProcessor

用于處理Distro協議的資料物件,

/**
 * Distro data processor.
 *
 * @author xiweng.yy
 */
public interface DistroDataProcessor {
    
    /**
     * Process type of this processor.
     * 當前處理器可處理的型別
     * @return type of this processor
     */
    String processType();
    
    /**
     * Process received data.
     * 處理接收到的資料
     * @param distroData received data	接收到的資料物件
     * @return true if process data successfully, otherwise false
     */
    boolean processData(DistroData distroData);
    
    /**
     * Process received verify data.
     * 處理接收到的驗證型別的資料
     * @param distroData    verify data	被處理的資料
     * @param sourceAddress source server address, might be get data from source server 被處理資料的來源服務器
     * @return true if the data is available, otherwise false
     */
    boolean processVerifyData(DistroData distroData, String sourceAddress);
    
    /**
     * Process snapshot data.
     * 處理快照資料
     * @param distroData snapshot data
     * @return true if process data successfully, otherwise false
     */
    boolean processSnapshot(DistroData distroData);
}

DiustroDataStorage

DistroData的存盤器

public interface DistroDataStorage {
    
    /**
     * Set this distro data storage has finished initial step.
	 * 設定當前存盤器已經初始化完畢它內部的DistroData
     */
    void finishInitial();
    
    /**
     * Whether this distro data is finished initial.
     * 當前存盤器是否已經初始化完畢內部的DistroData
     * <p>If not finished, this data storage should not send verify data to other node.
     *
     * @return {@code true} if finished, otherwise false
     */
    boolean isFinishInitial();
    
    /**
     * Get distro datum.
     * 獲取內部的DistroData
     * @param distroKey key of distro datum	資料對應的key
     * @return need to sync datum
     */
    DistroData getDistroData(DistroKey distroKey);
    
    /**
     * Get all distro datum snapshot.
     * 獲取內部存盤的所有DistroData
     * @return all datum
     */
    DistroData getDatumSnapshot();
    
    /**
     * Get verify datum.
     * 獲取所有的DistroData用于驗證
     * @return verify datum
     */
    List<DistroData> getVerifyData();
}

DistroFailedTaskHandler

用于Distro任務失敗重試

public interface DistroFailedTaskHandler {
    
    /**
     * Build retry task when distro task execute failed.
     * 當Distro任務執行失敗可以創建重試任務
     * @param distroKey distro key of failed task	失敗任務的distroKey
     * @param action action of task					任務的操作型別
     */
    void retry(DistroKey distroKey, DataOperation action);
}

DistroTransportAgent

DistroData的傳輸代理,用于發送請求,

public interface DistroTransportAgent {
    
    /**
     * Whether support transport data with callback.
     * 是否支持回呼
     * @return true if support, otherwise false
     */
    boolean supportCallbackTransport();
    
    /**
     * Sync data.
     * 同步資料
     * @param data         data			被同步的資料
     * @param targetServer target server同步的目標服務器
     * @return true is sync successfully, otherwise false
     */
    boolean syncData(DistroData data, String targetServer);
    
    /**
     * Sync data with callback.
     * 帶回呼的同步方法
     * @param data         data
     * @param targetServer target server
     * @param callback     callback
     * @throws UnsupportedOperationException if method supportCallbackTransport is false, should throw {@code
     *                                       UnsupportedOperationException}
     */
    void syncData(DistroData data, String targetServer, DistroCallback callback);
    
    /**
     * Sync verify data.
     * 同步驗證資料
     * @param verifyData   verify data
     * @param targetServer target server
     * @return true is verify successfully, otherwise false
     */
    boolean syncVerifyData(DistroData verifyData, String targetServer);
    
    /**
     * Sync verify data.
     * 帶回呼的同步驗證資料
     * @param verifyData   verify data
     * @param targetServer target server
     * @param callback     callback
     * @throws UnsupportedOperationException if method supportCallbackTransport is false, should throw {@code
     *                                       UnsupportedOperationException}
     */
    void syncVerifyData(DistroData verifyData, String targetServer, DistroCallback callback);
    
    /**
     * get Data from target server.
     * 從遠程節點獲取指定資料
     * @param key          key of data	需要獲取資料的key
     * @param targetServer target server遠端節點地址
     * @return distro data
     */
    DistroData getData(DistroKey key, String targetServer);
    
    /**
     * Get all datum snapshot from target server.
     * 從遠端節點獲取全量快照資料
     * @param targetServer target server.
     * @return distro data
     */
    DistroData getDatumSnapshot(String targetServer);
}

com.alibaba.nacos.core.distributed.distro.entity

這里存放了Distro協議的資料物件,

DistroData

Distro協議的核心物件,協議互動程序中的資料傳輸將使用此物件,它的設計也可以看做是一個容器,后期將會經常看見他,

public class DistroData {
    // 資料的key
    private DistroKey distroKey;
    // 資料的操作型別,也可以理解為是什么操作產生了此資料,或此資料用于什么操作
    private DataOperation type;
    // 資料的位元組陣列
    private byte[] content;
    
    public DistroData() {
    }
    
    public DistroData(DistroKey distroKey, byte[] content) {
        this.distroKey = distroKey;
        this.content = content;
    }
    
    public DistroKey getDistroKey() {
        return distroKey;
    }
    
    public void setDistroKey(DistroKey distroKey) {
        this.distroKey = distroKey;
    }
    
    public DataOperation getType() {
        return type;
    }
    
    public void setType(DataOperation type) {
        this.type = type;
    }
    
    public byte[] getContent() {
        return content;
    }
    
    public void setContent(byte[] content) {
        this.content = content;
    }
}

DistroKey

DistroData的key物件,可以包含較多的屬性,

public class DistroKey {
    // 資料本身的key
    private String resourceKey;
    // 資料的型別
    private String resourceType;
    // 資料傳輸的目標服務器
    private String targetServer;
    
    public DistroKey() {
    }
    
    public DistroKey(String resourceKey, String resourceType) {
        this.resourceKey = resourceKey;
        this.resourceType = resourceType;
    }
    
    public DistroKey(String resourceKey, String resourceType, String targetServer) {
        this.resourceKey = resourceKey;
        this.resourceType = resourceType;
        this.targetServer = targetServer;
    }
    
    public String getResourceKey() {
        return resourceKey;
    }
    
    public void setResourceKey(String resourceKey) {
        this.resourceKey = resourceKey;
    }
    
    public String getResourceType() {
        return resourceType;
    }
    
    public void setResourceType(String resourceType) {
        this.resourceType = resourceType;
    }
    
    public String getTargetServer() {
        return targetServer;
    }
    
    public void setTargetServer(String targetServer) {
        this.targetServer = targetServer;
    }
    
    @Override
    public boolean equals(Object o) {
        if (this == o) {
            return true;
        }
        if (o == null || getClass() != o.getClass()) {
            return false;
        }
        DistroKey distroKey = (DistroKey) o;
        return Objects.equals(resourceKey, distroKey.resourceKey) && Objects
                .equals(resourceType, distroKey.resourceType) && Objects.equals(targetServer, distroKey.targetServer);
    }
    
    @Override
    public int hashCode() {
        return Objects.hash(resourceKey, resourceType, targetServer);
    }
    
    @Override
    public String toString() {
        return "DistroKey{" + "resourceKey='" + resourceKey + '\'' + ", resourceType='" + resourceType + '\''
                + ", targetServer='" + targetServer + '\'' + '}';
    }
}

com.alibaba.nacos.core.distributed.distro.exception

com.alibaba.nacos.core.distributed.distro.task

DistroTaskEngineHolder

Distro任務引擎持有者,用于管理不同型別的任務執行引擎,

@Component
public class DistroTaskEngineHolder {
    // 延遲任務執行引擎
    private final DistroDelayTaskExecuteEngine delayTaskExecuteEngine = new DistroDelayTaskExecuteEngine();
    // 非延遲任務執行引擎
    private final DistroExecuteTaskExecuteEngine executeWorkersManager = new DistroExecuteTaskExecuteEngine();
    
    public DistroTaskEngineHolder(DistroComponentHolder distroComponentHolder) {
		// 為延遲任務執行引擎添加默認任務處理器
        DistroDelayTaskProcessor defaultDelayTaskProcessor = new DistroDelayTaskProcessor(this, distroComponentHolder);
        delayTaskExecuteEngine.setDefaultTaskProcessor(defaultDelayTaskProcessor);
    }
    
    public DistroDelayTaskExecuteEngine getDelayTaskExecuteEngine() {
        return delayTaskExecuteEngine;
    }
    
    public DistroExecuteTaskExecuteEngine getExecuteWorkersManager() {
        return executeWorkersManager;
    }
    
	/**
     * 為延遲任務添加默認任務處理器
     * @param key          處理器向容器保存時的key
     * @param nacosTaskProcessor 處理器物件
     */
    public void registerNacosTaskProcessor(Object key, NacosTaskProcessor nacosTaskProcessor) {
        this.delayTaskExecuteEngine.addProcessor(key, nacosTaskProcessor);
    }
}

com.alibaba.nacos.core.distributed.distro.task.delay

DistroDelayTask

Distro延遲任務

public class DistroDelayTask extends AbstractDelayTask {
    // 當前任務處理資料的key
    private final DistroKey distroKey;
    // 當前任務處理資料的操作型別
    private DataOperation action;
    // 當前任務創建的時間
    private long createTime;
    
    public DistroDelayTask(DistroKey distroKey, long delayTime) {
        this(distroKey, DataOperation.CHANGE, delayTime);
    }
    
	// 構造一個延遲任務
    public DistroDelayTask(DistroKey distroKey, DataOperation action, long delayTime) {
        this.distroKey = distroKey;
        this.action = action;
        this.createTime = System.currentTimeMillis();
		// 創建時設定上次處理的時間
        setLastProcessTime(createTime);
		// 設定間隔多久執行
        setTaskInterval(delayTime);
    }
    
    public DistroKey getDistroKey() {
        return distroKey;
    }
    
    public DataOperation getAction() {
        return action;
    }
    
    public long getCreateTime() {
        return createTime;
    }
    
    /**
     * 從字面意思是合并任務,實際的操作證明它是用于更新過時的任務
     * 在向任務串列添加新的任務時,使用新任務的key來從任務串列獲取,若結果不為空,表明此任務已經存在
     * 相同的任務再次添加的話,就重復了,因此再此合并
     * 為什么新的任務會過時?(新任務指的是當前類)
     * 想要理解此處邏輯,請參考{@link com.alibaba.nacos.common.task.engine.NacosDelayTaskExecuteEngine#addTask(Object,
     *  AbstractDelayTask)}.添加任務時是帶鎖操作的,因此添加的先后順序不能保證
     * @param task task 已存在的任務
     */
    @Override
    public void merge(AbstractDelayTask task) {
        if (!(task instanceof DistroDelayTask)) {
            return;
        }
        DistroDelayTask oldTask = (DistroDelayTask) task;
        // 若舊的任務和新的任務的操作型別不同,并且新任務的創建時間小于舊任務的創建時間,說明當前這個新任務還未被添加成功
        // 這個新的任務已經過時了,不需要再執行這個任務的操作,因此將舊的任務的操作型別和創建時間設定給新任務
        if (!action.equals(oldTask.getAction()) && createTime < oldTask.getCreateTime()) {
            action = oldTask.getAction();
            createTime = oldTask.getCreateTime();
        }
        setLastProcessTime(oldTask.getLastProcessTime());
    }
}

DistroDelayTaskExecuteEngine

延遲任務執行引擎

public class DistroDelayTaskExecuteEngine extends NacosDelayTaskExecuteEngine {
    
    public DistroDelayTaskExecuteEngine() {
        super(DistroDelayTaskExecuteEngine.class.getName(), Loggers.DISTRO);
    }
    
    @Override
    public void addProcessor(Object key, NacosTaskProcessor taskProcessor) {
		// 構建當前任務的key
        Object actualKey = getActualKey(key);
        super.addProcessor(actualKey, taskProcessor);
    }
    
    @Override
    public NacosTaskProcessor getProcessor(Object key) {
        Object actualKey = getActualKey(key);
        return super.getProcessor(actualKey);
    }
    
    private Object getActualKey(Object key) {
        return key instanceof DistroKey ? ((DistroKey) key).getResourceType() : key;
    }
}

DistroDelayTaskProcessor

延遲任務處理器

/**
 * Distro delay task processor.
 *
 * @author xiweng.yy
 */
public class DistroDelayTaskProcessor implements NacosTaskProcessor {
    // Distro任務引擎持有者
    private final DistroTaskEngineHolder distroTaskEngineHolder;
    // Distro組件持有者
    private final DistroComponentHolder distroComponentHolder;
    
    public DistroDelayTaskProcessor(DistroTaskEngineHolder distroTaskEngineHolder,
            DistroComponentHolder distroComponentHolder) {
        this.distroTaskEngineHolder = distroTaskEngineHolder;
        this.distroComponentHolder = distroComponentHolder;
    }
    
    @Override
    public boolean process(NacosTask task) {
		// 不處理非延遲任務
        if (!(task instanceof DistroDelayTask)) {
            return true;
        }
        DistroDelayTask distroDelayTask = (DistroDelayTask) task;
        DistroKey distroKey = distroDelayTask.getDistroKey();
		// 根據不同的操作型別創建具體的任務
        switch (distroDelayTask.getAction()) {
            case DELETE:
                DistroSyncDeleteTask syncDeleteTask = new DistroSyncDeleteTask(distroKey, distroComponentHolder);
                distroTaskEngineHolder.getExecuteWorkersManager().addTask(distroKey, syncDeleteTask);
                return true;
            case CHANGE:
            case ADD:
                DistroSyncChangeTask syncChangeTask = new DistroSyncChangeTask(distroKey, distroComponentHolder);
                distroTaskEngineHolder.getExecuteWorkersManager().addTask(distroKey, syncChangeTask);
                return true;
            default:
                return false;
        }
    }
}

com.alibaba.nacos.core.distributed.distro.task.execute

AbstractDistroExecuteTask

抽象的執行任務,定義了任務處理流程,

public abstract class AbstractDistroExecuteTask extends AbstractExecuteTask {
    
    private final DistroKey distroKey;
    
    private final DistroComponentHolder distroComponentHolder;
    
    protected AbstractDistroExecuteTask(DistroKey distroKey, DistroComponentHolder distroComponentHolder) {
        this.distroKey = distroKey;
        this.distroComponentHolder = distroComponentHolder;
    }
    
    protected DistroKey getDistroKey() {
        return distroKey;
    }
    
    protected DistroComponentHolder getDistroComponentHolder() {
        return distroComponentHolder;
    }
    
    @Override
    public void run() {
		// 獲取被處理的資料資源型別
        String type = getDistroKey().getResourceType();
		// 根據型別獲取資料傳輸代理
        DistroTransportAgent transportAgent = distroComponentHolder.findTransportAgent(type);
        if (null == transportAgent) {
            Loggers.DISTRO.warn("No found transport agent for type [{}]", type);
            return;
        }
        Loggers.DISTRO.info("[DISTRO-START] {}", toString());
		// 判斷代理物件是否支持回呼
        if (transportAgent.supportCallbackTransport()) {
            doExecuteWithCallback(new DistroExecuteCallback());
        } else {
            executeDistroTask();
        }
    }
    
	// 執行任務
    private void executeDistroTask() {
        try {
            boolean result = doExecute();
            if (!result) {
				// 執行失敗之后,進行失敗處理
                handleFailedTask();
            }
            Loggers.DISTRO.info("[DISTRO-END] {} result: {}", toString(), result);
        } catch (Exception e) {
            Loggers.DISTRO.warn("[DISTRO] Sync data change failed.", e);
			// 執行失敗任務,進行失敗處理
            handleFailedTask();
        }
    }
    
    /**
     * Get {@link DataOperation} for current task.
     *
     * @return data operation
     */
    protected abstract DataOperation getDataOperation();
    
    /**
     * Do execute for different sub class.
     *
     * @return result of execute
     */
    protected abstract boolean doExecute();
    
    /**
     * Do execute with callback for different sub class.
     *
     * @param callback callback
     */
    protected abstract void doExecuteWithCallback(DistroCallback callback);
    
    /**
     * Handle failed task.
	 * 處理失敗的任務
     */
    protected void handleFailedTask() {
        String type = getDistroKey().getResourceType();
		// 使用失敗任務處理器進行重試
        DistroFailedTaskHandler failedTaskHandler = distroComponentHolder.findFailedTaskHandler(type);
        if (null == failedTaskHandler) {
            Loggers.DISTRO.warn("[DISTRO] Can't find failed task for type {}, so discarded", type);
            return;
        }
        failedTaskHandler.retry(getDistroKey(), getDataOperation());
    }
    
    private class DistroExecuteCallback implements DistroCallback {
        
        @Override
        public void onSuccess() {
            Loggers.DISTRO.info("[DISTRO-END] {} result: true", getDistroKey().toString());
        }
        
        @Override
        public void onFailed(Throwable throwable) {
            if (null == throwable) {
                Loggers.DISTRO.info("[DISTRO-END] {} result: false", getDistroKey().toString());
            } else {
                Loggers.DISTRO.warn("[DISTRO] Sync data change failed.", throwable);
            }
            handleFailedTask();
        }
    }
}

DistroExecuteTaskExecuteEngine

Distro協議負責執行任務的執行引擎

package com.alibaba.nacos.common.task.engine;


public class DistroExecuteTaskExecuteEngine extends NacosExecuteTaskExecuteEngine {
    
	// 直接創建了一個新的NacosExecuteTaskExecuteEngine執行引擎
    public DistroExecuteTaskExecuteEngine() {
        super(DistroExecuteTaskExecuteEngine.class.getSimpleName(), Loggers.DISTRO);
    }
}

NacosExecuteTaskExecuteEngine

package com.alibaba.nacos.common.task.engine;

/**
 * Nacos execute task execute engine.
 * Nacos負責執行任務的執行引擎
 * @author xiweng.yy
 */
public class NacosExecuteTaskExecuteEngine extends AbstractNacosTaskExecuteEngine<AbstractExecuteTask> {
    
	// 任務執行者
    private final TaskExecuteWorker[] executeWorkers;
    
    public NacosExecuteTaskExecuteEngine(String name, Logger logger) {
		// 任務執行者的數量,取決于CPU的核數,默認為CPU核數的1.5~2倍,傳遞的引數是表示需要產生的執行緒數量是CPU核數的多少倍
        this(name, logger, ThreadUtils.getSuitableThreadCount(1));
    }
    
    public NacosExecuteTaskExecuteEngine(String name, Logger logger, int dispatchWorkerCount) {
        super(logger);
		// 創建一組任務執行者
        executeWorkers = new TaskExecuteWorker[dispatchWorkerCount];
        for (int mod = 0; mod < dispatchWorkerCount; ++mod) {
            executeWorkers[mod] = new TaskExecuteWorker(name, mod, dispatchWorkerCount, getEngineLog());
        }
    }
    
    @Override
    public int size() {
        int result = 0;
        for (TaskExecuteWorker each : executeWorkers) {
            result += each.pendingTaskCount();
        }
        return result;
    }
    
    @Override
    public boolean isEmpty() {
        return 0 == size();
    }
    
    @Override
    public void addTask(Object tag, AbstractExecuteTask task) {
		// 從父類獲取任務處理器
        NacosTaskProcessor processor = getProcessor(tag);
		// 若存在處理器,則用處理器來處理
        if (null != processor) {
            processor.process(task);
            return;
        }
		// 不存在處理器則使用worker處理
        TaskExecuteWorker worker = getWorker(tag);
        worker.process(task);
    }
    
    private TaskExecuteWorker getWorker(Object tag) {
		// 計算當前任務應該由哪個worker處理
        int idx = (tag.hashCode() & Integer.MAX_VALUE) % workersCount();
        return executeWorkers[idx];
    }
    
    private int workersCount() {
        return executeWorkers.length;
    }
    
    @Override
    public AbstractExecuteTask removeTask(Object key) {
        throw new UnsupportedOperationException("ExecuteTaskEngine do not support remove task");
    }
    
    @Override
    public Collection<Object> getAllTaskKeys() {
        throw new UnsupportedOperationException("ExecuteTaskEngine do not support get all task keys");
    }
    
    @Override
    public void shutdown() throws NacosException {
        for (TaskExecuteWorker each : executeWorkers) {
            each.shutdown();
        }
    }
    
    /**
     * Get workers status.
     *
     * @return workers status string
     */
    public String workersStatus() {
        StringBuilder sb = new StringBuilder();
        for (TaskExecuteWorker worker : executeWorkers) {
            sb.append(worker.status()).append("\n");
        }
        return sb.toString();
    }
}

TaskExecuteWorker

package com.alibaba.nacos.common.task.engine;

/**
 * Nacos execute task execute worker.
 * Nacos任務執行者,每個執行者在創建的時候會同時啟動一個執行緒InnerWorker,持續從內部佇列中獲取需要處理的任務
 * @author xiweng.yy
 */
public final class TaskExecuteWorker implements NacosTaskProcessor, Closeable {

    /**
     * Max task queue size 32768.
     * 佇列最大數量為32768
     */
    private static final int QUEUE_CAPACITY = 1 << 15;

    private final Logger log;

    /**
     * 當前執行者執行緒的名稱
     */
    private final String name;

    /**
     * 負責處理的執行緒佇列
     */
    private final BlockingQueue<Runnable> queue;

    /**
     * 作業狀態
     */
    private final AtomicBoolean closed;

    public TaskExecuteWorker(final String name, final int mod, final int total) {
        this(name, mod, total, null);
    }

    public TaskExecuteWorker(final String name, final int mod, final int total, final Logger logger) {
        /**
         * 執行執行緒的名稱,以DistroExecuteTaskExecuteEngine舉例:
         * DistroExecuteTaskExecuteEngine_0%8
         * DistroExecuteTaskExecuteEngine_1%8
         * DistroExecuteTaskExecuteEngine_2%8
         * DistroExecuteTaskExecuteEngine_3%8
         * DistroExecuteTaskExecuteEngine_4%8
         * DistroExecuteTaskExecuteEngine_5%8
         * DistroExecuteTaskExecuteEngine_6%8
         * DistroExecuteTaskExecuteEngine_7%8
         */
        this.name = name + "_" + mod + "%" + total;
        this.queue = new ArrayBlockingQueue<Runnable>(QUEUE_CAPACITY);
        this.closed = new AtomicBoolean(false);
        this.log = null == logger ? LoggerFactory.getLogger(TaskExecuteWorker.class) : logger;
        // 啟動一個新執行緒來消費佇列
        new InnerWorker(name).start();
    }

    public String getName() {
        return name;
    }

    @Override
    public boolean process(NacosTask task) {
        if (task instanceof AbstractExecuteTask) {
            putTask((Runnable) task);
        }
        return true;
    }

    private void putTask(Runnable task) {
        try {
            queue.put(task);
        } catch (InterruptedException ire) {
            log.error(ire.toString(), ire);
        }
    }

    public int pendingTaskCount() {
        return queue.size();
    }

    /**
     * Worker status.
     */
    public String status() {
        return name + ", pending tasks: " + pendingTaskCount();
    }

    @Override
    public void shutdown() throws NacosException {
        queue.clear();
        closed.compareAndSet(false, true);
    }

    /**
     * Inner execute worker.
     */
    private class InnerWorker extends Thread {

        InnerWorker(String name) {
            setDaemon(false);
            setName(name);
        }

        @Override
        public void run() {
            // 若執行緒還未中斷,則持續執行
            while (!closed.get()) {
                try {
                    // 從佇列獲取任務
                    Runnable task = queue.take();
                    long begin = System.currentTimeMillis();
                    // 在當前InnerWorker執行緒內執行任務
                    task.run();
                    long duration = System.currentTimeMillis() - begin;
                    // 若任務執行時間超過1秒,則警告
                    if (duration > 1000L) {
                        log.warn("task {} takes {}ms", task, duration);
                    }
                } catch (Throwable e) {
                    log.error("[TASK-FAILED] " + e.toString(), e);
                }
            }
        }
    }
}

DistroSyncChangeTask

Distro同步變更任務,此任務用于向其他節點發送本機資料

public class DistroSyncChangeTask extends AbstractDistroExecuteTask {
    
	// 此任務操作型別為變更
    private static final DataOperation OPERATION = DataOperation.CHANGE;
    
    public DistroSyncChangeTask(DistroKey distroKey, DistroComponentHolder distroComponentHolder) {
        super(distroKey, distroComponentHolder);
    }
    
    @Override
    protected DataOperation getDataOperation() {
        return OPERATION;
    }
    
	/**
     * 執行不帶回呼的任務
     * @return
     */
    @Override
    protected boolean doExecute() {
		// 獲取同步的資料型別
        String type = getDistroKey().getResourceType();
		// 獲取同步資料
        DistroData distroData = https://www.cnblogs.com/lukama/p/getDistroData(type);
        if (null == distroData) {
            Loggers.DISTRO.warn("[DISTRO] {} with null data to sync, skip", toString());
            return true;
        }
		// 使用DistroTransportAgent同步資料
        return getDistroComponentHolder().findTransportAgent(type).syncData(distroData, getDistroKey().getTargetServer());
    }
    
	/**
     * 執行帶回呼的任務
     * @param callback callback
     */
    @Override
    protected void doExecuteWithCallback(DistroCallback callback) {
        String type = getDistroKey().getResourceType();
        DistroData distroData = https://www.cnblogs.com/lukama/p/getDistroData(type);
        if (null == distroData) {
            Loggers.DISTRO.warn("[DISTRO] {} with null data to sync, skip", toString());
            return;
        }
        getDistroComponentHolder().findTransportAgent(type).syncData(distroData, getDistroKey().getTargetServer(), callback);
    }
    
    @Override
    public String toString() {
        return "DistroSyncChangeTask for " + getDistroKey().toString();
    }
    
    private DistroData getDistroData(String type) {
        DistroData result = getDistroComponentHolder().findDataStorage(type).getDistroData(getDistroKey());
        if (null != result) {
            result.setType(OPERATION);
        }
        return result;
    }
}

DistroSyncDeleteTask

Distro同步洗掉任務,用于向其他節點發送本機洗掉的資料

public class DistroSyncDeleteTask extends AbstractDistroExecuteTask {
    
	// 此任務操作型別為洗掉
    private static final DataOperation OPERATION = DataOperation.DELETE;
    
    public DistroSyncDeleteTask(DistroKey distroKey, DistroComponentHolder distroComponentHolder) {
        super(distroKey, distroComponentHolder);
    }
    
    @Override
    protected DataOperation getDataOperation() {
        return OPERATION;
    }
    
	/**
     * 執行不帶回呼的任務
     * @return
     */
    @Override
    protected boolean doExecute() {
		// 構建請求引數
        String type = getDistroKey().getResourceType();
        DistroData distroData = https://www.cnblogs.com/lukama/p/new DistroData();
        distroData.setDistroKey(getDistroKey());
        distroData.setType(OPERATION);
		// 使用DistroTransportAgent同步資料
        return getDistroComponentHolder().findTransportAgent(type).syncData(distroData, getDistroKey().getTargetServer());
    }
    
	/**
     * 執行帶回呼的任務
     * @param callback callback
     */
    @Override
    protected void doExecuteWithCallback(DistroCallback callback) {
        String type = getDistroKey().getResourceType();
        DistroData distroData = new DistroData();
        distroData.setDistroKey(getDistroKey());
        distroData.setType(OPERATION);
        getDistroComponentHolder().findTransportAgent(type).syncData(distroData, getDistroKey().getTargetServer(), callback);
    }
    
    @Override
    public String toString() {
        return"DistroSyncDeleteTask for " + getDistroKey().toString();
    }
}

提示:

  • DistroSyncChangeTask是將本機所有的服務發送到其他節點
  • DistroSyncDeleteTask是將本機洗掉的服務發送到其他節點

com.alibaba.nacos.core.distributed.distro.task.load

DistroLoadDataTask

Distro全量資料同步任務,用于在節點啟動后首次從其他節點同步服務資料到當前節點,

public class DistroLoadDataTask implements Runnable {

	// 節點管理器
    private final ServerMemberManager memberManager;
	// Distro協議組件持有者
    private final DistroComponentHolder distroComponentHolder;
	// Distro協議配置
    private final DistroConfig distroConfig;
	// 回呼函式
    private final DistroCallback loadCallback;
	// 已加載資料集合
    private final Map<String, Boolean> loadCompletedMap;

    public DistroLoadDataTask(ServerMemberManager memberManager, DistroComponentHolder distroComponentHolder, DistroConfig distroConfig, DistroCallback loadCallback) {
        this.memberManager = memberManager;
        this.distroComponentHolder = distroComponentHolder;
        this.distroConfig = distroConfig;
        this.loadCallback = loadCallback;
        loadCompletedMap = new HashMap<>(1);
    }

    @Override
    public void run() {
        try {
			// 首次加載
            load();
			// 若首次加載沒有完成,繼續加載
            if (!checkCompleted()) {
				// 繼續創建一個新的加載任務進行加載
                GlobalExecutor.submitLoadDataTask(this, distroConfig.getLoadDataRetryDelayMillis());
            } else {
				// 觸發回呼函式
                loadCallback.onSuccess();
                Loggers.DISTRO.info("[DISTRO-INIT] load snapshot data success");
            }
        } catch (Exception e) {
            loadCallback.onFailed(e);
            Loggers.DISTRO.error("[DISTRO-INIT] load snapshot data failed. ", e);
        }
    }

    private void load() throws Exception {
		// 若出自身之外沒有其他節點,則休眠1秒,可能其他節點還未啟動完畢
        while (memberManager.allMembersWithoutSelf().isEmpty()) {
            Loggers.DISTRO.info("[DISTRO-INIT] waiting server list init...");
            TimeUnit.SECONDS.sleep(1);
        }
		// 若資料型別為空,說明distroComponentHolder的組件注冊器還未初始化完畢(v1版本為DistroHttpRegistry, v2版本為DistroClientComponentRegistry)
        while (distroComponentHolder.getDataStorageTypes().isEmpty()) {
            Loggers.DISTRO.info("[DISTRO-INIT] waiting distro data storage register...");
            TimeUnit.SECONDS.sleep(1);
        }
		// 加載每個型別的資料
        for (String each : distroComponentHolder.getDataStorageTypes()) {
            if (!loadCompletedMap.containsKey(each) || !loadCompletedMap.get(each)) {
				// 呼叫加載方法,并標記已處理
                loadCompletedMap.put(each, loadAllDataSnapshotFromRemote(each));
            }
        }
    }

	/**
     * 從其他節點獲取同步資料
     * @param resourceType
     * @return
     */
    private boolean loadAllDataSnapshotFromRemote(String resourceType) {
		// 獲取資料傳輸物件
        DistroTransportAgent transportAgent = distroComponentHolder.findTransportAgent(resourceType);
		// 獲取資料處理器
        DistroDataProcessor dataProcessor = distroComponentHolder.findDataProcessor(resourceType);
        if (null == transportAgent || null == dataProcessor) {
            Loggers.DISTRO.warn("[DISTRO-INIT] Can't find component for type {}, transportAgent: {}, dataProcessor: {}", resourceType, transportAgent, dataProcessor);
            return false;
        }
		// 向每個節點請求資料
        for (Member each : memberManager.allMembersWithoutSelf()) {
            try {
                Loggers.DISTRO.info("[DISTRO-INIT] load snapshot {} from {}", resourceType, each.getAddress());
				// 獲取到資料
                DistroData distroData = https://www.cnblogs.com/lukama/p/transportAgent.getDatumSnapshot(each.getAddress());
				// 決議資料
                boolean result = dataProcessor.processSnapshot(distroData);
                Loggers.DISTRO.info("[DISTRO-INIT] load snapshot {} from {} result: {}", resourceType, each.getAddress(), result);
				// 若決議成功,標記此型別資料已加載完畢
                if (result) {
                    distroComponentHolder.findDataStorage(resourceType).finishInitial();
                    return true;
                }
            } catch (Exception e) {
                Loggers.DISTRO.error("[DISTRO-INIT] load snapshot {} from {} failed.", resourceType, each.getAddress(), e);
            }
        }
        return false;
    }

	// 判斷是否完成加載
    private boolean checkCompleted() {
		// 若待加載的資料型別數量和已經加載完畢的資料型別數量不一致,鐵定是未加載完成
        if (distroComponentHolder.getDataStorageTypes().size() != loadCompletedMap.size()) {
            return false;
        }
		// 若加載完畢串列內的狀態有false的,說明可能是決議失敗,還需要重新加載
        for (Boolean each : loadCompletedMap.values()) {
            if (!each) {
                return false;
            }
        }
        return true;
    }
}

com.alibaba.nacos.core.distributed.distro.task.verify

DistroVerifyExecuteTask

Distro資料驗證任務執行器,用于向其他節點發送當前節點負責的Client狀態報告,通知對方此Client正常服務,它的資料處理維度是DistroData,

/**
 * Execute distro verify task.
 * 執行Distro協議資料驗證的任務,為每個DistroData發送一個異步的rpc請求
 * @author xiweng.yy
 */
public class DistroVerifyExecuteTask extends AbstractExecuteTask {

    /**
     * 被驗證資料的傳輸物件
     */
    private final DistroTransportAgent transportAgent;

    /**
     * 被驗證資料
     */
    private final List<DistroData> verifyData;

    /**
     * 目標節點
     */
    private final String targetServer;

    /**
     * 被驗證資料的型別
     */
    private final String resourceType;

    public DistroVerifyExecuteTask(DistroTransportAgent transportAgent, List<DistroData> verifyData,
            String targetServer, String resourceType) {
        this.transportAgent = transportAgent;
        this.verifyData = https://www.cnblogs.com/lukama/p/verifyData;
        this.targetServer = targetServer;
        this.resourceType = resourceType;
    }

    @Override
    public void run() {
        for (DistroData each : verifyData) {
            try {
                // 判斷傳輸物件是否支持回呼(若是http的則不支持,實際上沒區別,當前2.0.1版本沒有實作回呼的實質內容)
                if (transportAgent.supportCallbackTransport()) {
                    doSyncVerifyDataWithCallback(each);
                } else {
                    doSyncVerifyData(each);
                }
            } catch (Exception e) {
                Loggers.DISTRO
                        .error("[DISTRO-FAILED] verify data for type {} to {} failed.", resourceType, targetServer, e);
            }
        }
    }

    /**
     * 支持回呼的同步資料驗證
     * @param data
     */
    private void doSyncVerifyDataWithCallback(DistroData data) {
        // 回呼實際上,也沒啥,,,基本算是空物件
        transportAgent.syncVerifyData(data, targetServer, new DistroVerifyCallback());
    }

    /**
     * 不支持回呼的同步資料驗證
     * @param data
     */
    private void doSyncVerifyData(DistroData data) {
        transportAgent.syncVerifyData(data, targetServer);
    }

    /**
     * TODO add verify monitor.
     */
    private class DistroVerifyCallback implements DistroCallback {

        @Override
        public void onSuccess() {
            if (Loggers.DISTRO.isDebugEnabled()) {
                Loggers.DISTRO.debug("[DISTRO] verify data for type {} to {} success", resourceType, targetServer);
            }
        }

        @Override
        public void onFailed(Throwable throwable) {
            if (Loggers.DISTRO.isDebugEnabled()) {
                Loggers.DISTRO
                        .debug("[DISTRO-FAILED] verify data for type {} to {} failed.", resourceType, targetServer,
                                throwable);
            }
        }
    }
}

DistroVerifyTimedTask

定時驗證任務,此任務在啟動時延遲5秒,間隔5秒執行,主要用于為每個節點創建一個資料驗證的執行任務DistroVerifyExecuteTask,它的資料處理維度是Member,

/**
 * Timed to start distro verify task.
 * 啟動Distro協議的資料驗證流程
 * @author xiweng.yy
 */
public class DistroVerifyTimedTask implements Runnable {

    private final ServerMemberManager serverMemberManager;

    private final DistroComponentHolder distroComponentHolder;

    private final DistroExecuteTaskExecuteEngine executeTaskExecuteEngine;

    public DistroVerifyTimedTask(ServerMemberManager serverMemberManager, DistroComponentHolder distroComponentHolder,
            DistroExecuteTaskExecuteEngine executeTaskExecuteEngine) {
        this.serverMemberManager = serverMemberManager;
        this.distroComponentHolder = distroComponentHolder;
        this.executeTaskExecuteEngine = executeTaskExecuteEngine;
    }

    @Override
    public void run() {
        try {
            List<Member> targetServer = serverMemberManager.allMembersWithoutSelf();
            if (Loggers.DISTRO.isDebugEnabled()) {
                Loggers.DISTRO.debug("server list is: {}", targetServer);
            }
            for (String each : distroComponentHolder.getDataStorageTypes()) {
                verifyForDataStorage(each, targetServer);
            }
        } catch (Exception e) {
            Loggers.DISTRO.error("[DISTRO-FAILED] verify task failed.", e);
        }
    }

    private void verifyForDataStorage(String type, List<Member> targetServer) {
        DistroDataStorage dataStorage = distroComponentHolder.findDataStorage(type);
        if (!dataStorage.isFinishInitial()) {
            Loggers.DISTRO.warn("data storage {} has not finished initial step, do not send verify data",
                    dataStorage.getClass().getSimpleName());
            return;
        }
        List<DistroData> verifyData = https://www.cnblogs.com/lukama/p/dataStorage.getVerifyData();
        if (null == verifyData || verifyData.isEmpty()) {
            return;
        }
        for (Member member : targetServer) {
            DistroTransportAgent agent = distroComponentHolder.findTransportAgent(type);
            if (null == agent) {
                continue;
            }
            executeTaskExecuteEngine.addTask(member.getAddress() + type,
                    new DistroVerifyExecuteTask(agent, verifyData, member.getAddress(), type));
        }
    }
}

DistroConfig

Distro協議的配置資訊,

public class DistroConfig {
    
    private static final DistroConfig INSTANCE = new DistroConfig();
    // 同步任務延遲時長(單位:毫秒)
    private long syncDelayMillis = DistroConstants.DEFAULT_DATA_SYNC_DELAY_MILLISECONDS;
    // 同步任務超時時長(單位:毫秒)
    private long syncTimeoutMillis = DistroConstants.DEFAULT_DATA_SYNC_TIMEOUT_MILLISECONDS;
    // 同步任務重試延遲時長(單位:毫秒)
    private long syncRetryDelayMillis = DistroConstants.DEFAULT_DATA_SYNC_RETRY_DELAY_MILLISECONDS;
    // 驗證任務執行間隔時長(單位:毫秒)
    private long verifyIntervalMillis = DistroConstants.DEFAULT_DATA_VERIFY_INTERVAL_MILLISECONDS;
    // 驗證任務超時時長(單位:毫秒)
    private long verifyTimeoutMillis = DistroConstants.DEFAULT_DATA_VERIFY_TIMEOUT_MILLISECONDS;
    // 首次同步資料重試延遲時長(單位:毫秒)
    private long loadDataRetryDelayMillis = DistroConstants.DEFAULT_DATA_LOAD_RETRY_DELAY_MILLISECONDS;
    
    private DistroConfig() {
        try {
			// 嘗試從環境資訊中獲取配置
            getDistroConfigFromEnv();
        } catch (Exception e) {
            Loggers.CORE.warn("Get Distro config from env failed, will use default value", e);
        }
    }
    
	/**
     * 從環境資訊中獲取配置,若沒有,則使用默認值
     */
    private void getDistroConfigFromEnv() {
		
		// 從常量物件中獲取key和default value 
		
        syncDelayMillis = EnvUtil.getProperty(DistroConstants.DATA_SYNC_DELAY_MILLISECONDS, Long.class,
                DistroConstants.DEFAULT_DATA_SYNC_DELAY_MILLISECONDS);
        syncTimeoutMillis = EnvUtil.getProperty(DistroConstants.DATA_SYNC_TIMEOUT_MILLISECONDS, Long.class,
                DistroConstants.DEFAULT_DATA_SYNC_TIMEOUT_MILLISECONDS);
        syncRetryDelayMillis = EnvUtil.getProperty(DistroConstants.DATA_SYNC_RETRY_DELAY_MILLISECONDS, Long.class,
                DistroConstants.DEFAULT_DATA_SYNC_RETRY_DELAY_MILLISECONDS);
        verifyIntervalMillis = EnvUtil.getProperty(DistroConstants.DATA_VERIFY_INTERVAL_MILLISECONDS, Long.class,
                DistroConstants.DEFAULT_DATA_VERIFY_INTERVAL_MILLISECONDS);
        verifyTimeoutMillis = EnvUtil.getProperty(DistroConstants.DATA_VERIFY_TIMEOUT_MILLISECONDS, Long.class,
                DistroConstants.DEFAULT_DATA_VERIFY_TIMEOUT_MILLISECONDS);
        loadDataRetryDelayMillis = EnvUtil.getProperty(DistroConstants.DATA_LOAD_RETRY_DELAY_MILLISECONDS, Long.class,
                DistroConstants.DEFAULT_DATA_LOAD_RETRY_DELAY_MILLISECONDS);
    }
    
    public static DistroConfig getInstance() {
        return INSTANCE;
    }
    
    public long getSyncDelayMillis() {
        return syncDelayMillis;
    }
    
    public void setSyncDelayMillis(long syncDelayMillis) {
        this.syncDelayMillis = syncDelayMillis;
    }
    
    public long getSyncTimeoutMillis() {
        return syncTimeoutMillis;
    }
    
    public void setSyncTimeoutMillis(long syncTimeoutMillis) {
        this.syncTimeoutMillis = syncTimeoutMillis;
    }
    
    public long getSyncRetryDelayMillis() {
        return syncRetryDelayMillis;
    }
    
    public void setSyncRetryDelayMillis(long syncRetryDelayMillis) {
        this.syncRetryDelayMillis = syncRetryDelayMillis;
    }
    
    public long getVerifyIntervalMillis() {
        return verifyIntervalMillis;
    }
    
    public void setVerifyIntervalMillis(long verifyIntervalMillis) {
        this.verifyIntervalMillis = verifyIntervalMillis;
    }
    
    public long getVerifyTimeoutMillis() {
        return verifyTimeoutMillis;
    }
    
    public void setVerifyTimeoutMillis(long verifyTimeoutMillis) {
        this.verifyTimeoutMillis = verifyTimeoutMillis;
    }
    
    public long getLoadDataRetryDelayMillis() {
        return loadDataRetryDelayMillis;
    }
    
    public void setLoadDataRetryDelayMillis(long loadDataRetryDelayMillis) {
        this.loadDataRetryDelayMillis = loadDataRetryDelayMillis;
    }
}

DistroConstants

Distro常量配置,主要定義了一些關于任務執行時長的可配置的配置名稱和對應的默認值,具體的使用,可以參考DistroConfig

public class DistroConstants {
    
    public static final String DATA_SYNC_DELAY_MILLISECONDS = "nacos.core.protocol.distro.data.sync.delayMs";
    
    public static final long DEFAULT_DATA_SYNC_DELAY_MILLISECONDS = 1000L;
    
    public static final String DATA_SYNC_TIMEOUT_MILLISECONDS = "nacos.core.protocol.distro.data.sync.timeoutMs";
    
    public static final long DEFAULT_DATA_SYNC_TIMEOUT_MILLISECONDS = 3000L;
    
    public static final String DATA_SYNC_RETRY_DELAY_MILLISECONDS = "nacos.core.protocol.distro.data.sync.retryDelayMs";
    
    public static final long DEFAULT_DATA_SYNC_RETRY_DELAY_MILLISECONDS = 3000L;
    
    public static final String DATA_VERIFY_INTERVAL_MILLISECONDS = "nacos.core.protocol.distro.data.verify.intervalMs";
    
    public static final long DEFAULT_DATA_VERIFY_INTERVAL_MILLISECONDS = 5000L;
    
    public static final String DATA_VERIFY_TIMEOUT_MILLISECONDS = "nacos.core.protocol.distro.data.verify.timeoutMs";
    
    public static final long DEFAULT_DATA_VERIFY_TIMEOUT_MILLISECONDS = 3000L;
    
    public static final String DATA_LOAD_RETRY_DELAY_MILLISECONDS = "nacos.core.protocol.distro.data.load.retryDelayMs";
    
    public static final long DEFAULT_DATA_LOAD_RETRY_DELAY_MILLISECONDS = 30000L;
    
}

DistroProtocol

Distro協議的真正入口,這里將使用上面定義的所有組件來共同完實作Distro協議,可以看到它使用了Spring的@Componet注解,意味著它將被Spring容器管理,執行到構造方法的時候將會啟動Distro協議的作業,

@Component
public class DistroProtocol {

    private Logger logger = LoggerFactory.getLogger(DistroProtocol.class);

    /**
     * 節點管理器
     */
    private final ServerMemberManager memberManager;

    /**
     * Distro組件持有者
     */
    private final DistroComponentHolder distroComponentHolder;

    /**
     * Distro任務引擎持有者
     */
    private final DistroTaskEngineHolder distroTaskEngineHolder;

    private volatile boolean isInitialized = false;

    public DistroProtocol(ServerMemberManager memberManager, DistroComponentHolder distroComponentHolder,
            DistroTaskEngineHolder distroTaskEngineHolder) {
        this.memberManager = memberManager;
        this.distroComponentHolder = distroComponentHolder;
        this.distroTaskEngineHolder = distroTaskEngineHolder;
        // 啟動Distro協議
        startDistroTask();
    }

    private void startDistroTask() {
        // 單機模式不進行資料同步操作
        if (EnvUtil.getStandaloneMode()) {
            isInitialized = true;
            return;
        }
        // 開啟節點Client狀態報告任務
        startVerifyTask();
        // 啟動資料同步任務
        startLoadTask();
    }

    /**
     * 從其他節點獲取資料到當前節點
     */
    private void startLoadTask() {
        DistroCallback loadCallback = new DistroCallback() {
            @Override
            public void onSuccess() {
                isInitialized = true;
            }

            @Override
            public void onFailed(Throwable throwable) {
                isInitialized = false;
            }
        };
        // 提交資料加載任務
        GlobalExecutor.submitLoadDataTask(new DistroLoadDataTask(memberManager, distroComponentHolder, DistroConfig.getInstance(), loadCallback));
    }

    private void startVerifyTask() {
        // 啟動資料報告的定時任務
        GlobalExecutor.schedulePartitionDataTimedSync(
            new DistroVerifyTimedTask(
                memberManager,
                distroComponentHolder,
                distroTaskEngineHolder.getExecuteWorkersManager()
            ),
        DistroConfig.getInstance().getVerifyIntervalMillis());
    }

    public boolean isInitialized() {
        return isInitialized;
    }

    /**
     * Start to sync by configured delay.
     * 按配置的延遲開始同步
     * @param distroKey distro key of sync data
     * @param action    the action of data operation
     */
    public void sync(DistroKey distroKey, DataOperation action) {
        sync(distroKey, action, DistroConfig.getInstance().getSyncDelayMillis());
    }

    /**
     * Start to sync data to all remote server.
     * 開始將資料同步到其他節點
     * @param distroKey distro key of sync data
     * @param action    the action of data operation
     * @param delay     delay time for sync
     */
    public void sync(DistroKey distroKey, DataOperation action, long delay) {
        for (Member each : memberManager.allMembersWithoutSelf()) {
            syncToTarget(distroKey, action, each.getAddress(), delay);
        }
    }

    /**
     * Start to sync to target server.
     *
     * @param distroKey    distro key of sync data
     * @param action       the action of data operation
     * @param targetServer target server
     * @param delay        delay time for sync
     */
    public void syncToTarget(DistroKey distroKey, DataOperation action, String targetServer, long delay) {
        DistroKey distroKeyWithTarget = new DistroKey(distroKey.getResourceKey(), distroKey.getResourceType(), targetServer);
        DistroDelayTask distroDelayTask = new DistroDelayTask(distroKeyWithTarget, action, delay);
        distroTaskEngineHolder.getDelayTaskExecuteEngine().addTask(distroKeyWithTarget, distroDelayTask);
        if (Loggers.DISTRO.isDebugEnabled()) {
            Loggers.DISTRO.debug("[DISTRO-SCHEDULE] {} to {}", distroKey, targetServer);
        }
    }

    /**
     * Query data from specified server.
     * 從指定節點查詢資料
     * @param distroKey data key
     * @return data
     */
    public DistroData queryFromRemote(DistroKey distroKey) {
        if (null == distroKey.getTargetServer()) {
            Loggers.DISTRO.warn("[DISTRO] Can't query data from empty server");
            return null;
        }
        String resourceType = distroKey.getResourceType();
        DistroTransportAgent transportAgent = distroComponentHolder.findTransportAgent(resourceType);
        if (null == transportAgent) {
            Loggers.DISTRO.warn("[DISTRO] Can't find transport agent for key {}", resourceType);
            return null;
        }
        return transportAgent.getData(distroKey, distroKey.getTargetServer());
    }

    /**
     * Receive synced distro data, find processor to process.
     * 接收到同步資料,并查找處理器進行處理
     * @param distroData Received data
     * @return true if handle receive data successfully, otherwise false
     */
    public boolean onReceive(DistroData distroData) {
        Loggers.DISTRO.info("[DISTRO] Receive distro data type: {}, key: {}", distroData.getType(),
                distroData.getDistroKey());
        String resourceType = distroData.getDistroKey().getResourceType();
        DistroDataProcessor dataProcessor = distroComponentHolder.findDataProcessor(resourceType);
        if (null == dataProcessor) {
            Loggers.DISTRO.warn("[DISTRO] Can't find data process for received data {}", resourceType);
            return false;
        }
        return dataProcessor.processData(distroData);
    }

    /**
     * Receive verify data, find processor to process.
     * 接收到驗證資料,并查找處理器進行處理
     * @param distroData    verify data
     * @param sourceAddress source server address, might be get data from source server
     * @return true if verify data successfully, otherwise false
     */
    public boolean onVerify(DistroData distroData, String sourceAddress) {
        if (Loggers.DISTRO.isDebugEnabled()) {
            Loggers.DISTRO.debug("[DISTRO] Receive verify data type: {}, key: {}", distroData.getType(), distroData.getDistroKey());
        }
        String resourceType = distroData.getDistroKey().getResourceType();
        DistroDataProcessor dataProcessor = distroComponentHolder.findDataProcessor(resourceType);
        if (null == dataProcessor) {
            Loggers.DISTRO.warn("[DISTRO] Can't find verify data process for received data {}", resourceType);
            return false;
        }
        return dataProcessor.processVerifyData(distroData, sourceAddress);
    }

    /**
     * Query data of input distro key.
     * 根據條件查詢資料
     * @param distroKey key of data
     * @return data
     */
    public DistroData onQuery(DistroKey distroKey) {
        String resourceType = distroKey.getResourceType();
        DistroDataStorage distroDataStorage = distroComponentHolder.findDataStorage(resourceType);
        if (null == distroDataStorage) {
            Loggers.DISTRO.warn("[DISTRO] Can't find data storage for received key {}", resourceType);
            return new DistroData(distroKey, new byte[0]);
        }
        return distroDataStorage.getDistroData(distroKey);
    }

    /**
     * Query all datum snapshot.
     * 查詢所有快照資料
     * @param type datum type
     * @return all datum snapshot
     */
    public DistroData onSnapshot(String type) {
        DistroDataStorage distroDataStorage = distroComponentHolder.findDataStorage(type);
        if (null == distroDataStorage) {
            Loggers.DISTRO.warn("[DISTRO] Can't find data storage for received key {}", type);
            return new DistroData(new DistroKey("snapshot", type), new byte[0]);
        }
        return distroDataStorage.getDatumSnapshot();
    }
}

如果您認真從頭看到這里,相信您腦海中會記住一些關鍵字,比如TasksyncprocessorDistroData,所謂的Distro協議,不就是同步資料嘛,沒錯,它就是同步資料,在多個節點之間同步資料,通過DistroProtocol這個類不難發現,它實作了定時向其他節點報告狀態、首次從其他節點加載資料、同步資料到指定節點、獲取當前節點的快照資料,將這些功能組合在一起便可以實作多節點同步,因為所有節點都會做這些操作,

Distro協議資料物件

在整個互動程序中,是使用DistroData物件作為資料載體,它可以保存多種操作型別的任意資料,結構如圖:

在DistroKey中,包含了資源的標識、資源的型別,以及該資源所屬的節點,因此任何DistroData資料都能夠確定它是來自于那臺機器的什么型別的資料,在DataOperation中則定義了該資料將被用于什么操作,至于真正的資料型別,位元組陣列保證了它的兼容性,實際上DistroKey和Operation也能確定它將會是什么型別,

Distro協議重要角色

我們知道DistroData是作為Distro協議的互動物件,剩下的還有負責保存資料的組件、處理資料的組件、發送資料的組件,它們共同協作來完成整個協議流程,

存盤DistroData

DistroDataStorage 用于保存DistroData, 它有多種實作,用于處理不同型別的資料,實際上就是處理不同版本中的資料,

  • v1版本的實作:DistroDataStorageImpl
  • v2版本的實作:DistroClientDataProcessor

提示:
后續將不再刻意提及v1或者是v2的實作,默認以v2實作來分析,

資料的獲取發生在DistroDataStorage介面的getDistroData(DistroKey distroKey)getDatumSnapshot()getVerifyData()三個方法中,在v2版本中DistroClientDataProcessor實作了DistroDataStorage介面,提供DistroData的獲取功能,

// DistroClientDataProcessor.java

@Override
public DistroData getDistroData(DistroKey distroKey) {
	// 從Client管理器中獲取指定Client
	Client client = clientManager.getClient(distroKey.getResourceKey());
	if (null == client) {
		return null;
	}
	byte[] data = https://www.cnblogs.com/lukama/p/ApplicationUtils.getBean(Serializer.class).serialize(client.generateSyncData());
	return new DistroData(distroKey, data);
}

@Override
public DistroData getDatumSnapshot() {
	List datum = new LinkedList<>();
	// 從Client管理器中獲取所有Client
	for (String each : clientManager.allClientId()) {
		Client client = clientManager.getClient(each);
		if (null == client || !client.isEphemeral()) {
			continue;
		}
		datum.add(client.generateSyncData());
	}
	ClientSyncDatumSnapshot snapshot = new ClientSyncDatumSnapshot();
	snapshot.setClientSyncDataList(datum);
	byte[] data = ApplicationUtils.getBean(Serializer.class).serialize(snapshot);
	return new DistroData(new DistroKey(DataOperation.SNAPSHOT.name(), TYPE), data);
}

@Override
public List getVerifyData() {
	List result = new LinkedList<>();
	// 從Client管理器中獲取所有Client
	for (String each : clientManager.allClientId()) {
		Client client = clientManager.getClient(each);
		if (null == client || !client.isEphemeral()) {
			continue;
		}
		if (clientManager.isResponsibleClient(client)) {
			// TODO add revision for client.
			DistroClientVerifyInfo verifyData = new DistroClientVerifyInfo(client.getClientId(), 0);
			DistroKey distroKey = new DistroKey(client.getClientId(), TYPE);
			DistroData data = new DistroData(distroKey,
					ApplicationUtils.getBean(Serializer.class).serialize(verifyData));
			data.setType(DataOperation.VERIFY);
			result.add(data);
		}
	}
	return result;
}

通過v2版本的資料存盤實作可以發現,它并沒有直接去保存資料,而是從ClientManager內部獲取,

處理DistroData

DistroDataProcessor 用于處理DistroData,資料的處理發生在processData(DistroData distroData)processVerifyData(DistroData distroData, String sourceAddress)processSnapshot(DistroData distroData)三個方法中,在v2版本中DistroClientDataProcessor實作了DistroDataProcessor介面,提供DistroData的處理能力,

// DistroClientDataProcessor.java

@Override
public boolean processData(DistroData distroData) {
	switch (distroData.getType()) {
		case ADD:
		case CHANGE:
			ClientSyncData clientSyncData = https://www.cnblogs.com/lukama/p/ApplicationUtils.getBean(Serializer.class).deserialize(distroData.getContent(), ClientSyncData.class);
			handlerClientSyncData(clientSyncData);
			return true;
		case DELETE:
			String deleteClientId = distroData.getDistroKey().getResourceKey();
			Loggers.DISTRO.info("[Client-Delete] Received distro client sync data {}", deleteClientId);
			clientManager.clientDisconnected(deleteClientId);
			return true;
		default:
			return false;
	}
}

@Override
public boolean processVerifyData(DistroData distroData, String sourceAddress) {
	DistroClientVerifyInfo verifyData = https://www.cnblogs.com/lukama/p/ApplicationUtils.getBean(Serializer.class).deserialize(distroData.getContent(), DistroClientVerifyInfo.class);
	if (clientManager.verifyClient(verifyData.getClientId())) {
		return true;
	}
	Loggers.DISTRO.info("client {} is invalid, get new client from {}", verifyData.getClientId(), sourceAddress);
	return false;
}

@Override
public boolean processSnapshot(DistroData distroData) {
	ClientSyncDatumSnapshot snapshot = ApplicationUtils.getBean(Serializer.class).deserialize(distroData.getContent(), ClientSyncDatumSnapshot.class);
	for (ClientSyncData each : snapshot.getClientSyncDataList()) {
		handlerClientSyncData(each);
	}
	return true;
}

發送DistroData

DistroTransportAgent用于傳輸DistroData,v2版本中DistroClientTransportAgent實作了DistroTransportAgent介面,提供DistroData的發送能力,

/**
 * Distro transport agent for v2.
 * v2版本的DistroData傳輸代理
 * @author xiweng.yy
 */
public class DistroClientTransportAgent implements DistroTransportAgent {

    private final ClusterRpcClientProxy clusterRpcClientProxy;

    private final ServerMemberManager memberManager;

    public DistroClientTransportAgent(ClusterRpcClientProxy clusterRpcClientProxy,
            ServerMemberManager serverMemberManager) {
        this.clusterRpcClientProxy = clusterRpcClientProxy;
        this.memberManager = serverMemberManager;
    }

    /**
     * 當前實作支持回呼
     * @return
     */
    @Override
    public boolean supportCallbackTransport() {
        return true;
    }

    /**
     * 向指定節點發送同步資料
     * @param data         data
     * @param targetServer target server
     * @return
     */
    @Override
    public boolean syncData(DistroData data, String targetServer) {
        if (isNoExistTarget(targetServer)) {
            return true;
        }
        DistroDataRequest request = new DistroDataRequest(data, data.getType());
        Member member = memberManager.find(targetServer);
        if (checkTargetServerStatusUnhealthy(member)) {
            Loggers.DISTRO.warn("[DISTRO] Cancel distro sync caused by target server {} unhealthy", targetServer);
            return false;
        }
        try {
            Response response = clusterRpcClientProxy.sendRequest(member, request);
            return checkResponse(response);
        } catch (NacosException e) {
            Loggers.DISTRO.error("[DISTRO-FAILED] Sync distro data failed! ", e);
        }
        return false;
    }

    /**
     * 向指定節點發送回同步資料(支持回呼)
     * @param data         data
     * @param targetServer target server
     * @param callback     callback
     */
    @Override
    public void syncData(DistroData data, String targetServer, DistroCallback callback) {
        if (isNoExistTarget(targetServer)) {
            callback.onSuccess();
        }
        DistroDataRequest request = new DistroDataRequest(data, data.getType());
        Member member = memberManager.find(targetServer);
        try {
            clusterRpcClientProxy.asyncRequest(member, request, new DistroRpcCallbackWrapper(callback, member));
        } catch (NacosException nacosException) {
            callback.onFailed(nacosException);
        }
    }

    /**
     * 向指定節點發送驗證資料
     * @param verifyData   verify data
     * @param targetServer target server
     * @return
     */
    @Override
    public boolean syncVerifyData(DistroData verifyData, String targetServer) {
        if (isNoExistTarget(targetServer)) {
            return true;
        }
        // replace target server as self server so that can callback.
        verifyData.getDistroKey().setTargetServer(memberManager.getSelf().getAddress());
        DistroDataRequest request = new DistroDataRequest(verifyData, DataOperation.VERIFY);
        Member member = memberManager.find(targetServer);
        if (checkTargetServerStatusUnhealthy(member)) {
            Loggers.DISTRO.warn("[DISTRO] Cancel distro verify caused by target server {} unhealthy", targetServer);
            return false;
        }
        try {
            Response response = clusterRpcClientProxy.sendRequest(member, request);
            return checkResponse(response);
        } catch (NacosException e) {
            Loggers.DISTRO.error("[DISTRO-FAILED] Verify distro data failed! ", e);
        }
        return false;
    }

    /**
     * 向指定節點發送驗證資料(支持回呼)
     * @param verifyData   verify data
     * @param targetServer target server
     * @param callback     callback
     */
    @Override
    public void syncVerifyData(DistroData verifyData, String targetServer, DistroCallback callback) {
        // 若此節點不在當前節點快取中,直接回傳,因為可能下線、或者過期,不需要驗證了
        if (isNoExistTarget(targetServer)) {
            callback.onSuccess();
        }
        // 構建請求物件
        DistroDataRequest request = new DistroDataRequest(verifyData, DataOperation.VERIFY);
        Member member = memberManager.find(targetServer);
        try {
            // 創建一個回呼物件(Wrapper實作了RequestCallBack介面)
            DistroVerifyCallbackWrapper wrapper = new DistroVerifyCallbackWrapper(targetServer,
                    verifyData.getDistroKey().getResourceKey(), callback, member);
            // 使用集群Rpc請求物件發送異步任務
            clusterRpcClientProxy.asyncRequest(member, request, wrapper);
        } catch (NacosException nacosException) {
            callback.onFailed(nacosException);
        }
    }

    /**
     * 從指定節點獲取資料
     * @param key          key of data
     * @param targetServer target server
     * @return
     */
    @Override
    public DistroData getData(DistroKey key, String targetServer) {
        Member member = memberManager.find(targetServer);
        if (checkTargetServerStatusUnhealthy(member)) {
            throw new DistroException(
                    String.format("[DISTRO] Cancel get snapshot caused by target server %s unhealthy", targetServer));
        }
        DistroDataRequest request = new DistroDataRequest();
        DistroData distroData = https://www.cnblogs.com/lukama/p/new DistroData();
        distroData.setDistroKey(key);
        distroData.setType(DataOperation.QUERY);
        request.setDistroData(distroData);
        request.setDataOperation(DataOperation.QUERY);
        try {
            Response response = clusterRpcClientProxy.sendRequest(member, request);
            if (checkResponse(response)) {
                return ((DistroDataResponse) response).getDistroData();
            } else {
                throw new DistroException(
                        String.format("[DISTRO-FAILED] Get data request to %s failed, code: %d, message: %s",
                                targetServer, response.getErrorCode(), response.getMessage()));
            }
        } catch (NacosException e) {
            throw new DistroException("[DISTRO-FAILED] Get distro data failed! ", e);
        }
    }

    /**
     * 從指定節點獲取快照資料
     * @param targetServer target server.
     * @return
     */
    @Override
    public DistroData getDatumSnapshot(String targetServer) {
        Member member = memberManager.find(targetServer);
        if (checkTargetServerStatusUnhealthy(member)) {
            throw new DistroException(
                    String.format("[DISTRO] Cancel get snapshot caused by target server %s unhealthy", targetServer));
        }
        DistroDataRequest request = new DistroDataRequest();
        request.setDataOperation(DataOperation.SNAPSHOT);
        try {
            Response response = clusterRpcClientProxy.sendRequest(member, request);
            if (checkResponse(response)) {
                return ((DistroDataResponse) response).getDistroData();
            } else {
                throw new DistroException(
                        String.format("[DISTRO-FAILED] Get snapshot request to %s failed, code: %d, message: %s",
                                targetServer, response.getErrorCode(), response.getMessage()));
            }
        } catch (NacosException e) {
            throw new DistroException("[DISTRO-FAILED] Get distro snapshot failed! ", e);
        }
    }

    private boolean isNoExistTarget(String target) {
        return !memberManager.hasMember(target);
    }

    private boolean checkTargetServerStatusUnhealthy(Member member) {
        return null == member || !NodeState.UP.equals(member.getState());
    }

    private boolean checkResponse(Response response) {
        return ResponseCode.SUCCESS.getCode() == response.getResultCode();
    }

    /**
     * rpc請求回呼包裝器
     */
    private class DistroRpcCallbackWrapper implements RequestCallBack<Response> {

        private final DistroCallback distroCallback;

        private final Member member;

        public DistroRpcCallbackWrapper(DistroCallback distroCallback, Member member) {
            this.distroCallback = distroCallback;
            this.member = member;
        }

        @Override
        public Executor getExecutor() {
            return GlobalExecutor.getCallbackExecutor();
        }

        @Override
        public long getTimeout() {
            return DistroConfig.getInstance().getSyncTimeoutMillis();
        }

        @Override
        public void onResponse(Response response) {
            if (checkResponse(response)) {
                NamingTpsMonitor.distroSyncSuccess(member.getAddress(), member.getIp());
                distroCallback.onSuccess();
            } else {
                NamingTpsMonitor.distroSyncFail(member.getAddress(), member.getIp());
                distroCallback.onFailed(null);
            }
        }

        @Override
        public void onException(Throwable e) {
            distroCallback.onFailed(e);
        }
    }

    /**
     * 驗證資料回呼包裝器
     */
    private class DistroVerifyCallbackWrapper implements RequestCallBack<Response> {

        private final String targetServer;

        private final String clientId;

        private final DistroCallback distroCallback;

        private final Member member;

        private DistroVerifyCallbackWrapper(String targetServer, String clientId, DistroCallback distroCallback,
                Member member) {
            this.targetServer = targetServer;
            this.clientId = clientId;
            this.distroCallback = distroCallback;
            this.member = member;
        }

        @Override
        public Executor getExecutor() {
            return GlobalExecutor.getCallbackExecutor();
        }

        @Override
        public long getTimeout() {
            return DistroConfig.getInstance().getVerifyTimeoutMillis();
        }

        @Override
        public void onResponse(Response response) {
            if (checkResponse(response)) {
                NamingTpsMonitor.distroVerifySuccess(member.getAddress(), member.getIp());
                distroCallback.onSuccess();
            } else {
                Loggers.DISTRO.info("Target {} verify client {} failed, sync new client", targetServer, clientId);
				// 驗證失敗之后發布事件
                NotifyCenter.publishEvent(new ClientEvent.ClientVerifyFailedEvent(clientId, targetServer));
                NamingTpsMonitor.distroVerifyFail(member.getAddress(), member.getIp());
                distroCallback.onFailed(null);
            }
        }

        @Override
        public void onException(Throwable e) {
            distroCallback.onFailed(e);
        }
    }
}


不管同步資料的操作型別是什么,最終發送資料使用的是ClusterRpcClientProxy物件,

以上3個組件是實作Distro協議中重要的一環,后續關于Distro協議的邏輯將全部圍繞這三個組件進行

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/296178.html

標籤:其他

上一篇:Nacos 2.0原始碼分析-健康檢查機制

下一篇:搭建前端監控系統

標籤雲
其他(157675) Python(38076) JavaScript(25376) Java(17977) C(15215) 區塊鏈(8255) C#(7972) AI(7469) 爪哇(7425) MySQL(7132) html(6777) 基礎類(6313) sql(6102) 熊猫(6058) PHP(5869) 数组(5741) R(5409) Linux(5327) 反应(5209) 腳本語言(PerlPython)(5129) 非技術區(4971) Android(4554) 数据框(4311) css(4259) 节点.js(4032) C語言(3288) json(3245) 列表(3129) 扑(3119) C++語言(3117) 安卓(2998) 打字稿(2995) VBA(2789) Java相關(2746) 疑難問題(2699) 细绳(2522) 單片機工控(2479) iOS(2429) ASP.NET(2402) MongoDB(2323) 麻木的(2285) 正则表达式(2254) 字典(2211) 循环(2198) 迅速(2185) 擅长(2169) 镖(2155) 功能(1967) .NET技术(1958) Web開發(1951) python-3.x(1918) HtmlCss(1915) 弹簧靴(1913) C++(1909) xml(1889) PostgreSQL(1872) .NETCore(1853) 谷歌表格(1846) Unity3D(1843) for循环(1842)

熱門瀏覽
  • 網閘典型架構簡述

    網閘架構一般分為兩種:三主機的三系統架構網閘和雙主機的2+1架構網閘。 三主機架構分別為內端機、外端機和仲裁機。三機無論從軟體和硬體上均各自獨立。首先從硬體上來看,三機都用各自獨立的主板、記憶體及存盤設備。從軟體上來看,三機有各自獨立的作業系統。這樣能達到完全的三機獨立。對于“2+1”系統,“2”分為 ......

    uj5u.com 2020-09-10 02:00:44 more
  • 如何從xshell上傳檔案到centos linux虛擬機里

    如何從xshell上傳檔案到centos linux虛擬機里及:虛擬機CentOs下執行 yum -y install lrzsz命令,出現錯誤:鏡像無法找到軟體包 前言 一、安裝lrzsz步驟 二、上傳檔案 三、遇到的問題及解決方案 總結 前言 提示:其實很簡單,往虛擬機上安裝一個上傳檔案的工具 ......

    uj5u.com 2020-09-10 02:00:47 more
  • 一、SQLMAP入門

    一、SQLMAP入門 1、判斷是否存在注入 sqlmap.py -u 網址/id=1 id=1不可缺少。當注入點后面的引數大于兩個時。需要加雙引號, sqlmap.py -u "網址/id=1&uid=1" 2、判斷文本中的請求是否存在注入 從文本中加載http請求,SQLMAP可以從一個文本檔案中 ......

    uj5u.com 2020-09-10 02:00:50 more
  • Metasploit 簡單使用教程

    metasploit 簡單使用教程 浩先生, 2020-08-28 16:18:25 分類專欄: kail 網路安全 linux 文章標簽: linux資訊安全 編輯 著作權 metasploit 使用教程 前言 一、Metasploit是什么? 二、準備作業 三、具體步驟 前言 Msfconsole ......

    uj5u.com 2020-09-10 02:00:53 more
  • 游戲逆向之驅動層與用戶層通訊

    驅動層代碼: #pragma once #include <ntifs.h> #define add_code CTL_CODE(FILE_DEVICE_UNKNOWN,0x800,METHOD_BUFFERED,FILE_ANY_ACCESS) /* 更多游戲逆向視頻www.yxfzedu.com ......

    uj5u.com 2020-09-10 02:00:56 more
  • 北斗電力時鐘(北斗授時服務器)讓網路資料更精準

    北斗電力時鐘(北斗授時服務器)讓網路資料更精準 北斗電力時鐘(北斗授時服務器)讓網路資料更精準 京準電子科技官微——ahjzsz 近幾年,資訊技術的得了快速發展,互聯網在逐漸普及,其在人們生活和生產中都得到了廣泛應用,并且取得了不錯的應用效果。計算機網路資訊在電力系統中的應用,一方面使電力系統的運行 ......

    uj5u.com 2020-09-10 02:01:03 more
  • 【CTF】CTFHub 技能樹 彩蛋 writeup

    ?碎碎念 CTFHub:https://www.ctfhub.com/ 筆者入門CTF時時剛開始刷的是bugku的舊平臺,后來才有了CTFHub。 感覺不論是網頁UI設計,還是題目質量,賽事跟蹤,工具軟體都做得很不錯。 而且因為獨到的金幣制度的確讓人有一種想去刷題賺金幣的感覺。 個人還是非常喜歡這個 ......

    uj5u.com 2020-09-10 02:04:05 more
  • 02windows基礎操作

    我學到了一下幾點 Windows系統目錄結構與滲透的作用 常見Windows的服務詳解 Windows埠詳解 常用的Windows注冊表詳解 hacker DOS命令詳解(net user / type /md /rd/ dir /cd /net use copy、批處理 等) 利用dos命令制作 ......

    uj5u.com 2020-09-10 02:04:18 more
  • 03.Linux基礎操作

    我學到了以下幾點 01Linux系統介紹02系統安裝,密碼啊破解03Linux常用命令04LAMP 01LINUX windows: win03 8 12 16 19 配置不繁瑣 Linux:redhat,centos(紅帽社區版),Ubuntu server,suse unix:金融機構,證券,銀 ......

    uj5u.com 2020-09-10 02:04:30 more
  • 05HTML

    01HTML介紹 02頭部標簽講解03基礎標簽講解04表單標簽講解 HTML前段語言 js1.了解代碼2.根據代碼 懂得挖掘漏洞 (POST注入/XSS漏洞上傳)3.黑帽seo 白帽seo 客戶網站被黑帽植入劫持代碼如何處理4.熟悉html表單 <html><head><title>TDK標題,描述 ......

    uj5u.com 2020-09-10 02:04:36 more
最新发布
  • 2023年最新微信小程式抓包教程

    01 開門見山 隔一個月發一篇文章,不過分。 首先回顧一下《微信系結手機號資料庫被脫庫事件》,我也是第一時間得知了這個訊息,然后跟蹤了整件事情的經過。下面是這起事件的相關截圖以及近日流出的一萬條資料樣本: 個人認為這件事也沒什么,還不如關注一下之前45億快遞資料查詢渠道疑似在近日復活的訊息。 訊息是 ......

    uj5u.com 2023-04-20 08:48:24 more
  • web3 產品介紹:metamask 錢包 使用最多的瀏覽器插件錢包

    Metamask錢包是一種基于區塊鏈技術的數字貨幣錢包,它允許用戶在安全、便捷的環境下管理自己的加密資產。Metamask錢包是以太坊生態系統中最流行的錢包之一,它具有易于使用、安全性高和功能強大等優點。 本文將詳細介紹Metamask錢包的功能和使用方法。 一、 Metamask錢包的功能 數字資 ......

    uj5u.com 2023-04-20 08:47:46 more
  • vulnhub_Earth

    前言 靶機地址->>>vulnhub_Earth 攻擊機ip:192.168.20.121 靶機ip:192.168.20.122 參考文章 https://www.cnblogs.com/Jing-X/archive/2022/04/03/16097695.html https://www.cnb ......

    uj5u.com 2023-04-20 07:46:20 more
  • 從4k到42k,軟體測驗工程師的漲薪史,給我看哭了

    清明節一過,盲猜大家已經無心上班,在數著日子準備過五一,但一想到銀行卡里的余額……瞬間心情就不美麗了。最近,2023年高校畢業生就業調查顯示,本科畢業月平均起薪為5825元。調查一出,便有很多同學表示自己又被平均了。看著這一資料,不免讓人想到前不久中國青年報的一項調查:近六成大學生認為畢業10年內會 ......

    uj5u.com 2023-04-20 07:44:00 more
  • 最新版本 Stable Diffusion 開源 AI 繪畫工具之中文自動提詞篇

    🎈 標簽生成器 由于輸入正向提示詞 prompt 和反向提示詞 negative prompt 都是使用英文,所以對學習母語的我們非常不友好 使用網址:https://tinygeeker.github.io/p/ai-prompt-generator 這個網址是為了讓大家在使用 AI 繪畫的時候 ......

    uj5u.com 2023-04-20 07:43:36 more
  • 漫談前端自動化測驗演進之路及測驗工具分析

    隨著前端技術的不斷發展和應用程式的日益復雜,前端自動化測驗也在不斷演進。隨著 Web 應用程式變得越來越復雜,自動化測驗的需求也越來越高。如今,自動化測驗已經成為 Web 應用程式開發程序中不可或缺的一部分,它們可以幫助開發人員更快地發現和修復錯誤,提高應用程式的性能和可靠性。 ......

    uj5u.com 2023-04-20 07:43:16 more
  • CANN開發實踐:4個DVPP記憶體問題的典型案例解讀

    摘要:由于DVPP媒體資料處理功能對存放輸入、輸出資料的記憶體有更高的要求(例如,記憶體首地址128位元組對齊),因此需呼叫專用的記憶體申請介面,那么本期就分享幾個關于DVPP記憶體問題的典型案例,并給出原因分析及解決方法。 本文分享自華為云社區《FAQ_DVPP記憶體問題案例》,作者:昇騰CANN。 DVPP ......

    uj5u.com 2023-04-20 07:43:03 more
  • msf學習

    msf學習 以kali自帶的msf為例 一、msf核心模塊與功能 msf模塊都放在/usr/share/metasploit-framework/modules目錄下 1、auxiliary 輔助模塊,輔助滲透(埠掃描、登錄密碼爆破、漏洞驗證等) 2、encoders 編碼器模塊,主要包含各種編碼 ......

    uj5u.com 2023-04-20 07:42:59 more
  • Halcon軟體安裝與界面簡介

    1. 下載Halcon17版本到到本地 2. 雙擊安裝包后 3. 步驟如下 1.2 Halcon軟體安裝 界面分為四大塊 1. Halcon的五個助手 1) 影像采集助手:與相機連接,設定相機引數,采集影像 2) 標定助手:九點標定或是其它的標定,生成標定檔案及內參外參,可以將像素單位轉換為長度單位 ......

    uj5u.com 2023-04-20 07:42:17 more
  • 在MacOS下使用Unity3D開發游戲

    第一次發博客,先發一下我的游戲開發環境吧。 去年2月份買了一臺MacBookPro2021 M1pro(以下簡稱mbp),這一年來一直在用mbp開發游戲。我大致分享一下我的開發工具以及使用體驗。 1、Unity 官網鏈接: https://unity.cn/releases 我一般使用的Apple ......

    uj5u.com 2023-04-20 07:40:19 more