主頁 >  其他 > Nacos 2.0原始碼分析-事件發布機制

Nacos 2.0原始碼分析-事件發布機制

2021-08-31 20:06:17 其他

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

Nacos的服務注冊、服務變更等功能都是通過事件發布來通知的,搞清楚事件發布訂閱的機制,有利于理解業務的流程走向,本文將淺顯的分析Nacos中的事件發布訂閱實作,

事件(Event)

常規事件(Event)

package com.alibaba.nacos.common.notify;

public abstract class Event implements Serializable {
    
    private static final AtomicLong SEQUENCE = new AtomicLong(0);
    
    private final long sequence = SEQUENCE.getAndIncrement();
    
    /**
     * Event sequence number, which can be used to handle the sequence of events.
     *
     * @return sequence num, It's best to make sure it's monotone.
     */
    public long sequence() {
        return sequence;
    }
}

在事件抽象類中定義了一個事件的序列號,它是自增的,用于區分事件執行的前后順序,它是由DefaultPublisher來處理,

慢事件(SlowEvent)

之所以稱之為慢事件,可能因為所有的事件都共享同一個佇列吧,

package com.alibaba.nacos.common.notify;

/**
 * This event share one event-queue.
 * @author <a href="mailto:[email protected]">liaochuntao</a>
 * @author zongtanghu
 */
@SuppressWarnings("PMD.AbstractClassShouldStartWithAbstractNamingRule")
public abstract class SlowEvent extends Event {
    
    @Override
    public long sequence() {
        return 0;
    }
}

提示:
SlowEvent可以共享一個事件佇列,也就是一個發布者可以同時管理多個事件的發布(區別于DefaultPublisher只能管理一個事件),

訂閱者(Subscriber)

單事件訂閱者

這里的單事件訂閱者指的是當前的訂閱者只能訂閱一種型別的事件,

package com.alibaba.nacos.common.notify.listener;

/**
 * An abstract subscriber class for subscriber interface.
 * @author <a href="mailto:[email protected]">liaochuntao</a>
 * @author zongtanghu
 */
@SuppressWarnings("PMD.AbstractClassShouldStartWithAbstractNamingRule")
public abstract class Subscriber<T extends Event> {
    
    /**
     * Event callback.
     * 事件處理入口,由對應的事件發布器呼叫
     * @param event {@link Event}
     */
    public abstract void onEvent(T event);
    
    /**
     * Type of this subscriber's subscription.
     * 訂閱的事件型別
     * @return Class which extends {@link Event}
     */
    public abstract Class<? extends Event> subscribeType();
    
    /**
     * It is up to the listener to determine whether the callback is asynchronous or synchronous.
     * 執行緒執行器,由具體的實作類來決定是異步還是同步呼叫
     * @return {@link Executor}
     */
    public Executor executor() {
        return null;
    }
    
    /**
     * Whether to ignore expired events.
     * 是否忽略過期事件
     * @return default value is {@link Boolean#FALSE}
     */
    public boolean ignoreExpireEvent() {
        return false;
    }
}

這是默認的訂閱者物件,默認情況下一個訂閱者只能訂閱一個型別的事件,

多事件訂閱者

package com.alibaba.nacos.common.notify.listener;


/**
 * Subscribers to multiple events can be listened to.
 *
 * @author <a href="mailto:[email protected]">liaochuntao</a>
 * @author zongtanghu
 */
@SuppressWarnings("PMD.AbstractClassShouldStartWithAbstractNamingRule")
public abstract class SmartSubscriber extends Subscriber {
    
    /**
     * Returns which event type are smartsubscriber interested in.
     * 區別于父類,這里支持多個事件型別
     * @return The interestd event types.
     */
    public abstract List<Class<? extends Event>> subscribeTypes();
    
    @Override
    public final Class<? extends Event> subscribeType() {
		// 采用final修飾,禁止使用單一事件屬性
        return null;
    }
    
    @Override
    public final boolean ignoreExpireEvent() {
        return false;
    }
}

