Skip to content

观察者模式 (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);
}

问题:

  1. 主题与观察者紧密耦合,新增观察者需修改主题代码
  2. 轮询方式浪费 CPU 资源
  3. 无法动态添加/移除观察者
  4. 通知逻辑分散,难以维护

三、结构

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")));
    }
}

五、优缺点

优点

  1. 松耦合:Subject 只知道 Observer 接口,不知道具体实现
  2. 支持广播通信:一个 Subject 可以通知多个 Observer
  3. 动态订阅:可以在运行时动态添加/移除观察者
  4. 符合开闭原则:新增观察者无需修改 Subject 代码

缺点

  1. 通知顺序不可控:Observer 之间的通知顺序不确定(除非使用 @Order
  2. 可能导致循环依赖:Observer 和 Subject 之间相互调用可能导致死循环
  3. 内存泄漏风险:忘记移除 Observer 可能导致内存泄漏
  4. 性能问题:如果 Observer 很多,通知所有 Observer 可能耗时
  5. 异常处理:一个 Observer 抛出异常可能影响其他 Observer

六、适用场景

  1. 事件驱动架构:Spring 事件、消息队列、EventBus
  2. GUI 组件:按钮点击、文本框变化等事件监听
  3. 消息推送:用户关注、订阅通知
  4. 数据同步:缓存更新、配置变更通知
  5. 监控告警:系统指标变化时触发告警
  6. 实时数据展示:股票行情、天气预报等实时更新的面板
  7. 微服务事件通知:订单状态变更、库存变化等

七、JDK / Spring 框架中的实际应用

框架应用位置说明
JDKjava.util.Observer / ObservableJDK 内置观察者模式(已废弃)
JDKjava.beans.PropertyChangeListener属性变更监听器
JDKjava.util.EventListener事件监听器接口
JDKjavax.servlet.ServletContextListenerServlet 上下文监听器
SpringApplicationEvent / ApplicationListenerSpring 事件机制
Spring@EventListener注解驱动的事件监听
Spring BootApplicationRunner / CommandLineRunner启动事件监听
GuavaEventBus轻量级事件总线
RxJavaObservable / 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) 没有提供有序通知机制。推荐使用 PropertyChangeListenerFlow API(Java 9+)或第三方库如 Guava EventBus。

Q5:如何处理观察者模式中的异步通知?

A:(1) 使用线程池异步执行通知:executor.submit(() -> observer.update());(2) 使用 Spring 的 @Async 注解;(3) 使用消息队列(RabbitMQ、Kafka)解耦;(4) 使用 Reactor/RxJava 响应式编程。异步通知可以提高吞吐量,但需要注意事务一致性和异常处理。