Appearance
观察者模式 (Observer)
一、定义
一句话概括:定义对象之间的一对多依赖关系,当一个对象状态改变时,所有依赖它的对象都会收到通知并自动更新。
官方定义:Define a one-to-many dependency between objects so that when one object changes state, all its dependents are notified and updated automatically.
二、解决的问题
2.1 问题场景
在系统中,当一个对象的状态发生变化时,需要通知其他相关对象做出响应。例如:
- 天气预报:气象站数据变化时,多个显示面板需要更新
- 股票行情:股票价格变化时,多个投资人的持仓需要更新
- 消息通知:用户发布文章时,所有粉丝收到通知
- 事件驱动:按钮点击时,多个监听器响应
2.2 不用观察者模式会怎样?
java
// 反例:轮询检查状态变化
class WeatherData {
private float temperature;
private float humidity;
public void setMeasurements(float temp, float humidity) {
this.temperature = temp;
this.humidity = humidity;
// 硬编码通知所有显示面板
display1.update(temp, humidity);
display2.update(temp, humidity);
display3.update(temp, humidity);
// 新增显示面板需要修改此代码
}
}
// 或者:客户端轮询
while (true) {
if (weatherData.hasChanged()) {
display.update(weatherData);
}
Thread.sleep(1000);
}问题:
- 主题与观察者紧密耦合,新增观察者需修改主题代码
- 轮询方式浪费 CPU 资源
- 无法动态添加/移除观察者
- 通知逻辑分散,难以维护
三、结构
3.1 角色组成
| 角色 | 说明 |
|---|---|
| Subject(主题/被观察者) | 维护观察者列表,提供注册/移除/通知方法 |
| Observer(观察者接口) | 定义更新接口,当 Subject 状态变化时被调用 |
| ConcreteSubject(具体主题) | 实现 Subject,状态变化时通知所有观察者 |
| ConcreteObserver(具体观察者) | 实现 Observer 接口,定义收到通知后的具体行为 |
3.2 类图(ASCII)
┌──────────────┐ ┌──────────────┐
│ Subject │────────►│ Observer │
├──────────────┤ ├──────────────┤
│ - observers │ │ + update() │
│ + attach() │ └──────┬───────┘
│ + detach() │ │
│ + notify() │ ┌──────┴──────┐
└──────┬───────┘ │ │
│ ▼ ▼
▼ ┌──────────┐ ┌──────────┐
┌──────────────┐ │Concrete │ │Concrete │
│ConcreteSubject│ │ObserverA │ │ObserverB │
│ - state │ └──────────┘ └──────────┘
│ + getState() │
│ + setState() │
└──────────────┘3.3 时序图(ASCII)
Subject Observer1 Observer2
│ │ │
│ attach(o1) │ │
│◄─────────────────│ │
│ attach(o2) │ │
│◄───────────────────────────────│
│ │ │
│ setState() │ │
│ notify() │ │
│───update()──────►│ │
│───update()─────────────────────►│
│ │ │四、代码实现
4.1 基础实现
java
// ==================== 观察者接口 ====================
interface Observer {
void update(float temperature, float humidity, float pressure);
}
// ==================== 主题接口 ====================
interface Subject {
void registerObserver(Observer observer);
void removeObserver(Observer observer);
void notifyObservers();
}
// ==================== 具体主题:气象站 ====================
class WeatherData implements Subject {
private List<Observer> observers = new ArrayList<>();
private float temperature;
private float humidity;
private float pressure;
@Override
public void registerObserver(Observer observer) {
observers.add(observer);
}
@Override
public void removeObserver(Observer observer) {
observers.remove(observer);
}
@Override
public void notifyObservers() {
for (Observer observer : observers) {
observer.update(temperature, humidity, pressure);
}
}
public void setMeasurements(float temperature, float humidity, float pressure) {
this.temperature = temperature;
this.humidity = humidity;
this.pressure = pressure;
measurementsChanged();
}
private void measurementsChanged() {
notifyObservers();
}
public float getTemperature() { return temperature; }
public float getHumidity() { return humidity; }
public float getPressure() { return pressure; }
}
// ==================== 具体观察者:当前天气显示 ====================
class CurrentConditionsDisplay implements Observer {
private float temperature;
private float humidity;
private Subject weatherData;
public CurrentConditionsDisplay(Subject weatherData) {
this.weatherData = weatherData;
weatherData.registerObserver(this);
}
@Override
public void update(float temperature, float humidity, float pressure) {
this.temperature = temperature;
this.humidity = humidity;
display();
}
public void display() {
System.out.printf("当前天气: 温度 %.1f°C, 湿度 %.1f%%\n", temperature, humidity);
}
}
// ==================== 具体观察者:统计显示 ====================
class StatisticsDisplay implements Observer {
private float maxTemp = Float.MIN_VALUE;
private float minTemp = Float.MAX_VALUE;
private float tempSum = 0;
private int numReadings = 0;
@Override
public void update(float temperature, float humidity, float pressure) {
tempSum += temperature;
numReadings++;
if (temperature > maxTemp) maxTemp = temperature;
if (temperature < minTemp) minTemp = temperature;
display();
}
public void display() {
System.out.printf("统计: 平均 %.1f°C, 最高 %.1f°C, 最低 %.1f°C\n",
tempSum / numReadings, maxTemp, minTemp);
}
}
// ==================== 具体观察者:预报显示 ====================
class ForecastDisplay implements Observer {
private float currentPressure = 29.92f;
private float lastPressure;
@Override
public void update(float temperature, float humidity, float pressure) {
lastPressure = currentPressure;
currentPressure = pressure;
display();
}
public void display() {
if (currentPressure > lastPressure) {
System.out.println("预报: 天气转好");
} else if (currentPressure < lastPressure) {
System.out.println("预报: 天气转凉,可能下雨");
} else {
System.out.println("预报: 天气不变");
}
}
}
// ==================== 客户端 ====================
public class ObserverDemo {
public static void main(String[] args) {
WeatherData weatherData = new WeatherData();
CurrentConditionsDisplay currentDisplay = new CurrentConditionsDisplay(weatherData);
StatisticsDisplay statisticsDisplay = new StatisticsDisplay();
weatherData.registerObserver(statisticsDisplay);
ForecastDisplay forecastDisplay = new ForecastDisplay();
weatherData.registerObserver(forecastDisplay);
weatherData.setMeasurements(28, 65, 30.4f);
weatherData.setMeasurements(30, 70, 29.2f);
weatherData.setMeasurements(26, 90, 29.2f);
}
}4.2 进阶实现
4.2.1 推模型 vs 拉模型
推模型(Push):Subject 将变化的数据直接推送给 Observer(如上面的示例)。
java
// 推模型:Subject 主动推送数据
interface PushObserver {
void update(float temperature, float humidity, float pressure);
}
// Subject 调用: observer.update(temperature, humidity, pressure);拉模型(Pull):Subject 只通知 Observer 发生了变化,Observer 需要时再从 Subject 中拉取数据。
java
// 拉模型:Observer 根据需要拉取数据
interface PullObserver {
void update(Subject subject); // 只传递 Subject 引用
}
class PullCurrentConditionsDisplay implements PullObserver {
private WeatherData weatherData;
@Override
public void update(Subject subject) {
if (subject instanceof WeatherData) {
WeatherData wd = (WeatherData) subject;
// 按需拉取
System.out.printf("温度: %.1f, 湿度: %.1f\n",
wd.getTemperature(), wd.getHumidity());
}
}
}对比:
| 维度 | 推模型 | 拉模型 |
|---|---|---|
| 数据传输 | Subject 主动推送 | Observer 按需拉取 |
| 灵活性 | 低,Observer 必须接收所有数据 | 高,Observer 选择性获取 |
| 耦合度 | 较高,Observer 需知道数据结构 | 较低,Observer 依赖 Subject 接口 |
| 适用场景 | 数据量小,Observer 需要全部数据 | 数据量大,Observer 只需部分数据 |
4.2.2 Java Observable 类(JDK 内置)
java
// JDK 观察者模式(Java 9 已废弃,但仍值得了解)
import java.util.Observable;
import java.util.Observer;
class WeatherDataObservable extends Observable {
private float temperature;
public void setTemperature(float temperature) {
this.temperature = temperature;
setChanged(); // 标记状态已改变
notifyObservers(temperature); // 通知所有观察者
}
public float getTemperature() { return temperature; }
}
class WeatherObserver implements Observer {
@Override
public void update(Observable o, Object arg) {
if (o instanceof WeatherDataObservable) {
float temp = (float) arg;
System.out.println("温度更新: " + temp);
}
}
}
// 使用
public class ObservableDemo {
public static void main(String[] args) {
WeatherDataObservable observable = new WeatherDataObservable();
observable.addObserver(new WeatherObserver());
observable.setTemperature(30.5f);
}
}注意:Observable 是一个类而非接口,限制了其复用性,且 Java 9 已标记为 @Deprecated。推荐使用 PropertyChangeListener 或自定义实现。
4.2.3 事件驱动实现
java
// 自定义事件对象
class OrderEvent extends EventObject {
private final Long orderId;
private final String action;
public OrderEvent(Object source, Long orderId, String action) {
super(source);
this.orderId = orderId;
this.action = action;
}
public Long getOrderId() { return orderId; }
public String getAction() { return action; }
}
// 事件监听器接口
interface OrderEventListener extends EventListener {
void onOrderEvent(OrderEvent event);
}
// 事件源
class OrderEventSource {
private final List<OrderEventListener> listeners = new CopyOnWriteArrayList<>();
public void addListener(OrderEventListener listener) {
listeners.add(listener);
}
public void removeListener(OrderEventListener listener) {
listeners.remove(listener);
}
public void fireEvent(OrderEvent event) {
for (OrderEventListener listener : listeners) {
listener.onOrderEvent(event);
}
}
}
// 具体监听器:发送短信
class SmsOrderListener implements OrderEventListener {
@Override
public void onOrderEvent(OrderEvent event) {
System.out.println("发送短信通知: 订单 " + event.getOrderId()
+ " " + event.getAction());
}
}
// 具体监听器:记录日志
class LogOrderListener implements OrderEventListener {
@Override
public void onOrderEvent(OrderEvent event) {
System.out.println("记录日志: 订单 " + event.getOrderId()
+ " " + event.getAction());
}
}
// 使用
public class EventDrivenDemo {
public static void main(String[] args) {
OrderEventSource source = new OrderEventSource();
source.addListener(new SmsOrderListener());
source.addListener(new LogOrderListener());
source.fireEvent(new OrderEvent(source, 1001L, "已创建"));
source.fireEvent(new OrderEvent(source, 1001L, "已支付"));
}
}4.3 生产级实现
Spring Boot 事件驱动架构
java
// ==================== 自定义事件 ====================
class OrderCreatedEvent extends ApplicationEvent {
private final Long orderId;
private final Long userId;
private final BigDecimal amount;
public OrderCreatedEvent(Object source, Long orderId, Long userId, BigDecimal amount) {
super(source);
this.orderId = orderId;
this.userId = userId;
this.amount = amount;
}
public Long getOrderId() { return orderId; }
public Long getUserId() { return userId; }
public BigDecimal getAmount() { return amount; }
}
// ==================== 事件监听器 ====================
@Component
class SmsNotificationListener {
@EventListener
@Async // 异步执行
public void handleOrderCreated(OrderCreatedEvent event) {
System.out.println(Thread.currentThread().getName()
+ " [短信服务] 发送订单确认短信: 订单" + event.getOrderId()
+ ", 金额" + event.getAmount());
}
}
@Component
class InventoryUpdateListener {
@EventListener
@Order(1) // 优先执行
public void handleOrderCreated(OrderCreatedEvent event) {
System.out.println("[库存服务] 扣减库存: 订单" + event.getOrderId());
}
}
@Component
class PointsRewardListener {
@EventListener
@Async
public void handleOrderCreated(OrderCreatedEvent event) {
System.out.println("[积分服务] 发放积分: 用户" + event.getUserId()
+ ", 订单" + event.getOrderId());
}
}
@Component
class AuditLogListener {
@EventListener
public void handleOrderCreated(OrderCreatedEvent event) {
System.out.println("[审计日志] 记录订单创建: " + event.getOrderId());
}
}
// ==================== 事件发布器 ====================
@Service
class OrderService {
@Autowired
private ApplicationEventPublisher eventPublisher;
public Order createOrder(Long userId, BigDecimal amount) {
// 创建订单逻辑
Order order = new Order();
order.setId(1001L);
order.setUserId(userId);
order.setAmount(amount);
// 发布事件
eventPublisher.publishEvent(
new OrderCreatedEvent(this, order.getId(), userId, amount));
return order;
}
}
// ==================== 异步配置 ====================
@Configuration
@EnableAsync
class AsyncConfig {
@Bean
public Executor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(10);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("event-");
executor.initialize();
return executor;
}
}
// ==================== Controller ====================
@RestController
@RequestMapping("/api/orders")
class OrderController {
@Autowired private OrderService orderService;
@PostMapping
public String createOrder(@RequestParam Long userId,
@RequestParam BigDecimal amount) {
Order order = orderService.createOrder(userId, amount);
return "订单创建成功: " + order.getId();
}
}基于 Guava EventBus 的实现
java
// 使用 Google Guava EventBus(轻量级事件总线)
import com.google.common.eventbus.EventBus;
import com.google.common.eventbus.Subscribe;
class EventBusDemo {
public static void main(String[] args) {
EventBus eventBus = new EventBus();
// 注册监听器
eventBus.register(new Object() {
@Subscribe
public void handleOrderCreated(OrderCreatedEvent event) {
System.out.println("处理订单: " + event.getOrderId());
}
});
// 发布事件
eventBus.post(new OrderCreatedEvent("source", 1001L, 1L, new BigDecimal("99.99")));
}
}五、优缺点
优点
- 松耦合:Subject 只知道 Observer 接口,不知道具体实现
- 支持广播通信:一个 Subject 可以通知多个 Observer
- 动态订阅:可以在运行时动态添加/移除观察者
- 符合开闭原则:新增观察者无需修改 Subject 代码
缺点
- 通知顺序不可控:Observer 之间的通知顺序不确定(除非使用
@Order) - 可能导致循环依赖:Observer 和 Subject 之间相互调用可能导致死循环
- 内存泄漏风险:忘记移除 Observer 可能导致内存泄漏
- 性能问题:如果 Observer 很多,通知所有 Observer 可能耗时
- 异常处理:一个 Observer 抛出异常可能影响其他 Observer
六、适用场景
- 事件驱动架构:Spring 事件、消息队列、EventBus
- GUI 组件:按钮点击、文本框变化等事件监听
- 消息推送:用户关注、订阅通知
- 数据同步:缓存更新、配置变更通知
- 监控告警:系统指标变化时触发告警
- 实时数据展示:股票行情、天气预报等实时更新的面板
- 微服务事件通知:订单状态变更、库存变化等
七、JDK / Spring 框架中的实际应用
| 框架 | 应用位置 | 说明 |
|---|---|---|
| JDK | java.util.Observer / Observable | JDK 内置观察者模式(已废弃) |
| JDK | java.beans.PropertyChangeListener | 属性变更监听器 |
| JDK | java.util.EventListener | 事件监听器接口 |
| JDK | javax.servlet.ServletContextListener | Servlet 上下文监听器 |
| Spring | ApplicationEvent / ApplicationListener | Spring 事件机制 |
| Spring | @EventListener | 注解驱动的事件监听 |
| Spring Boot | ApplicationRunner / CommandLineRunner | 启动事件监听 |
| Guava | EventBus | 轻量级事件总线 |
| RxJava | Observable / Observer | 响应式编程中的观察者模式 |
八、与其他模式的关系
与发布-订阅模式(Pub-Sub)
- 相似:都是一对多通知
- 区别:观察者模式中 Subject 和 Observer 直接通信(松耦合);发布-订阅模式中通过消息代理(Broker/EventBus)解耦,Publisher 和 Subscriber 完全不知道对方的存在
与中介者模式
- 观察者模式可以结合中介者模式:Mediator 作为 Subject,Colleague 作为 Observer,实现多对多通信
与责任链模式
- 观察者模式是广播通知(所有 Observer 都收到),责任链模式是接力传递(找到第一个能处理的)
与 MVC 模式
- MVC 中的 Model 是 Subject,View 是 Observer,Model 变化时通知 View 更新
九、面试常见问题
Q1:观察者模式和发布-订阅模式有什么区别?
A:观察者模式中,Subject 和 Observer 直接交互(Subject 维护 Observer 列表),两者是松耦合的。发布-订阅模式中,Publisher 和 Subscriber 不直接通信,而是通过消息代理(EventBus、消息队列)进行通信,两者完全解耦。观察者模式通常用于单进程内的事件通知,发布-订阅模式通常用于跨进程/跨系统的消息通信。
Q2:Spring 的事件机制是如何实现的?
A:Spring 事件机制基于观察者模式。ApplicationEventPublisher 发布事件,ApplicationListener 或 @EventListener 监听事件。ApplicationContext 在初始化时会注册所有监听器,发布事件时遍历监听器列表并调用。@Async 注解可以让监听器异步执行。@Order 注解可以控制监听器执行顺序。
Q3:观察者模式可能导致哪些问题?如何解决?
A:(1) 内存泄漏:Observer 被注册后未移除,导致 GC 无法回收。解决:使用弱引用(WeakReference)或确保在适当时机移除。(2) 循环通知:A 通知 B,B 又通知 A,导致死循环。解决:在通知前检查是否已处理过。(3) 异常传播:一个 Observer 抛异常导致后续 Observer 无法执行。解决:用 try-catch 包裹每个 Observer 的调用。
Q4:Java 9 为什么要废弃 Observable 和 Observer?
A:(1) Observable 是类而非接口,限制了灵活性(Java 是单继承);(2) 不支持序列化;(3) 线程不安全;(4) setChanged() 是 protected 方法,不能被组合使用;(5) 没有提供有序通知机制。推荐使用 PropertyChangeListener、Flow API(Java 9+)或第三方库如 Guava EventBus。
Q5:如何处理观察者模式中的异步通知?
A:(1) 使用线程池异步执行通知:executor.submit(() -> observer.update());(2) 使用 Spring 的 @Async 注解;(3) 使用消息队列(RabbitMQ、Kafka)解耦;(4) 使用 Reactor/RxJava 响应式编程。异步通知可以提高吞吐量,但需要注意事务一致性和异常处理。
