欢迎你来读这篇博客,这篇博客主要是关于Guava EventBus。
其中包括 EventBus 基本使用、事件发布与订阅、@Subscribe、同步事件分发、事件继承、DeadEvent、异常处理、AsyncEventBus、线程池分发、WatchService 文件监听实战、手写 EventBus、源码结构、优缺点,以及与 Spring Event、MQ 的区别。
序言
事件驱动是一种非常常见的软件设计思想。
在后端系统中,我们经常会遇到这样的场景:
某个动作发生后,希望多个模块都能感知并执行自己的逻辑。
比如:
1 2 3 4 5 6
| 订单创建成功 -> 发送站内信 -> 记录审计日志 -> 推送运营数据 -> 刷新缓存 -> 触发风控分析
|
如果订单服务直接调用所有模块,代码会变成:
1 2 3 4 5
| messageService.send(orderId); auditLogService.record(orderId); operationDataService.push(orderId); cacheService.refresh(orderId); riskService.analyze(orderId);
|
这样会导致订单服务依赖太多下游模块。
EventBus 的思想是:
发布者只发布事件,订阅者自己监听事件并处理,发布者和订阅者不需要直接互相依赖。
Guava 提供了一个进程内事件总线:
1
| com.google.common.eventbus.EventBus
|
它可以让我们用非常简单的方式做事件发布和订阅。
例如:
1
| eventBus.post(new OrderCreatedEvent(orderId));
|
订阅者:
1 2 3 4
| @Subscribe public void onOrderCreated(OrderCreatedEvent event) { }
|
但是,EventBus 也不是万能工具。
它曾经很流行,但现在新项目中不再是首选。
原因包括:
- 订阅关系隐式,代码跳转不直观;
- 基于反射扫描
@Subscribe;
- 事件流难追踪;
- 异常处理能力有限;
- 不适合跨进程;
- 不适合可靠消息;
- 不适合复杂业务流程;
- 官方文档也明确提醒新项目要谨慎使用。
所以这一章不是为了鼓励你在所有项目里使用 EventBus。
真正目标是:
通过 EventBus 学会观察者模式、事件分发、订阅者注册、反射调用、异步分发和事件总线源码设计。
这对于理解 Spring Event、领域事件、MQ、插件系统、回调机制都很有帮助。
本章对应课程结构:
1 2 3 4 5 6 7 8
| 第 7 章:事件总线 EventBus 7.1 EventBus 使用详解第一部分 7.2 EventBus 使用详解第二部分 7.3 EventBus 和 NIO 2.0 WatchService 综合实战 7.4 手动实现 EventBus:程序结构搭建 7.5 手动实现 EventBus:快速实现程序功能 7.6 手动实现 EventBus:总结与查缺补漏 7.7 EventBus 源码剖析以及优缺点总结
|
正文
chapter 1:EventBus 是什么
Guava EventBus 是一个进程内事件发布订阅工具。
它主要包含三个动作:
基本流程:
1 2 3 4
| Subscriber 注册到 EventBus Publisher 发布 Event EventBus 找到匹配的 Subscriber EventBus 调用 @Subscribe 方法
|
代码上大概是:
1 2 3 4 5
| EventBus eventBus = new EventBus();
eventBus.register(new OrderEventListener());
eventBus.post(new OrderCreatedEvent(1001L));
|
订阅者:
1 2 3 4 5 6 7
| public class OrderEventListener {
@Subscribe public void onOrderCreated(OrderCreatedEvent event) { System.out.println("收到订单创建事件:" + event.getOrderId()); } }
|
EventBus 是观察者模式的一种工程实现。
被观察者不再直接维护观察者列表,而是交给 EventBus 统一管理。
chapter 2:EventBus 的核心角色
EventBus 中常见角色如下:
| 角色 |
说明 |
| EventBus |
事件总线,负责注册、注销、发布、分发 |
| Event |
事件对象,例如 OrderCreatedEvent |
| Subscriber |
订阅者对象 |
@Subscribe |
标记订阅方法 |
| Dispatcher |
分发器,负责调用订阅者 |
| SubscriberRegistry |
订阅者注册表 |
| Executor |
执行器,同步或异步执行 |
| DeadEvent |
没有订阅者处理的事件包装对象 |
可以把 EventBus 理解成一个本地中转站:
1
| 事件发布者 -> EventBus -> 事件订阅者
|
发布者不需要知道订阅者是谁。
订阅者也不需要知道发布者是谁。
双方只通过事件类型产生关系。
chapter 3:添加 Guava 依赖
Maven 依赖示例:
1 2 3 4 5
| <dependency> <groupId>com.google.guava</groupId> <artifactId>guava</artifactId> <version>33.0.0-jre</version> </dependency>
|
真实项目中,版本号以项目实际依赖管理为准。
如果你的 Spring Boot 项目已经通过依赖管理引入 Guava,也可以不单独写版本。
chapter 4:定义第一个事件
定义一个订单创建事件:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19
| public class OrderCreatedEvent {
private final Long orderId;
private final Long userId;
public OrderCreatedEvent(Long orderId, Long userId) { this.orderId = orderId; this.userId = userId; }
public Long getOrderId() { return orderId; }
public Long getUserId() { return userId; } }
|
事件对象建议设计成不可变对象。
事件表达的是已经发生的事实。
所以命名建议是:
1 2 3 4
| OrderCreatedEvent PaymentSucceededEvent UserRegisteredEvent FileChangedEvent
|
不建议:
1 2 3
| CreateOrderEvent DoPaymentEvent HandleUserEvent
|
因为事件不是命令。
事件是事实。
chapter 5:定义订阅者
订阅者通过 @Subscribe 标记方法。
1 2 3 4 5 6 7 8 9 10 11 12
| import com.google.common.eventbus.Subscribe;
public class OrderEventListener {
@Subscribe public void onOrderCreated(OrderCreatedEvent event) { System.out.println("收到订单创建事件,orderId = " + event.getOrderId() + ",userId = " + event.getUserId()); } }
|
@Subscribe 方法有几个要求:
- 必须是实例方法;
- 通常是 public;
- 只能有一个参数;
- 参数类型就是订阅的事件类型;
- EventBus 会通过反射调用它。
一个订阅者类可以有多个 @Subscribe 方法。
1 2 3 4 5 6 7 8 9 10
| public class MultiEventListener {
@Subscribe public void onOrderCreated(OrderCreatedEvent event) { }
@Subscribe public void onPaymentSucceeded(PaymentSucceededEvent event) { } }
|
chapter 6:发布事件
基本用法:
1 2 3 4 5 6 7 8 9 10 11 12
| import com.google.common.eventbus.EventBus;
public class EventBusBasicDemo {
public static void main(String[] args) { EventBus eventBus = new EventBus();
eventBus.register(new OrderEventListener());
eventBus.post(new OrderCreatedEvent(1001L, 2001L)); } }
|
输出:
1
| 收到订单创建事件,orderId = 1001,userId = 2001
|
这就是 EventBus 最基本的使用方式。
chapter 7:同步事件分发
普通 EventBus 是同步分发。
也就是说:
会在当前线程中调用所有订阅者方法。
示例:
1 2 3 4 5 6 7 8 9
| import com.google.common.eventbus.Subscribe;
public class ThreadPrintListener {
@Subscribe public void onOrderCreated(OrderCreatedEvent event) { System.out.println("listener thread = " + Thread.currentThread().getName()); } }
|
发布:
1 2 3 4 5 6 7 8 9 10 11 12
| public class SyncEventBusDemo {
public static void main(String[] args) { EventBus eventBus = new EventBus();
eventBus.register(new ThreadPrintListener());
System.out.println("publisher thread = " + Thread.currentThread().getName());
eventBus.post(new OrderCreatedEvent(1001L, 2001L)); } }
|
输出类似:
1 2
| publisher thread = main listener thread = main
|
说明订阅者在发布事件的同一个线程中执行。
同步分发的优点是简单。
缺点是:
- 一个订阅者很慢,会拖慢发布者;
- 订阅者异常处理需要注意;
- 不适合耗时任务;
- 不适合强隔离场景。
chapter 8:多个订阅者
一个事件可以被多个订阅者处理。
1 2 3 4 5 6 7 8 9
| import com.google.common.eventbus.Subscribe;
public class MessageListener {
@Subscribe public void onOrderCreated(OrderCreatedEvent event) { System.out.println("发送站内信,orderId = " + event.getOrderId()); } }
|
1 2 3 4 5 6 7 8 9
| import com.google.common.eventbus.Subscribe;
public class AuditLogListener {
@Subscribe public void onOrderCreated(OrderCreatedEvent event) { System.out.println("记录审计日志,orderId = " + event.getOrderId()); } }
|
使用:
1 2 3 4 5 6
| EventBus eventBus = new EventBus();
eventBus.register(new MessageListener()); eventBus.register(new AuditLogListener());
eventBus.post(new OrderCreatedEvent(1001L, 2001L));
|
输出:
1 2
| 发送站内信,orderId = 1001 记录审计日志,orderId = 1001
|
EventBus 会把事件分发给所有匹配的订阅方法。
chapter 9:注销订阅者
可以使用:
1
| eventBus.unregister(listener);
|
示例:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
| public class UnregisterDemo {
public static void main(String[] args) { EventBus eventBus = new EventBus();
OrderEventListener listener = new OrderEventListener();
eventBus.register(listener);
eventBus.post(new OrderCreatedEvent(1001L, 2001L));
eventBus.unregister(listener);
eventBus.post(new OrderCreatedEvent(1002L, 2002L)); } }
|
第二次事件不会再被这个 listener 处理。
需要注意:
unregister 必须传入已经注册过的对象,否则可能抛异常。
在 Spring 项目中,如果用单例 Bean 注册到 EventBus,通常在启动时注册,关闭时注销。
chapter 10:事件继承
EventBus 支持事件类型继承。
假设有一个父类事件:
1 2 3 4 5 6 7 8 9 10 11 12
| public class BaseEvent {
private final String source;
public BaseEvent(String source) { this.source = source; }
public String getSource() { return source; } }
|
订单事件继承它:
1 2 3 4 5 6 7 8 9 10 11 12 13
| public class OrderPaidEvent extends BaseEvent {
private final Long orderId;
public OrderPaidEvent(String source, Long orderId) { super(source); this.orderId = orderId; }
public Long getOrderId() { return orderId; } }
|
监听父类事件:
1 2 3 4 5 6 7 8 9
| import com.google.common.eventbus.Subscribe;
public class BaseEventListener {
@Subscribe public void onBaseEvent(BaseEvent event) { System.out.println("收到 BaseEvent,source = " + event.getSource()); } }
|
监听子类事件:
1 2 3 4 5 6 7 8 9
| import com.google.common.eventbus.Subscribe;
public class OrderPaidListener {
@Subscribe public void onOrderPaid(OrderPaidEvent event) { System.out.println("收到 OrderPaidEvent,orderId = " + event.getOrderId()); } }
|
发布子类事件:
1
| eventBus.post(new OrderPaidEvent("order-service", 1001L));
|
父类订阅者和子类订阅者都可能收到事件。
这说明 EventBus 按事件类型层级查找订阅者。
chapter 11:事件继承的风险
事件继承虽然灵活,但也容易造成隐式订阅范围过大。
例如:
1 2 3
| @Subscribe public void onObject(Object event) { }
|
这个方法几乎能接收所有事件。
这会带来几个问题:
- 所有事件都进入这个方法;
- DeadEvent 可能不会生成;
- 调试时很难判断谁处理了事件;
- 订阅关系变得更隐蔽。
所以不建议随便订阅过于宽泛的父类型。
推荐订阅明确事件类型:
1 2 3
| OrderCreatedEvent PaymentSucceededEvent FileChangedEvent
|
chapter 12:DeadEvent 是什么
如果发布的事件没有任何订阅者处理,EventBus 会把它包装成 DeadEvent 再发布一次。
示例:
1 2 3 4 5 6 7 8 9 10 11
| import com.google.common.eventbus.DeadEvent; import com.google.common.eventbus.Subscribe;
public class DeadEventListener {
@Subscribe public void onDeadEvent(DeadEvent deadEvent) { System.out.println("发现无人处理的事件:" + deadEvent.getEvent().getClass().getName()); } }
|
使用:
1 2 3 4 5
| EventBus eventBus = new EventBus();
eventBus.register(new DeadEventListener());
eventBus.post(new OrderCreatedEvent(1001L, 2001L));
|
如果没有 OrderCreatedEvent 的订阅者,DeadEventListener 会收到 DeadEvent。
输出:
1
| 发现无人处理的事件:OrderCreatedEvent
|
chapter 13:DeadEvent 的用途
DeadEvent 常用于:
- 调试事件没有被消费的问题;
- 发现漏注册订阅者;
- 监控无人处理事件;
- 记录异常事件流。
但它不是可靠消息机制。
如果你发布一个事件,没有订阅者,DeadEvent 只是本地提示。
它不会帮你重试,也不会帮你持久化。
chapter 14:异常处理
订阅者方法如果抛出异常,EventBus 会捕获异常并交给异常处理器。
默认情况下,它通常会记录日志,而不是把异常直接抛回发布者。
示例:
1 2 3 4 5 6 7 8 9
| import com.google.common.eventbus.Subscribe;
public class BrokenListener {
@Subscribe public void onOrderCreated(OrderCreatedEvent event) { throw new RuntimeException("listener failed"); } }
|
发布:
1 2
| eventBus.register(new BrokenListener()); eventBus.post(new OrderCreatedEvent(1001L, 2001L));
|
这里需要注意:
EventBus 的异常处理不适合承载复杂业务错误处理。
如果某个监听器失败后必须重试、补偿、告警、阻断主流程,EventBus 默认机制并不够。
chapter 15:自定义 SubscriberExceptionHandler
可以在创建 EventBus 时传入异常处理器。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25
| import com.google.common.eventbus.EventBus; import com.google.common.eventbus.SubscriberExceptionContext; import com.google.common.eventbus.SubscriberExceptionHandler;
public class CustomExceptionHandlerDemo {
public static void main(String[] args) { SubscriberExceptionHandler handler = new SubscriberExceptionHandler() { @Override public void handleException(Throwable exception, SubscriberExceptionContext context) { System.out.println("订阅者异常:" + "event = " + context.getEvent() + ",subscriber = " + context.getSubscriber() + ",method = " + context.getSubscriberMethod().getName() + ",error = " + exception.getMessage()); } };
EventBus eventBus = new EventBus(handler);
eventBus.register(new BrokenListener());
eventBus.post(new OrderCreatedEvent(1001L, 2001L)); } }
|
自定义异常处理器适合:
- 记录日志;
- 统计失败次数;
- 上报告警;
- 输出订阅者信息;
- 辅助开发排查。
但它依旧不是完整的错误恢复机制。
chapter 16:AsyncEventBus 是什么
AsyncEventBus 是 Guava 提供的异步事件总线。
它需要一个 Executor。
示例:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22
| import com.google.common.eventbus.AsyncEventBus; import com.google.common.eventbus.EventBus;
import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors;
public class AsyncEventBusDemo {
public static void main(String[] args) { ExecutorService executorService = Executors.newFixedThreadPool(4);
EventBus eventBus = new AsyncEventBus(executorService);
eventBus.register(new ThreadPrintListener());
System.out.println("publisher thread = " + Thread.currentThread().getName());
eventBus.post(new OrderCreatedEvent(1001L, 2001L));
executorService.shutdown(); } }
|
输出可能是:
1 2
| publisher thread = main listener thread = pool-1-thread-1
|
说明订阅者在独立线程池中执行。
chapter 17:AsyncEventBus 使用注意事项
异步事件分发会带来新的问题。
1. 线程池要合理配置
不要直接使用无界线程池。
推荐根据业务设置:
- 核心线程数;
- 最大线程数;
- 队列大小;
- 拒绝策略;
- 线程命名;
- 异常处理。
2. 事件顺序不一定可靠
多个事件异步执行时,不一定按发布顺序完成。
3. 异常不会抛回发布线程
订阅者异常需要单独处理。
4. 需要考虑上下文传递
例如:
- traceId;
- requestId;
- MDC;
- 用户上下文。
异步线程默认拿不到调用线程的 ThreadLocal。
5. 业务一致性更复杂
异步不是魔法,只是把复杂度搬到了另一个线程。
如果监听器失败要重试,EventBus 不是最好的工具。
MQ 更合适。
chapter 18:EventBus 和 WatchService 综合实战
接下来做一个综合实战:
使用 JDK NIO 2.0 WatchService 监听目录文件变化,再用 EventBus 分发文件变更事件。
业务场景:
- 监听配置目录;
- 文件新增时触发导入;
- 文件修改时刷新缓存;
- 文件删除时记录日志;
- 事件分发给多个监听器。
技术点:
1 2
| WatchService 负责监听文件系统变化 EventBus 负责把文件变化事件分发给订阅者
|
chapter 19:定义文件变更事件
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42
| import java.nio.file.Path; import java.nio.file.WatchEvent;
public class FileChangedEvent {
private final Path directory;
private final Path fileName;
private final WatchEvent.Kind<?> kind;
public FileChangedEvent(Path directory, Path fileName, WatchEvent.Kind<?> kind) { this.directory = directory; this.fileName = fileName; this.kind = kind; }
public Path getDirectory() { return directory; }
public Path getFileName() { return fileName; }
public WatchEvent.Kind<?> getKind() { return kind; }
public Path getFullPath() { return directory.resolve(fileName); }
@Override public String toString() { return "FileChangedEvent{" + "directory=" + directory + ", fileName=" + fileName + ", kind=" + kind + '}'; } }
|
chapter 20:定义文件监听器订阅者
1 2 3 4 5 6 7 8 9 10 11 12 13 14
| import com.google.common.eventbus.Subscribe;
import java.nio.file.StandardWatchEventKinds;
public class FileChangeLogListener {
@Subscribe public void onFileChanged(FileChangedEvent event) { System.out.println("文件变化:" + event.getKind() + ",path = " + event.getFullPath()); } }
|
导入监听器:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20
| import com.google.common.eventbus.Subscribe;
import java.nio.file.StandardWatchEventKinds;
public class CsvImportListener {
@Subscribe public void onFileChanged(FileChangedEvent event) { if (!StandardWatchEventKinds.ENTRY_CREATE.equals(event.getKind())) { return; }
if (!event.getFileName().toString().endsWith(".csv")) { return; }
System.out.println("发现新的 CSV 文件,准备导入:" + event.getFullPath()); } }
|
配置刷新监听器:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20
| import com.google.common.eventbus.Subscribe;
import java.nio.file.StandardWatchEventKinds;
public class ConfigRefreshListener {
@Subscribe public void onFileChanged(FileChangedEvent event) { if (!StandardWatchEventKinds.ENTRY_MODIFY.equals(event.getKind())) { return; }
if (!event.getFileName().toString().endsWith(".properties")) { return; }
System.out.println("配置文件变化,刷新配置:" + event.getFullPath()); } }
|
chapter 21:实现 WatchService 文件监听
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63
| import com.google.common.eventbus.EventBus;
import java.io.IOException; import java.nio.file.*;
public class DirectoryWatchService implements Runnable {
private final Path directory;
private final EventBus eventBus;
private volatile boolean running = true;
public DirectoryWatchService(Path directory, EventBus eventBus) { this.directory = directory; this.eventBus = eventBus; }
@Override public void run() { try (WatchService watchService = FileSystems.getDefault().newWatchService()) { directory.register( watchService, StandardWatchEventKinds.ENTRY_CREATE, StandardWatchEventKinds.ENTRY_MODIFY, StandardWatchEventKinds.ENTRY_DELETE );
System.out.println("开始监听目录:" + directory);
while (running) { WatchKey key = watchService.take();
for (WatchEvent<?> event : key.pollEvents()) { WatchEvent.Kind<?> kind = event.kind();
if (StandardWatchEventKinds.OVERFLOW.equals(kind)) { continue; }
Path fileName = (Path) event.context();
eventBus.post(new FileChangedEvent(directory, fileName, kind)); }
boolean valid = key.reset();
if (!valid) { break; } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("文件监听线程被中断"); } catch (IOException e) { throw new RuntimeException("watch directory failed: " + directory, e); } }
public void stop() { running = false; } }
|
注意:
WatchService 需要关闭;
watchService.take() 会阻塞;
- 线程中断要恢复中断标记;
OVERFLOW 表示事件可能丢失;
- 真实项目要考虑递归监听子目录。
chapter 22:文件监听通知系统客户端
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28
| import com.google.common.eventbus.EventBus;
import java.nio.file.Files; import java.nio.file.Path;
public class FileWatchEventBusDemo {
public static void main(String[] args) throws Exception { Path directory = Path.of("watch-demo");
if (!Files.exists(directory)) { Files.createDirectories(directory); }
EventBus eventBus = new EventBus();
eventBus.register(new FileChangeLogListener()); eventBus.register(new CsvImportListener()); eventBus.register(new ConfigRefreshListener());
DirectoryWatchService watchService = new DirectoryWatchService(directory, eventBus);
Thread watchThread = new Thread(watchService, "directory-watch-thread"); watchThread.start();
System.out.println("请在 watch-demo 目录下创建、修改或删除文件测试。"); } }
|
这个例子说明:
1 2 3
| WatchService 负责产生事件 EventBus 负责分发事件 监听器负责处理事件
|
这就是一个简单的文件监听通知系统。
chapter 23:WatchService 实战注意事项
1. WatchService 不等于可靠文件队列
它适合监听变化,不适合做强可靠任务队列。
2. 文件刚创建时可能还没写完
监听到 ENTRY_CREATE 后,文件可能仍在写入。
导入系统要考虑延迟或文件完成标记。
3. 事件可能重复
修改文件可能触发多次 ENTRY_MODIFY。
要做防抖或去重。
4. 子目录需要单独注册
默认只监听当前目录。
5. 生产环境要处理异常和重启
监听线程挂了要能恢复。
6. 重要文件处理建议配合任务表
如果文件导入很重要,不要只靠 WatchService。
可以扫描目录 + 任务表 + 状态机 + 重试。
chapter 24:手动实现 EventBus:程序结构搭建
接下来我们手写一个简化版 EventBus。
目标不是完全复制 Guava,而是理解核心结构。
我们设计几个核心对象:
1 2 3 4 5
| MiniEventBus MiniSubscribe SubscriberRegistry Subscriber Dispatcher
|
职责:
| 对象 |
职责 |
| MiniEventBus |
对外提供 register、unregister、post |
| MiniSubscribe |
标记订阅方法 |
| SubscriberRegistry |
扫描和保存订阅者 |
| Subscriber |
封装订阅对象和订阅方法 |
| Dispatcher |
负责事件分发 |
| SubscriberExceptionHandler |
处理订阅异常 |
程序结构:
1 2 3 4 5 6 7
| eventbus ├── MiniEventBus ├── MiniSubscribe ├── Subscriber ├── SubscriberRegistry ├── Dispatcher └── SubscriberExceptionHandler
|
chapter 25:定义 MiniSubscribe 注解
1 2 3 4 5 6 7 8 9
| import java.lang.annotation.ElementType; import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target;
@Target(ElementType.METHOD) @Retention(RetentionPolicy.RUNTIME) public @interface MiniSubscribe { }
|
这个注解用于标记订阅方法。
和 Guava 的:
类似。
chapter 26:定义 Subscriber
Subscriber 封装订阅对象和方法。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41
| import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method;
public class Subscriber {
private final Object target;
private final Method method;
public Subscriber(Object target, Method method) { this.target = target; this.method = method; this.method.setAccessible(true); }
public void invoke(Object event) throws Exception { try { method.invoke(target, event); } catch (InvocationTargetException e) { Throwable targetException = e.getTargetException();
if (targetException instanceof Exception exception) { throw exception; }
if (targetException instanceof Error error) { throw error; }
throw new RuntimeException(targetException); } }
public Object getTarget() { return target; }
public Method getMethod() { return method; } }
|
这里使用反射调用订阅方法。
注意:
1
| InvocationTargetException
|
是反射调用时的包装异常。
真正的业务异常在:
中。
chapter 27:定义异常处理器
1 2 3 4
| public interface SubscriberExceptionHandler {
void handleException(Throwable exception, Object event, Subscriber subscriber); }
|
默认实现:
1 2 3 4 5 6 7 8 9 10 11 12 13 14
| public class LoggingSubscriberExceptionHandler implements SubscriberExceptionHandler {
@Override public void handleException(Throwable exception, Object event, Subscriber subscriber) { System.out.println("订阅者执行异常,event = " + event.getClass().getName() + ",subscriber = " + subscriber.getTarget().getClass().getName() + ",method = " + subscriber.getMethod().getName() + ",error = " + exception.getMessage()); } }
|
chapter 28:定义 SubscriberRegistry
订阅者注册表负责扫描 @MiniSubscribe 方法,并按事件类型保存。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47
| import java.lang.reflect.Method; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap;
public class SubscriberRegistry {
private final Map<Class<?>, List<Subscriber>> subscriberMap = new ConcurrentHashMap<>();
public void register(Object listener) { Method[] methods = listener.getClass().getDeclaredMethods();
for (Method method : methods) { if (!method.isAnnotationPresent(MiniSubscribe.class)) { continue; }
Class<?>[] parameterTypes = method.getParameterTypes();
if (parameterTypes.length != 1) { throw new IllegalArgumentException("@MiniSubscribe method must have exactly one parameter: " + method); }
Class<?> eventType = parameterTypes[0];
subscriberMap.computeIfAbsent(eventType, key -> new ArrayList<>()) .add(new Subscriber(listener, method)); } }
public List<Subscriber> getSubscribers(Object event) { List<Subscriber> result = new ArrayList<>();
Class<?> eventClass = event.getClass();
for (Map.Entry<Class<?>, List<Subscriber>> entry : subscriberMap.entrySet()) { Class<?> subscribedType = entry.getKey();
if (subscribedType.isAssignableFrom(eventClass)) { result.addAll(entry.getValue()); } }
return result; } }
|
这里有一个关键判断:
1
| subscribedType.isAssignableFrom(eventClass)
|
它表示:
如果订阅类型是事件真实类型的父类或接口,也可以接收事件。
这就是事件继承的基础。
chapter 29:定义 Dispatcher
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20
| import java.util.List;
public class Dispatcher {
private final SubscriberExceptionHandler exceptionHandler;
public Dispatcher(SubscriberExceptionHandler exceptionHandler) { this.exceptionHandler = exceptionHandler; }
public void dispatch(Object event, List<Subscriber> subscribers) { for (Subscriber subscriber : subscribers) { try { subscriber.invoke(event); } catch (Throwable e) { exceptionHandler.handleException(e, event, subscriber); } } } }
|
Dispatcher 负责分发事件。
这里是同步分发。
后面可以扩展异步分发。
chapter 30:定义 MiniEventBus
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28
| import java.util.List;
public class MiniEventBus {
private final SubscriberRegistry registry;
private final Dispatcher dispatcher;
public MiniEventBus() { this.registry = new SubscriberRegistry(); this.dispatcher = new Dispatcher(new LoggingSubscriberExceptionHandler()); }
public void register(Object listener) { registry.register(listener); }
public void post(Object event) { List<Subscriber> subscribers = registry.getSubscribers(event);
if (subscribers.isEmpty()) { System.out.println("没有订阅者处理事件:" + event.getClass().getName()); return; }
dispatcher.dispatch(event, subscribers); } }
|
这就是一个最小可用 EventBus。
chapter 31:使用 MiniEventBus
定义事件:
1 2 3 4 5 6 7 8 9 10 11 12
| public class UserRegisteredEvent {
private final Long userId;
public UserRegisteredEvent(Long userId) { this.userId = userId; }
public Long getUserId() { return userId; } }
|
定义监听器:
1 2 3 4 5 6 7
| public class UserEventListener {
@MiniSubscribe public void onUserRegistered(UserRegisteredEvent event) { System.out.println("欢迎新用户,userId = " + event.getUserId()); } }
|
客户端:
1 2 3 4 5 6 7 8 9 10
| public class MiniEventBusDemo {
public static void main(String[] args) { MiniEventBus eventBus = new MiniEventBus();
eventBus.register(new UserEventListener());
eventBus.post(new UserRegisteredEvent(1001L)); } }
|
输出:
chapter 32:手写 EventBus 功能补全:unregister
当前版本没有注销功能。
可以在 SubscriberRegistry 中增加:
1 2 3 4 5
| public void unregister(Object listener) { for (List<Subscriber> subscribers : subscriberMap.values()) { subscribers.removeIf(subscriber -> subscriber.getTarget() == listener); } }
|
然后在 MiniEventBus 增加:
1 2 3
| public void unregister(Object listener) { registry.unregister(listener); }
|
注意:
- 这里用对象引用判断;
- 真实实现要考虑并发安全;
- 还要考虑同一个对象重复注册问题。
chapter 33:线程安全问题
上面的简化版有线程安全问题。
例如:
1
| Map<Class<?>, List<Subscriber>>
|
虽然 Map 是 ConcurrentHashMap,但 List 是 ArrayList。
多线程 register 和 post 时可能出问题。
可以改成:
例如:
1 2
| subscriberMap.computeIfAbsent(eventType, key -> new CopyOnWriteArrayList<>()) .add(new Subscriber(listener, method));
|
或者注册阶段完成后不再修改。
真实 EventBus 需要认真处理:
- 注册并发;
- 分发并发;
- 注销并发;
- 订阅者集合快照;
- 重复注册;
- reentrant event;
- 执行顺序。
chapter 34:异步分发扩展
可以定义异步 Dispatcher。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28
| import java.util.List; import java.util.concurrent.Executor;
public class AsyncDispatcher extends Dispatcher {
private final Executor executor;
private final SubscriberExceptionHandler exceptionHandler;
public AsyncDispatcher(Executor executor, SubscriberExceptionHandler exceptionHandler) { super(exceptionHandler); this.executor = executor; this.exceptionHandler = exceptionHandler; }
@Override public void dispatch(Object event, List<Subscriber> subscribers) { for (Subscriber subscriber : subscribers) { executor.execute(() -> { try { subscriber.invoke(event); } catch (Throwable e) { exceptionHandler.handleException(e, event, subscriber); } }); } } }
|
然后 EventBus 可以根据构造参数选择同步或异步 Dispatcher。
注意:
异步分发一定要考虑线程池和异常处理。
chapter 35:手写 EventBus 设计缺陷总结
我们的 MiniEventBus 很简单,但它有不少缺陷。
1. 反射扫描不够完善
没有处理继承方法、桥接方法、泛型等复杂情况。
2. 线程安全不足
注册、注销、分发并发时可能有问题。
3. 异常处理简单
只是打印日志,没有重试、告警、隔离。
4. 没有 DeadEvent
无人处理事件只是打印提示。
5. 没有订阅者执行顺序控制
多个订阅者顺序不明确。
6. 没有上下文传递
异步模式下 traceId、MDC 无法自动传递。
7. 没有背压
事件太多时可能压垮线程池。
8. 没有事件持久化
进程重启事件丢失。
这说明:
手写 EventBus 可以学习原理,但不要轻易把简化版用于生产关键链路。
chapter 36:EventBus 源码结构概览
Guava EventBus 的源码大致可以拆成几部分:
1 2 3 4 5 6
| EventBus -> SubscriberRegistry -> Dispatcher -> Subscriber -> SubscriberExceptionHandler -> Executor
|
1. EventBus
对外入口。
提供:
1 2 3
| register() unregister() post()
|
2. SubscriberRegistry
负责扫描订阅者对象中的 @Subscribe 方法,并维护事件类型到订阅者的映射。
3. Subscriber
封装订阅者对象和订阅方法。
负责反射调用。
4. Dispatcher
负责具体事件分发策略。
5. Executor
决定订阅方法在哪个线程执行。
普通 EventBus 使用直接执行。
AsyncEventBus 使用传入的 Executor。
6. SubscriberExceptionHandler
处理订阅者异常。
chapter 37:SubscriberRegistry 的作用
SubscriberRegistry 的核心职责是:
建立事件类型到订阅方法的映射。
大概结构:
1
| Map<Class<?>, Collection<Subscriber>>
|
当调用:
1
| eventBus.register(listener)
|
它会扫描 listener 中所有带 @Subscribe 的方法。
例如:
1 2 3
| @Subscribe public void onOrderCreated(OrderCreatedEvent event) { }
|
注册表会记录:
1
| OrderCreatedEvent -> listener.onOrderCreated
|
发布事件时:
1
| eventBus.post(new OrderCreatedEvent(...))
|
EventBus 会从注册表中找到所有匹配的 Subscriber。
chapter 38:Dispatcher 的作用
Dispatcher 负责把事件分发给订阅者。
它关心的是:
- 是否立即分发;
- 是否排队分发;
- 是否异步分发;
- 如何避免递归分发问题;
- 如何调用 Subscriber。
普通同步分发可以理解为:
1 2 3
| for (Subscriber subscriber : subscribers) { subscriber.dispatchEvent(event); }
|
异步分发可以理解为:
1
| executor.execute(() -> subscriber.dispatchEvent(event));
|
Dispatcher 是 EventBus 内部非常关键的一层。
它把“找订阅者”和“执行订阅者”分开了。
chapter 39:Executor 的作用
普通 EventBus 通常在当前线程执行订阅者。
AsyncEventBus 通过 Executor 执行订阅者。
1
| EventBus eventBus = new AsyncEventBus(executor);
|
这意味着:
Executor 的选择非常重要。
不要随便用:
1
| Executors.newCachedThreadPool()
|
生产环境建议使用有界线程池。
chapter 40:EventBus 的优点
1. 使用简单
只需要:
1 2 3
| @Subscribe eventBus.register(listener) eventBus.post(event)
|
2. 解耦发布者和订阅者
发布者不用直接依赖订阅者。
3. 支持多个订阅者
一个事件可以触发多个监听器。
4. 支持事件继承
订阅父类事件可以接收子类事件。
5. 支持同步和异步
普通 EventBus 同步。
AsyncEventBus 异步。
6. 适合进程内轻量事件通知
例如:
- 本地缓存刷新;
- 插件通知;
- 开发工具;
- 非核心事件分发;
- 测试 demo。
chapter 41:EventBus 的缺点
1. 订阅关系隐式
代码里搜索调用链不直观。
你看到:
很难立刻知道谁会处理它。
2. 调试困难
事件流隐藏在反射和注册表里。
排查问题要找所有 @Subscribe。
3. 反射调用
EventBus 基于反射调用订阅方法。
这会带来一定运行时成本,也对代码优化器、混淆器不友好。
4. 异常处理能力有限
默认不适合复杂业务错误处理。
5. 不支持跨进程
EventBus 是本地内存事件总线。
服务重启事件就没了。
6. 没有持久化和重试
不像 MQ 有消息存储、重试、死信队列。
7. 容易滥用
所有业务都靠 EventBus 串起来后,业务流程会变得很隐蔽。
chapter 42:为什么官方不再强烈推荐 EventBus
Guava 官方文档现在已经明确提醒:
主要原因可以总结为:
1. 现代库有更好的替代方案
例如:
- Spring Event;
- Reactive Streams;
- RxJava;
- Reactor;
- Kotlin Flow;
- MQ;
- 依赖注入框架;
- 明确的接口回调。
2. 订阅关系难追踪
生产者和消费者之间没有显式引用。
这会让代码阅读和调试变难。
3. 反射对优化工具不友好
例如 Proguard、R8、代码混淆和裁剪场景。
4. 异常处理不是强项
复杂业务错误处理、重试、补偿不适合依赖 EventBus。
5. 需要启动时注册订阅者
这可能导致应用急切初始化所有订阅者。
6. 不适合现代大型后端系统关键链路
大型系统更需要可观测性、可靠性、事务边界、消息持久化和明确依赖关系。
所以:
EventBus 可以学,可以维护老代码,但新项目要谨慎选择。
chapter 43:EventBus 与 Spring Event 对比
| 对比项 |
Guava EventBus |
Spring Event |
| 所属生态 |
Guava |
Spring |
| 注册方式 |
手动 register |
Spring Bean 自动管理 |
| 订阅方式 |
@Subscribe |
@EventListener |
| 异步支持 |
AsyncEventBus |
@Async |
| 事务事件 |
不支持 |
@TransactionalEventListener |
| Spring Bean 注入 |
不直接管理 |
天然支持 |
| 推荐场景 |
轻量本地事件、老项目 |
Spring 项目内部事件 |
如果是 Spring Boot 项目,通常优先考虑 Spring Event。
尤其是你需要事务提交后发布事件:
1
| @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
|
这点 EventBus 不具备。
chapter 44:EventBus 与 MQ 对比
| 对比项 |
EventBus |
MQ |
| 通信范围 |
单 JVM 内 |
跨进程、跨服务 |
| 是否持久化 |
否 |
通常支持 |
| 失败重试 |
弱 |
支持 |
| 消息堆积 |
不适合 |
支持 |
| 可靠性 |
低 |
高 |
| 复杂度 |
低 |
较高 |
| 适合场景 |
本地轻量事件 |
分布式可靠事件 |
如果事件只是进程内通知,EventBus 可以。
如果事件涉及:
- 跨服务;
- 可靠投递;
- 重试;
- 削峰;
- 审计;
- 任务异步;
- 业务最终一致性;
应该考虑 MQ。
例如:
1 2 3 4
| Kafka RabbitMQ RocketMQ Pulsar
|
chapter 45:EventBus 使用建议
1. 新项目不要默认选择 EventBus
尤其是 Spring Boot 项目,优先考虑 Spring Event 或明确接口调用。
2. 事件命名要清晰
推荐:
1 2 3
| OrderCreatedEvent PaymentSucceededEvent FileChangedEvent
|
3. 不要订阅过宽类型
避免:
1 2
| @Subscribe public void onEvent(Object event)
|
4. 订阅方法不要做太重逻辑
耗时任务建议异步或 MQ。
5. 异步必须配置合理线程池
不要无界线程池。
6. 异常处理要明确
不要指望默认日志就够。
7. 关键业务不要依赖本地 EventBus
比如扣库存、支付、退款。
8. 需要可靠事件就用 MQ
EventBus 不是可靠消息系统。
9. 老项目要做好事件地图
维护 EventBus 项目时,建议整理:
否则代码很难读。
10. 学源码,不滥用
EventBus 很适合学习观察者模式和事件分发源码,但不一定适合新项目核心链路。
chapter 46:完整 MiniEventBus 代码汇总
MiniSubscribe
1 2 3 4 5 6 7 8 9
| import java.lang.annotation.ElementType; import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target;
@Target(ElementType.METHOD) @Retention(RetentionPolicy.RUNTIME) public @interface MiniSubscribe { }
|
Subscriber
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41
| import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method;
public class Subscriber {
private final Object target;
private final Method method;
public Subscriber(Object target, Method method) { this.target = target; this.method = method; this.method.setAccessible(true); }
public void invoke(Object event) throws Exception { try { method.invoke(target, event); } catch (InvocationTargetException e) { Throwable targetException = e.getTargetException();
if (targetException instanceof Exception exception) { throw exception; }
if (targetException instanceof Error error) { throw error; }
throw new RuntimeException(targetException); } }
public Object getTarget() { return target; }
public Method getMethod() { return method; } }
|
SubscriberExceptionHandler
1 2 3 4
| public interface SubscriberExceptionHandler {
void handleException(Throwable exception, Object event, Subscriber subscriber); }
|
LoggingSubscriberExceptionHandler
1 2 3 4 5 6 7 8 9 10 11 12 13 14
| public class LoggingSubscriberExceptionHandler implements SubscriberExceptionHandler {
@Override public void handleException(Throwable exception, Object event, Subscriber subscriber) { System.out.println("订阅者执行异常,event = " + event.getClass().getName() + ",subscriber = " + subscriber.getTarget().getClass().getName() + ",method = " + subscriber.getMethod().getName() + ",error = " + exception.getMessage()); } }
|
SubscriberRegistry
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54
| import java.lang.reflect.Method; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList;
public class SubscriberRegistry {
private final Map<Class<?>, CopyOnWriteArrayList<Subscriber>> subscriberMap = new ConcurrentHashMap<>();
public void register(Object listener) { Method[] methods = listener.getClass().getDeclaredMethods();
for (Method method : methods) { if (!method.isAnnotationPresent(MiniSubscribe.class)) { continue; }
Class<?>[] parameterTypes = method.getParameterTypes();
if (parameterTypes.length != 1) { throw new IllegalArgumentException("@MiniSubscribe method must have exactly one parameter: " + method); }
Class<?> eventType = parameterTypes[0];
subscriberMap.computeIfAbsent(eventType, key -> new CopyOnWriteArrayList<>()) .add(new Subscriber(listener, method)); } }
public void unregister(Object listener) { for (List<Subscriber> subscribers : subscriberMap.values()) { subscribers.removeIf(subscriber -> subscriber.getTarget() == listener); } }
public List<Subscriber> getSubscribers(Object event) { List<Subscriber> result = new ArrayList<>();
Class<?> eventClass = event.getClass();
for (Map.Entry<Class<?>, CopyOnWriteArrayList<Subscriber>> entry : subscriberMap.entrySet()) { Class<?> subscribedType = entry.getKey();
if (subscribedType.isAssignableFrom(eventClass)) { result.addAll(entry.getValue()); } }
return result; } }
|
Dispatcher
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20
| import java.util.List;
public class Dispatcher {
private final SubscriberExceptionHandler exceptionHandler;
public Dispatcher(SubscriberExceptionHandler exceptionHandler) { this.exceptionHandler = exceptionHandler; }
public void dispatch(Object event, List<Subscriber> subscribers) { for (Subscriber subscriber : subscribers) { try { subscriber.invoke(event); } catch (Throwable e) { exceptionHandler.handleException(e, event, subscriber); } } } }
|
MiniEventBus
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32
| import java.util.List;
public class MiniEventBus {
private final SubscriberRegistry registry;
private final Dispatcher dispatcher;
public MiniEventBus() { this.registry = new SubscriberRegistry(); this.dispatcher = new Dispatcher(new LoggingSubscriberExceptionHandler()); }
public void register(Object listener) { registry.register(listener); }
public void unregister(Object listener) { registry.unregister(listener); }
public void post(Object event) { List<Subscriber> subscribers = registry.getSubscribers(event);
if (subscribers.isEmpty()) { System.out.println("没有订阅者处理事件:" + event.getClass().getName()); return; }
dispatcher.dispatch(event, subscribers); } }
|
使用示例
1 2 3 4 5 6 7 8 9 10 11 12
| public class UserRegisteredEvent {
private final Long userId;
public UserRegisteredEvent(Long userId) { this.userId = userId; }
public Long getUserId() { return userId; } }
|
1 2 3 4 5 6 7
| public class UserEventListener {
@MiniSubscribe public void onUserRegistered(UserRegisteredEvent event) { System.out.println("欢迎新用户,userId = " + event.getUserId()); } }
|
1 2 3 4 5 6 7 8 9 10
| public class MiniEventBusDemo {
public static void main(String[] args) { MiniEventBus eventBus = new MiniEventBus();
eventBus.register(new UserEventListener());
eventBus.post(new UserRegisteredEvent(1001L)); } }
|
chapter 47:一句话总结
这一章的核心是:
EventBus 是一个进程内发布订阅工具,它通过注册订阅者、扫描 @Subscribe 方法、按事件类型匹配订阅者,再通过 Dispatcher 调用订阅方法,实现了本地事件分发。
EventBus 的优点是:
- 简单;
- 轻量;
- 解耦;
- 支持多订阅者;
- 支持同步和异步;
- 适合学习观察者模式和事件分发源码。
EventBus 的缺点是:
- 订阅关系隐式;
- 调试困难;
- 基于反射;
- 异常处理能力有限;
- 不支持跨进程;
- 不支持持久化;
- 不适合可靠消息;
- 新项目中已经不再是首选。
如果你在维护老项目,EventBus 必须能看懂。
如果你在做新项目,尤其是 Spring Boot 项目,优先考虑:
1 2 3 4 5
| Spring Event 明确接口调用 领域事件 MQ Reactor / Reactive Streams
|
如果只是学习源码设计,EventBus 非常值得拆。
它能让你理解:
1 2 3 4 5 6 7 8
| 注解扫描 订阅者注册表 事件类型匹配 反射调用 同步分发 异步分发 异常处理 DeadEvent
|
一句话:
EventBus 可以作为学习事件驱动设计的好教材,但不要把它当成现代后端系统的万能事件架构。
参考资料
- Google Guava API: EventBus.
- Google Guava API: AsyncEventBus.
- Google Guava API: Subscribe.
- Google Guava API: DeadEvent.
- Google Guava Wiki: EventBusExplained.
- Java Documentation: WatchService.
- Java Documentation: WatchKey.
- Java Documentation: StandardWatchEventKinds.
- Spring Framework Documentation: Application Events.
- Spring Framework Documentation: Transaction-bound Events.
- Reactive Streams Specification.
- Oracle Java Reflection Documentation.
启示录
富贵岂由人,时会高志须酬。
能成功于千载者,必以近察远。