提示
SmartSubscriber和Subscriber的區別是一個可以訂閱多個事件,一個只能訂閱一個事件,處理它們的發布者也不同,

發布者(Publisher)

發布者指的是Nacos中的事件發布者,頂級介面為EventPublisher,

package com.alibaba.nacos.common.notify;

/**
 * Event publisher.
 *
 * @author <a href="mailto:[email protected]">liaochuntao</a>
 * @author zongtanghu
 */
public interface EventPublisher extends Closeable {
    
    /**
     * Initializes the event publisher.
     * 初始化事件發布者
     * @param type       {@link Event >}
     * @param bufferSize Message staging queue size
     */
    void init(Class<? extends Event> type, int bufferSize);
    
    /**
     * The number of currently staged events.
     * 當前暫存的事件數量
     * @return event size
     */
    long currentEventSize();
    
    /**
     * Add listener.
     * 添加訂閱者
     * @param subscriber {@link Subscriber}
     */
    void addSubscriber(Subscriber subscriber);
    
    /**
     * Remove listener.
     * 移除訂閱者
     * @param subscriber {@link Subscriber}
     */
    void removeSubscriber(Subscriber subscriber);
    
    /**
     * publish event.
     * 發布事件
     * @param event {@link Event}
     * @return publish event is success
     */
    boolean publish(Event event);
    
    /**
     * Notify listener.
     * 通知訂閱者
     * @param subscriber {@link Subscriber}
     * @param event      {@link Event}
     */
    void notifySubscriber(Subscriber subscriber, Event event);
    
}

發布者的主要功能就是新增訂閱者、通知訂閱者,目前有兩種型別的發布者分別是DefaultPublisher和DefaultSharePublisher,

單事件發布者(DefaultPublisher)

一個發布者實體只能處理一種型別的事件,

public class DefaultPublisher extends Thread implements EventPublisher {
	
	// 發布者是否初始化完畢
	private volatile boolean initialized = false;
	// 是否關閉了發布者
	private volatile boolean shutdown = false;
	// 事件的型別
	private Class<? extends Event> eventType;
	// 訂閱者串列
	protected final ConcurrentHashSet<Subscriber> subscribers = new ConcurrentHashSet<Subscriber>();
	// 佇列最大容量
	private int queueMaxSize = -1;
	// 佇列型別
	private BlockingQueue<Event> queue;
	// 最后一個事件的序列號
	protected volatile Long lastEventSequence = -1L;
	// 事件序列號更新物件,用于更新原子屬性lastEventSequence
	private static final AtomicReferenceFieldUpdater<DefaultPublisher, Long> UPDATER = AtomicReferenceFieldUpdater.newUpdater(DefaultPublisher.class, Long.class, "lastEventSequence");
}

發布者的初始化

public void init(Class<? extends Event> type, int bufferSize) {
	setDaemon(true);
	setName("nacos.publisher-" + type.getName());
	this.eventType = type;
	this.queueMaxSize = bufferSize;
	this.queue = new ArrayBlockingQueue<Event>(bufferSize);
	start();
}

在初始化方法中,將其設定為了守護執行緒,意味著它將持續運行(它需要持續監控內部的事件佇列),傳入的type屬性為當前發布者需要處理的事件型別,設定當前執行緒的名稱以事件型別為區分,它將會以多個執行緒的形式存在,每個執行緒代表一種事件型別的發布者,后面初始化了佇列的長度,最后呼叫啟動方法完成當前執行緒的啟動,

發布者執行緒啟動

public synchronized void start() {
	if (!initialized) {
		// start just called once
		super.start();
		if (queueMaxSize == -1) {
			queueMaxSize = ringBufferSize;
		}
		initialized = true;
	}
}

直接呼叫了Thread的start方法開啟守護執行緒,并設定初始化狀態為true,根據java執行緒的啟動方式,呼叫start方法之后start方法是會呼叫run方法的,

public void run() {
	openEventHandler();
}

