diff --git a/event/src/main/java/com/zfoo/event/manager/EventBus.java b/event/src/main/java/com/zfoo/event/manager/EventBus.java index 6f2753b8..fd11a3ee 100644 --- a/event/src/main/java/com/zfoo/event/manager/EventBus.java +++ b/event/src/main/java/com/zfoo/event/manager/EventBus.java @@ -36,14 +36,15 @@ public abstract class EventBus { private static final Logger logger = LoggerFactory.getLogger(EventBus.class); - // 线程池的大小 + /** + * 线程池的大小. event的线程池比较大 + */ public static final int EXECUTORS_SIZE = Runtime.getRuntime().availableProcessors() * 2; private static final ExecutorService[] executors = new ExecutorService[EXECUTORS_SIZE]; private static final Map, List> receiverMap = new HashMap<>(); - static { for (int i = 0; i < executors.length; i++) { var namedThreadFactory = new EventThreadFactory(i + 1); @@ -85,16 +86,28 @@ public abstract class EventBus { } /** - * 随机获取一个线程池 + * 随机获取一个线程 */ public static Executor asyncExecute() { return executors[RandomUtils.randomInt(EXECUTORS_SIZE)]; } + /** + * 用指定线程执行 + * + * @param hashcode + * @return + */ public static Executor execute(int hashcode) { return executors[Math.abs(hashcode % EXECUTORS_SIZE)]; } + /** + * 执行方法调用 + * + * @param event 事件 + * @param receiverList 所有的观察者 + */ private static void doSubmit(IEvent event, List receiverList) { for (var receiver : receiverList) { try { @@ -107,6 +120,12 @@ public abstract class EventBus { } } + /** + * 注册事件及其对应观察者 + * + * @param eventType + * @param receiver + */ public static void registerEventReceiver(Class eventType, IEventReceiver receiver) { receiverMap.computeIfAbsent(eventType, it -> new LinkedList<>()).add(receiver); } diff --git a/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java b/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java index 003b1510..24163eb8 100644 --- a/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java +++ b/event/src/main/java/com/zfoo/event/schema/EventRegisterProcessor.java @@ -81,6 +81,8 @@ public class EventRegisterProcessor implements BeanPostProcessor { var receiverDefinition = new EventReceiverDefinition(bean, method, eventClazz); var enhanceReceiverDefinition = EnhanceUtils.createEventReceiver(receiverDefinition); + + // key:class类型 value:观察者 注册Event的receiverMap中 EventBus.registerEventReceiver(eventClazz, enhanceReceiverDefinition); } } catch (Throwable t) { diff --git a/net/src/main/java/com/zfoo/net/consumer/IConsumer.java b/net/src/main/java/com/zfoo/net/consumer/IConsumer.java index eb2c5cb7..c1f33e5c 100644 --- a/net/src/main/java/com/zfoo/net/consumer/IConsumer.java +++ b/net/src/main/java/com/zfoo/net/consumer/IConsumer.java @@ -32,6 +32,8 @@ public interface IConsumer { /** * 直接发送,不需要任何返回值 + *

+ * 例子:参考 com.zfoo.app.zapp.chat.controller。FrinedController 的 atApplyFriendRequest方法,客户端发起申请请求,chat服务处理后,再把消息直接发给网关 * * @param packet 需要发送的包 * @param argument 计算负载均衡的参数,比如用户的id