EventBus原理深度解析

一、问题描述

在工作中,经常会遇见使用异步的方式来发送事件,或者触发另外一个动作:经常用到的框架是MQ(分布式方式通知)。如果是同一个jvm里面通知的话,就可以使用EventBus。由于EventBus使用起来简单、便捷,因此,工作中会经常用到。深入理解该框架的原理就很有必要。

二、框架解析

2.1、组织结构

eventbus的组织结构如下:

eventbus主要有以下几部分组成:

1、eventbus、asyncEventBus:事件发送器。

2、event:事件承载单元。

3、SubscriberRegistry:订阅者注册器,将订阅者注册到event上,即将有注解Subscribe的方法和event绑定起来。

4、Dispatcher:事件分发器,将事件的订阅者调用来执行。

5、Subscriber、SynchronizedSubscriber:订阅者,并发订阅还是同步订阅。

2.2、运行原理

1、eventbus是基于注册监听的方式来运行的,因此,首先需要将eventbus,然后才会有事件及监听者。新建eventbus或者AsyncEventBus的方式如下:

 EventBus eventBus = new EventBus();

或者

1 BlockingQueue<Runnable> workQueue = new LinkedBlockingQueue<>(20); 2 ThreadPoolExecutor executor = new ThreadPoolExecutor(5, 20, 3 30, TimeUnit.SECONDS, workQueue); 4 AsyncEventBus asyncEventBus = new AsyncEventBus(executor);

2、注册监听者。

eventBus.register(eventListener);

底层就是将类eventListener中所有注解有Subscribe的方法与其Event对放在一个map中(一个event可以对应多个Subscribe的方法)。实现如下:

1 void register(Object listener) { 2 Multimap<Class<?>, Subscriber> listenerMethods = findAllSubscribers(listener); 3 4 for (Entry<Class<?>, Collection<Subscriber>> entry : listenerMethods.asMap().entrySet()) { 5 Class<?> eventType = entry.getKey(); 6 Collection<Subscriber> eventMethodsInListener = entry.getValue(); 7 8 CopyOnWriteArraySet<Subscriber> eventSubscribers = subscribers.get(eventType); 9 10 if (eventSubscribers == null) { 11 CopyOnWriteArraySet<Subscriber> newSet = new CopyOnWriteArraySet<>(); 12 eventSubscribers = 13 MoreObjects.firstNonNull(subscribers.putIfAbsent(eventType, newSet), newSet); 14 } 15 16 eventSubscribers.addAll(eventMethodsInListener); 17 } 18 }

3、事件发送:执行指定事件类型的订阅者(包含了method),从订阅者中获取指定事件的订阅者,然后按照规则(同步、异步)执行指定的方法。