void openEventHandler() {
	try {
		
		// This variable is defined to resolve the problem which message overstock in the queue.
		int waitTimes = 60;
		// To ensure that messages are not lost, enable EventHandler when
		// waiting for the first Subscriber to register
		for (; ; ) {
			// 執行緒終止條件判斷
			if (shutdown || hasSubscriber() || waitTimes <= 0) {
				break;
			}
			// 執行緒休眠1秒
			ThreadUtils.sleep(1000L);
			// 等待次數減1
			waitTimes--;
		}
		
		for (; ; ) {
			// 執行緒終止條件判斷
			if (shutdown) {
				break;
			}
			// 從佇列取出事件
			final Event event = queue.take();
			// 接收事件
			receiveEvent(event);
			// 更新事件序列號
			UPDATER.compareAndSet(this, lastEventSequence, Math.max(lastEventSequence, event.sequence()));
		}
	} catch (Throwable ex) {
		LOGGER.error("Event listener exception : {}", ex);
	}
}

在run方法中呼叫了openEventHandler()方法,那發布者的實際作業原理就存在于這個方法內部,在首次啟動的時候會等待1分鐘,然后再進行訊息消費,

接收并發布事件

這里的接收事件指的是接收通知中心發過來的事件,發布給訂閱者,

void receiveEvent(Event event) {
	// 獲取當前事件的序列號,它是自增的
	final long currentEventSequence = event.sequence();
	
	// 通知所有訂閱了該事件的訂閱者
	// Notification single event listener
	for (Subscriber subscriber : subscribers) {
		// 判斷訂閱者是否忽略事件過期,判斷當前事件是否被處理過(lastEventSequence初始化的值為-1,而Event的sequence初始化的值為0)
		// Whether to ignore expiration events
		if (subscriber.ignoreExpireEvent() && lastEventSequence > currentEventSequence) {
			LOGGER.debug("[NotifyCenter] the {} is unacceptable to this subscriber, because had expire", event.getClass());
			continue;
		}
		
		// Because unifying smartSubscriber and subscriber, so here need to think of compatibility.
		// Remove original judge part of codes.
		notifySubscriber(subscriber, event);
	}
}

public void notifySubscriber(final Subscriber subscriber, final Event event) {
	
	LOGGER.debug("[NotifyCenter] the {} will received by {}", event, subscriber);
	
	// 為每個訂閱者創建一個Runnable物件
	final Runnable job = () -> subscriber.onEvent(event);
	// 使用訂閱者的執行緒執行器
	final Executor executor = subscriber.executor();
	// 若訂閱者沒有自己的執行器,則直接執行run方法啟動訂閱者消費執行緒
	if (executor != null) {
		executor.execute(job);
	} else {
		try {
			job.run();
		} catch (Throwable e) {
			LOGGER.error("Event callback exception: ", e);
		}
	}
}

外部呼叫發布事件

前面的發布事件是指從佇列內部獲取事件并通知訂閱者,這里的發布事件區別在于它是開放給外部呼叫者,接收統一通知中心的事件并放入佇列中的,

public boolean publish(Event event) {
	checkIsStart();
	boolean success = this.queue.offer(event);
	if (!success) {
		LOGGER.warn("Unable to plug in due to interruption, synchronize sending time, event : {}", event);
		receiveEvent(event);
		return true;
	}
	return true;
}

在放入佇列成功的時候直接回傳,若放入佇列失敗,則是直接同步發送事件給訂閱者,不經過佇列,這里的同步我認為的是從呼叫者到發布者呼叫訂閱者之間是同步的,若佇列可用,則是呼叫者到入佇列就完成了本次呼叫,不需要等待回圈通知訂閱者,使用佇列解耦無疑會提升通知中心的作業效率,

總體來說就是一個發布者內部維護一個BlockingQueue,在實作上使用了ArrayBlockingQueue,它是一個有界阻塞佇列,元素先進先出,并且使用非公平模式提升性能,意味著等待消費的訂閱者執行順序將得不到保障(業務需求沒有這種順序性要求),同時也維護了一個訂閱者集合(他們都訂閱了同一個事件型別),在死回圈中不斷從ArrayBlockingQueue中獲取資料來回圈通知每一個訂閱者,也就是呼叫訂閱者的onEvent()方法,

多事件發布者(DefaultSharePublisher)

用于發布SlowEvent事件并通知所有訂閱了該事件的訂閱者,

public class DefaultSharePublisher extends DefaultPublisher {
	// 用于保存事件型別為SlowEvent的訂閱者,一個事件型別對應多個訂閱者
	private final Map<Class<? extends SlowEvent>, Set<Subscriber>> subMappings = new ConcurrentHashMap<Class<? extends SlowEvent>, Set<Subscriber>>();
    // 可重入鎖
    private final Lock lock = new ReentrantLock();
}

它繼承了DefaultPublisher,意味著它將擁有其所有的特性,從subMappings屬性來看,這個發布器是支持多個SlowEvent事件的,DefaultSharePublisher多載了DefaultPublisher的addSubscriber()和removeSubscriber()方法,用于處理多事件型別的情形,

添加訂閱者:

public void addSubscriber(Subscriber subscriber, Class<? extends Event> subscribeType) {
	
	// 將事件型別轉換為當前發布者支持的型別
	// Actually, do a classification based on the slowEvent type.
	Class<? extends SlowEvent> subSlowEventType = (Class<? extends SlowEvent>) subscribeType;
	// 添加到父類的訂閱者串列中,為何要添加呢?因為它需要使用父類的佇列消費邏輯
	// For adding to parent class attributes synchronization.
	subscribers.add(subscriber);
	// 為多個操作加鎖
	lock.lock();
	try {
		// 首先從事件訂閱串列里面獲取當前事件對應的訂閱者集合
		Set<Subscriber> sets = subMappings.get(subSlowEventType);
		// 若沒有訂閱者,則新增當前訂閱者
		if (sets == null) {
			Set<Subscriber> newSet = new ConcurrentHashSet<Subscriber>();
			newSet.add(subscriber);
			subMappings.put(subSlowEventType, newSet);
			return;
		}
		// 若當前事件訂閱者串列不為空,則插入,因為使用的是Set集合因此可以避免重復資料
		sets.add(subscriber);
	} finally {
		// 別忘了解鎖
		lock.unlock();
	}
}

提示:
Set newSet = new ConcurrentHashSet(); 它這里實際上使用的是自己實作的ConcurrentHashSet,它內部使用了ConcurrentHashMap來實作存盤,
在ConcurrentHashSet.add()方法的實作上,它以當前插入的Subscriber物件為key,以一個Boolean值占位:map.putIfAbsent(o, Boolean.TRUE),

事件型別和訂閱者的存盤狀態為:
EventType1 -> {Subscriber1, Subscriber2, Subscriber3...}
EventType2 -> {Subscriber1, Subscriber2, Subscriber3...}
EventType3 -> {Subscriber1, Subscriber2, Subscriber3...}
感興趣的可以自己查閱一下原始碼,

移除訂閱者

public void removeSubscriber(Subscriber subscriber, Class<? extends Event> subscribeType) {
	// 轉換型別
	// Actually, do a classification based on the slowEvent type.
	Class<? extends SlowEvent> subSlowEventType = (Class<? extends SlowEvent>) subscribeType;
	// 先移除父類中的訂閱者
	// For removing to parent class attributes synchronization.
	subscribers.remove(subscriber);
	// 加鎖
	lock.lock();
	try {
		// 移除指定事件的指定訂閱者
		Set<Subscriber> sets = subMappings.get(subSlowEventType);
		
		if (sets != null) {
			sets.remove(subscriber);
		}
	} finally {
		// 解鎖
		lock.unlock();
	}
}

接收事件