1 public void post(Object event) { 2 Iterator<Subscriber> eventSubscribers = subscribers.getSubscribers(event); 3 if (eventSubscribers.hasNext()) { 4 dispatcher.dispatch(event, eventSubscribers); 5 } else if (!(event instanceof DeadEvent)) { 6 // the event had no subscribers and was not itself a DeadEvent 7 post(new DeadEvent(this, event)); 8 } 9 }

上述代码说明,如果事件没有监听者,就当作死亡事件来对待。

1 /** Dispatches {@code event} to this subscriber using the proper executor. */ 2 final void dispatchEvent(final Object event) { 3 executor.execute( 4 new Runnable() { 5 @Override 6 public void run() { 7 try { 8 invokeSubscriberMethod(event); 9 } catch (InvocationTargetException e) { 10 bus.handleSubscriberException(e.getCause(), context(event)); 11 } 12 } 13 }); 14 } 15 void invokeSubscriberMethod(Object event) throws InvocationTargetException { 16 try { 17 method.invoke(target, checkNotNull(event)); 18 } catch (IllegalArgumentException e) { 19 throw new Error("Method rejected target/argument: " + event, e); 20 } catch (IllegalAccessException e) { 21 throw new Error("Method became inaccessible: " + event, e); 22 } catch (InvocationTargetException e) { 23 if (e.getCause() instanceof Error) { 24 throw (Error) e.getCause(); 25 } 26 throw e; 27 } 28 }

这里就说明,最后就是被订阅的方法被调用。

4、EventBus与AsyncEventBus的区别

从字面上看,AsyncEventBus是异步的EventBus,那么EventBus应该就是同步的了。EventBus的executor为MoreExecutors.directExecutor(),其实现如下:

1 public static Executor directExecutor() { 2 return DirectExecutor.INSTANCE; 3 } 4 5 /** See {@link #directExecutor} for behavioral notes. */ 6 private enum DirectExecutor implements Executor { 7 INSTANCE; 8 9 @Override 10 public void execute(Runnable command) { 11 command.run(); 12 } 13 14 @Override 15 public String toString() { 16 return "MoreExecutors.directExecutor()"; 17 } 18 }

其execute方法直接执行线程的run方法,即同步调用run方法执行。EventBus的dispatcher为PerThreadQueuedDispatcher。其dispatch方法如下:

1 @Override 2 void dispatch(Object event, Iterator<Subscriber> subscribers) { 3 checkNotNull(event); 4 checkNotNull(subscribers); 5 Queue<Event> queueForThread = queue.get(); 6 queueForThread.offer(new Event(event, subscribers)); 7 8 if (!dispatching.get()) { 9 dispatching.set(true); 10 try { 11 Event nextEvent; 12 while ((nextEvent = queueForThread.poll()) != null) { 13 while (nextEvent.subscribers.hasNext()) { 14 nextEvent.subscribers.next().dispatchEvent(nextEvent.event); 15 } 16 } 17 } finally { 18 dispatching.remove(); 19 queue.remove(); 20 } 21 } 22 }

dispatchEvent的实现如下:

1 final void dispatchEvent(final Object event) { 2 executor.execute( 3 new Runnable() { 4 @Override 5 public void run() { 6 try { 7 invokeSubscriberMethod(event); 8 } catch (InvocationTargetException e) { 9 bus.handleSubscriberException(e.getCause(), context(event)); 10 } 11 } 12 }); 13 }

因此,整个执行过程如下:

整个过程都是同步方式执行,因此,EventBus是同步的。

AsyncEventBus的dispatcher为LegacyAsyncDispatcher,executor为自己指定的线程池。运行流程如下:

虚线为线程池异步调度,因此,AsyncEventBus为异步方式。

5、AllowConcurrentEvents的作用

它所在的代码为:

1 static Subscriber create(EventBus bus, Object listener, Method method) { 2 return isDeclaredThreadSafe(method) 3 ? new Subscriber(bus, listener, method) 4 : new SynchronizedSubscriber(bus, listener, method); 5 } 6 7 private static boolean isDeclaredThreadSafe(Method method) { 8 return method.getAnnotation(AllowConcurrentEvents.class) != null; 9 }

即如果订阅者方法上有注解AllowConcurrentEvents,则返回Subscriber,否则,返回SynchronizedSubscriber。SynchronizedSubscriber的字面意思为同步订阅者,它的实现代码为:

1 @Override 2 void invokeSubscriberMethod(Object event) throws InvocationTargetException { 3 synchronized (this) { 4 super.invokeSubscriberMethod(event); 5 } 6 }

即没有使用注解AllowConcurrentEvents的订阅者,在并发环境中,都是串行执行。这在高并发环境中,会严重影响性能。

三、使用案例

3.1、eventbus定义

1@Configuration 2public class ConfigBean { 3 4 @Bean 5 public EventBus executorService() { 6 BlockingQueue<Runnable> workQueue = new LinkedBlockingQueue<>(20); 7 ThreadPoolExecutor executor = new ThreadPoolExecutor(5, 20, 8 30, TimeUnit.SECONDS, workQueue); 9 return new AsyncEventBus(executor); 10 } 11}

3.2、注册与事件发送

1@Service 2public class TestService implements InitializingBean { 3 4 @Autowired 5 private EventListener eventListener ; 6 7 @Autowired 8 private EventBus eventBus ; 9 10 public void postEvent(){ 11 eventBus.post(new LoginEvent("iwill","123456")); 12 } 13 14 @Override 15 public void afterPropertiesSet() throws Exception { 16 eventBus.register(eventListener); 17 } 18}

3.3、订阅者定义

1package com.iwill.eventBus.listener; 2 3import com.google.common.eventbus.Subscribe; 4import com.iwill.eventBus.event.LoginEvent; 5import com.iwill.eventBus.event.RegisterEvent; 6import org.springframework.stereotype.Component; 7 8@Component 9public class EventListener { 10 11 @Subscribe 12 public void subscribeLoginEvent1(LoginEvent event){ 13 System.out.println("method 1 : receive login event "); 14 } 15 16 @Subscribe 17 public void subscribeLoginEvent2(LoginEvent event){ 18 System.out.println("method 2 : receive login event "); 19 } 20 21 @Subscribe 22 public void subscribeRegisterEvent(RegisterEvent event){ 23 try{ 24 Thread.sleep(10000L); 25 }catch (Exception exp){ 26 exp.printStackTrace(); 27 } 28 System.out.println("method : receive register event "); 29 } 30}

四、注意事项

1、在高并发的环境下使用AsyncEventBus时,发送事件可能会出现异常,因为它使用的线程池,当线程池的线程不够用时,会拒绝接收任务,就会执行线程池的拒绝策略,如果需要关注是否提交事件成功,就需要将线程池的拒绝策略设为抛出异常,并且try-catch来捕获异常。如下:

1 try { 2 eventBus.post(new LoginEvent("iwill", "123456")); 3 }catch (Exception exp){ 4 //TODO 落表或者其他处理 5 }

2、本文用到的guava版本如下:

1 <dependency> 2 <groupId>com.google.guava</groupId> 3 <artifactId>guava</artifactId> 4 <version>26.0-jre</version> 5 </dependency>
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )

SpringBoot整合Redis乱码原因及解决方案

问题描述:springboot使用springdataredis存储数据时乱码rediskey/value出现\\xAC\\xED\\x00\\x05t\\x00\\x05问题分析:查看RedisTemplate类!(https://oscimg.oschina.net/oscnet/0a85565fa