@Override
public void receiveEvent(Event event) {
	// 獲取當前事件的序列號
	final long currentEventSequence = event.sequence();
	// 獲取事件的型別,轉換為當前發布器支持的事件
	// get subscriber set based on the slow EventType.
	final Class<? extends SlowEvent> slowEventType = (Class<? extends SlowEvent>) event.getClass();
	
	// 獲取當前事件的訂閱者串列
	// Get for Map, the algorithm is O(1).
	Set<Subscriber> subscribers = subMappings.get(slowEventType);
	if (null == subscribers) {
		LOGGER.debug("[NotifyCenter] No subscribers for slow event {}", slowEventType.getName());
		return;
	}
	
	// 回圈通知所有訂閱者
	// Notification single event subscriber
	for (Subscriber subscriber : subscribers) {
		// Whether to ignore expiration events
		if (subscriber.ignoreExpireEvent() && lastEventSequence > currentEventSequence) {
			LOGGER.debug("[NotifyCenter] the {} is unacceptable to this subscriber, because had expire", event.getClass());
			continue;
		}
		// 通知邏輯和父類是共用的
		// Notify single subscriber for slow event.
		notifySubscriber(subscriber, event);
	}
}

提示:
DefaultPublisher是一個發布器只負責發布一個事件,并通知訂閱了這個事件的所有訂閱者;DefaultSharePublisher則是一個發布器可以發布多個事件,并通知訂閱了這個事件的所有訂閱者,

通知中心(NotifyCenter)

NotifyCenter 在Nacos中主要用于注冊發布者、呼叫發布者發布事件、為發布者注冊訂閱者、為指定的事件增加指定的訂閱者等操作,可以說它完全接管了訂閱者、發布者和事件他們的組合程序,直接呼叫通知中心的相關方法即可實作事件發布訂閱者注冊等功能,

初始化資訊

package com.alibaba.nacos.common.notify;

public class NotifyCenter {
	
    /**
     * 單事件發布者內部的事件佇列初始容量
     */
    public static int ringBufferSize = 16384;

    /**
     * 多事件發布者內部的事件佇列初始容量
     */
    public static int shareBufferSize = 1024;

    /**
     * 發布者的狀態
     */
    private static final AtomicBoolean CLOSED = new AtomicBoolean(false);

    /**
     * 構造發布者的工廠
     */
    private static BiFunction<Class<? extends Event>, Integer, EventPublisher> publisherFactory = null;

    /**
     * 通知中心的實體
     */
    private static final NotifyCenter INSTANCE = new NotifyCenter();

    /**
     * 默認的多事件發布者
     */
    private DefaultSharePublisher sharePublisher;

    /**
     * 默認的單事件發布者型別
     * 此處并未直接指定單事件發布者是誰,只是限定了它的類別
     * 因為單事件發布者一個發布者只負責一個事件,因此會存在
     * 多個發布者實體,后面按需創建,并快取在publisherMap
     */
    private static Class<? extends EventPublisher> clazz = null;

    /**
     * Publisher management container.
     * 單事件發布者存盤容器
     */
    private final Map<String, EventPublisher> publisherMap = new ConcurrentHashMap<String, EventPublisher>(16);
	
	// 省略部分代碼
}

可以看到它初始化了一個通知中心的實體,這里是單例模式,定義了發布者,訂閱者是保存在發布者的內部,而發布者又保存在通知者的內部,這樣就組成了一套完整的事件發布機制,

靜態代碼塊

static {

	// 初始化DefaultPublisher的queue容量值
	// Internal ArrayBlockingQueue buffer size. For applications with high write throughput,
	// this value needs to be increased appropriately. default value is 16384
	String ringBufferSizeProperty = "nacos.core.notify.ring-buffer-size";
	ringBufferSize = Integer.getInteger(ringBufferSizeProperty, 16384);
	
	// 初始化DefaultSharePublisher的queue容量值
	// The size of the public publisher's message staging queue buffer
	String shareBufferSizeProperty = "nacos.core.notify.share-buffer-size";
	shareBufferSize = Integer.getInteger(shareBufferSizeProperty, 1024);
	
	// 使用Nacos SPI機制獲取事件發布者
	final Collection<EventPublisher> publishers = NacosServiceLoader.load(EventPublisher.class);
	
	// 獲取迭代器
	Iterator<EventPublisher> iterator = publishers.iterator();
	
	if (iterator.hasNext()) {
		clazz = iterator.next().getClass();
	} else {
		// 若為空,則使用默認的發布器(單事件發布者)
		clazz = DefaultPublisher.class;
	}
	
	// 宣告發布者工廠為一個函式,用于創建發布者實體
	publisherFactory = new BiFunction<Class<? extends Event>, Integer, EventPublisher>() {
		
		/**
		 * 為指定型別的事件創建一個單事件發布者物件
		 * @param cls       事件型別
		 * @param buffer    發布者內部佇列初始容量
		 * @return
		 */
		@Override
		public EventPublisher apply(Class<? extends Event> cls, Integer buffer) {
			try {
				// 實體化發布者
				EventPublisher publisher = clazz.newInstance();
				// 初始化
				publisher.init(cls, buffer);
				return publisher;
			} catch (Throwable ex) {
				LOGGER.error("Service class newInstance has error : {}", ex);
				throw new NacosRuntimeException(SERVER_ERROR, ex);
			}
		}
	};
	
	try {
		// 初始化多事件發布者
		// Create and init DefaultSharePublisher instance.
		INSTANCE.sharePublisher = new DefaultSharePublisher();
		INSTANCE.sharePublisher.init(SlowEvent.class, shareBufferSize);
		
	} catch (Throwable ex) {
		LOGGER.error("Service class newInstance has error : {}", ex);
	}

	// 增加關閉鉤子,用于關閉Publisher
	ThreadUtils.addShutdownHook(new Runnable() {
		@Override
		public void run() {
			shutdown();
		}
	});

    }

在靜態代碼塊中主要就做了兩件事:

初始化單事件發布者:可以由用戶擴展指定(通過Nacos SPI機制),也可以是Nacos默認的(DefaultPublisher),

初始化多事件發布者:DefaultSharePublisher,

注冊訂閱者

注冊訂閱者實際上就是將Subscriber添加到Publisher中,因為事件的發布是靠發布者來通知它內部的所有訂閱者,

/**
 * Register a Subscriber. If the Publisher concerned by the Subscriber does not exist, then PublihserMap will
 * preempt a placeholder Publisher first.
 *
 * @param consumer subscriber
 * @param <T>      event type
 */
public static <T> void registerSubscriber(final Subscriber consumer) {
	
	// 若想監聽多個事件,實作SmartSubscriber.subscribeTypes()方法,在里面回傳多個事件的串列即可
	// If you want to listen to multiple events, you do it separately,
	// based on subclass's subscribeTypes method return list, it can register to publisher.
	
	// 多事件訂閱者注冊
	if (consumer instanceof SmartSubscriber) {
		// 獲取事件串列
		for (Class<? extends Event> subscribeType : ((SmartSubscriber) consumer).subscribeTypes()) {
			// 判斷它的事件型別來決定采用哪種Publisher,多事件訂閱者由多事件發布者調度
			// For case, producer: defaultSharePublisher -> consumer: smartSubscriber.
			if (ClassUtils.isAssignableFrom(SlowEvent.class, subscribeType)) {
				//注冊到多事件發布者中
				INSTANCE.sharePublisher.addSubscriber(consumer, subscribeType);
			} else {
				// 注冊到單事件發布者中
				// For case, producer: defaultPublisher -> consumer: subscriber.
				addSubscriber(consumer, subscribeType);
			}
		}
		return;
	}
	
	// 單事件的訂閱者注冊
	final Class<? extends Event> subscribeType = consumer.subscribeType();
	// 防止誤使用,萬一有人在使用單事件訂閱者Subscriber的時候傳入了SlowEvent則可以在此避免
	if (ClassUtils.isAssignableFrom(SlowEvent.class, subscribeType)) {
		INSTANCE.sharePublisher.addSubscriber(consumer, subscribeType);
		// 添加完畢回傳
		return;
	}

	// 注冊到單事件發布者中
	addSubscriber(consumer, subscribeType);
}

/**
 * 單事件發布者添加訂閱者
 * Add a subscriber to publisher.
 * @param consumer      subscriber instance.
 * @param subscribeType subscribeType.
 */
private static void addSubscriber(final Subscriber consumer, Class<? extends Event> subscribeType) {
	// 獲取類的規范名稱,實際上就是包名加類名,作為topic
	final String topic = ClassUtils.getCanonicalName(subscribeType);
	synchronized (NotifyCenter.class) {
		// MapUtils.computeIfAbsent is a unsafe method.
		
		/**
		 * 生成指定型別的發布者,并將其放入publisherMap中
		 * 使用topic為key從publisherMap獲取資料,若為空則使用publisherFactory函式并傳遞subscribeType和ringBufferSize來實體
		 * 化一個clazz型別的發布者物件,使用topic為key放入publisherMap中,實際上就是為每一個型別的事件創建一個發布者,具體
		 * 可查看publisherFactory的邏輯,
		 */
		MapUtil.computeIfAbsent(INSTANCE.publisherMap, topic, publisherFactory, subscribeType, ringBufferSize);
	}
	// 獲取生成的發布者物件,將訂閱者添加進去
	EventPublisher publisher = INSTANCE.publisherMap.get(topic);
	publisher.addSubscriber(consumer);
}

提示:
單事件發布者容器內的存盤狀態為: 事件型別的完整限定名 -> DefaultPublisher.
例如:
com.alibaba.nacos.core.cluster.MembersChangeEvent -> {DefaultPublisher@6839} "Thread[nacos.publisher-com.alibaba.nacos.core.cluster.MembersChangeEvent,5,main]"

注冊發布者

實際上并沒有直接的注冊發布者這個概念,通過前面的章節你肯定知道發布者就兩種型別:單事件發布者、多事件發布者,單事件發布者直接就一個實體,多事件發布者會根據事件型別創建不同的實體,存盤于publisherMap中,它已經在通知中心了,因此并不需要有刻意的注冊動作,需要使用的時候
直接取即可,

注冊事件

注冊事件實際上就是將具體的事件和具體的發布者進行關聯,發布者有2種型別,那么事件也一定是兩種型別了(事件的型別這里說的是分類,服務于單事件發布者的事件和服務于多事件發布者的事件),

/**
 * Register publisher.
 *
 * @param eventType    class Instances type of the event type.
 * @param queueMaxSize the publisher's queue max size.
 */
public static EventPublisher registerToPublisher(final Class<? extends Event> eventType, final int queueMaxSize) {

	// 慢事件由多事件發布者處理
	if (ClassUtils.isAssignableFrom(SlowEvent.class, eventType)) {
		return INSTANCE.sharePublisher;
	}
	// 若不是慢事件,因為它可以存在多個不同的型別,因此需要判斷對應的發布者是否存在
	final String topic = ClassUtils.getCanonicalName(eventType);
	synchronized (NotifyCenter.class) {
		// 當前傳入的事件型別對應的發布者,有則忽略無則新建
		MapUtil.computeIfAbsent(INSTANCE.publisherMap, topic, publisherFactory, eventType, queueMaxSize);
	}
	return INSTANCE.publisherMap.get(topic);
}

這里并未有注冊動作,若是SlowEvent則直接回傳了,為何呢?這里再理一下關系,事件的實際用途是由訂閱者來決定的,由訂閱者來執行對應事件觸發后的操作,事件和發布者并沒有直接關系,而多事件發布者呢,它是一個發布者來處理所有的事件和訂閱者(事件:訂閱者,一對多的關系),這個事件都沒人訂閱何談發布呢?因此單純的注冊事件并沒有實際意義,反觀一次只能處理一個事件的單事件處理器(DefaultPublisher)則需要一個事件對應一個發布者,即便這個事件沒有人訂閱,也可以快取起來,

注銷訂閱者

注銷的操作基本上就是注冊的反向操作,

public static <T> void deregisterSubscriber(final Subscriber consumer) {
	// 若是多事件訂閱者
	if (consumer instanceof SmartSubscriber) {
		// 獲取事件串列
		for (Class<? extends Event> subscribeType : ((SmartSubscriber) consumer).subscribeTypes()) {
			// 若是慢事件
			if (ClassUtils.isAssignableFrom(SlowEvent.class, subscribeType)) {
				// 從多事件發布者中移除
				INSTANCE.sharePublisher.removeSubscriber(consumer, subscribeType);
			} else {
				// 從單事件發布者中移除
				removeSubscriber(consumer, subscribeType);
			}
		}
		return;
	}
	
	// 若是單事件訂閱者
	final Class<? extends Event> subscribeType = consumer.subscribeType();
	// 判斷是否是慢事件
	if (ClassUtils.isAssignableFrom(SlowEvent.class, subscribeType)) {
		INSTANCE.sharePublisher.removeSubscriber(consumer, subscribeType);
		return;
	}
	
	// 呼叫移除方法
	if (removeSubscriber(consumer, subscribeType)) {
		return;
	}
	throw new NoSuchElementException("The subscriber has no event publisher");
}

private static boolean removeSubscriber(final Subscriber consumer, Class<? extends Event> subscribeType) {
	// 獲取topic
	final String topic = ClassUtils.getCanonicalName(subscribeType);
	// 根據topic獲取對應的發布者
	EventPublisher eventPublisher = INSTANCE.publisherMap.get(topic);
	if (eventPublisher != null) {
		// 從發布者中移除訂閱者
		eventPublisher.removeSubscriber(consumer);
		return true;
	}
	return false;
}

注銷發布者

注銷發布者主要針對于單事件發布者來說的,因為多事件發布者只有一個實體,它需要處理多個事件型別,因此發布者不能移除,而單事件發布者一個發布者對應一個事件型別,因此某個型別的事件不需要處理的時候則需要將對應的發布者移除,

public static void deregisterPublisher(final Class<? extends Event> eventType) {
	// 獲取topic
	final String topic = ClassUtils.getCanonicalName(eventType);
	// 根據topic移除對應的發布者
	EventPublisher publisher = INSTANCE.publisherMap.remove(topic);
	try {
		// 呼叫關閉方法
		publisher.shutdown();
	} catch (Throwable ex) {
		LOGGER.error("There was an exception when publisher shutdown : {}", ex);
	}
}

public void shutdown() {
	// 標記關閉
	this.shutdown = true;
	// 清空快取
	this.queue.clear();
}

發布事件

發布事件的本質就是不同型別的發布者來呼叫內部維護的訂閱者的onEvent()方法,

private static boolean publishEvent(final Class<? extends Event> eventType, final Event event) {
	
	// 慢事件處理
	if (ClassUtils.isAssignableFrom(SlowEvent.class, eventType)) {
		return INSTANCE.sharePublisher.publish(event);
	}
	
	// 常規事件處理
	final String topic = ClassUtils.getCanonicalName(eventType);
	
	EventPublisher publisher = INSTANCE.publisherMap.get(topic);
	if (publisher != null) {
		return publisher.publish(event);
	}
	LOGGER.warn("There are no [{}] publishers for this event, please register", topic);
	return false;
}

總結

在Nacos中的事件發布分為兩條線:單一事件處理、多事件處理,圍繞這兩條線又有負責單一型別事件的訂閱者、發布者,也有負責多事件的訂閱者、發布者,區分開來兩種型別便很容易理解,

上圖展示了在通知中心中不同型別的事件、訂閱者、發布者的存盤狀態,

多事件發布者:

  • 發布者和事件的關系是一對多
  • 事件和訂閱者的關系是一對多
  • 發布者和訂閱者的關系是一對多
  • 事件型別為SlowEvent, 訂閱者型別是SmartSubscriber

單事件發布者

  • 發布者和事件的關系是一對一
  • 事件和訂閱者的關系是一對多
  • 發布者和訂閱者的關系是一對多
  • 事件型別為Event,訂閱者型別是Subscriber

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

標籤:其他

上一篇:給potplayer配置iptv源,看所有你想看的電視

下一篇: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