1. 观察者模式核心概念解析
观察者模式(Observer Pattern)是行为型设计模式中最常用的模式之一,它定义了对象间一对多的依赖关系。当目标对象(Subject)状态发生改变时,所有依赖它的观察者对象(Observer)都会自动收到通知并更新。这种模式在事件处理系统、消息队列、GUI组件交互等场景中应用广泛。
1.1 模式组成要素
典型的观察者模式包含四个核心角色:
Subject(目标对象):
- 维护观察者列表(List )
- 提供attach/detach方法管理观察者
- 状态变更时调用notify方法通知观察者
ConcreteSubject(具体目标):
- 存储具体状态信息
- 状态改变时触发通知逻辑
Observer(观察者接口):
- 定义update方法接口
- 约定通知时的回调规范
ConcreteObserver(具体观察者):
- 实现update方法
- 维护与Subject的引用关系
- 存储自身需要同步的状态
1.2 典型应用场景
观察者模式特别适合以下场景:
- 当一个对象的改变需要同时改变其他对象时
- 当对象需要在运行时动态决定通知对象时
- 当系统需要实现"广播通信"机制时
- 当不同对象间存在触发链关系时
实际开发中的典型用例包括:
- GUI事件处理(按钮点击、键盘输入)
- 发布-订阅系统(消息队列、事件总线)
- 数据监控系统(股票价格变动)
- 自动化任务触发(CI/CD流水线)
2. 观察者模式实现详解
2.1 Java标准实现方式
Java类库本身提供了Observable类和Observer接口,但自Java 9后已被标记为@Deprecated。现代Java开发更推荐自行实现:
// 观察者接口 public interface Observer { void update(String event); } // 具体观察者 public class ConcreteObserver implements Observer { private String id; public ConcreteObserver(String id) { this.id = id; } @Override public void update(String event) { System.out.println("Observer " + id + " received: " + event); } } // 目标对象 public class Subject { private List<Observer> observers = new ArrayList<>(); private String state; public void attach(Observer observer) { observers.add(observer); } public void detach(Observer observer) { observers.remove(observer); } public void setState(String newState) { this.state = newState; notifyObservers("State changed to: " + newState); } private void notifyObservers(String event) { for (Observer observer : observers) { observer.update(event); } } }2.2 推模型 vs 拉模型
观察者模式有两种主要实现方式:
推模型(Push Model):
- Subject将变更数据直接推送给Observer
- update方法携带详细变更信息
- 优点:Observer无需主动查询状态
- 缺点:可能推送不必要的数据
拉模型(Pull Model):
- Subject仅通知状态变更
- Observer需要主动查询所需数据
- 优点:Observer按需获取数据
- 缺点:增加耦合度
实际开发中推荐使用推拉结合的方式:
// 改进的Observer接口 public interface Observer<T> { void update(Subject<T> subject, T data); } // 使用时Observer可自主决定使用推送数据还是主动拉取3. 观察者模式高级应用
3.1 线程安全实现
在多线程环境下,观察者模式需要特别注意线程安全问题:
public class ConcurrentSubject { private final CopyOnWriteArrayList<Observer> observers = new CopyOnWriteArrayList<>(); public void attach(Observer observer) { observers.addIfAbsent(observer); } public void detach(Observer observer) { observers.remove(observer); } public void notifyObservers(String event) { for (Observer observer : observers) { executorService.submit(() -> observer.update(event)); } } }重要提示:在分布式系统中,需要考虑使用消息中间件(如Kafka、RabbitMQ)来实现跨进程的观察者模式。
3.2 基于注解的简化实现
现代框架通常提供注解方式来简化观察者模式:
// 使用Spring Event实现 @Component public class MyEventListener { @EventListener public void handleStateChange(StateChangeEvent event) { // 处理事件逻辑 } } // 发布事件 applicationContext.publishEvent(new StateChangeEvent(this, newState));3.3 响应式编程中的观察者
响应式流(如Reactor、RxJava)本质上是观察者模式的升级版:
Flux<String> source = Flux.just("data1", "data2", "data3"); source.subscribe( data -> System.out.println("Received: " + data), // onNext error -> System.err.println("Error: " + error), // onError () -> System.out.println("Done") // onComplete );4. 实战经验与性能优化
4.1 常见问题解决方案
内存泄漏:
- 场景:Observer未正确注销导致无法回收
- 解决:明确生命周期,在不再需要时调用detach
通知顺序不可控:
- 场景:多个Observer需要按特定顺序响应
- 解决:使用PriorityObserver或维护有序列表
过度通知:
- 场景:高频状态变更导致性能问题
- 解决:添加节流(throttling)或防抖(debounce)机制
4.2 性能优化技巧
- 差异化通知:
public void notifyObservers(ChangeType type) { observers.stream() .filter(o -> o.getInterestTypes().contains(type)) .forEach(o -> o.update(this)); }- 批量更新:
public void batchUpdate(List<Change> changes) { // 合并变化 Change merged = mergeChanges(changes); notifyObservers(merged); }- 异步处理:
private final ExecutorService executor = Executors.newFixedThreadPool(4); public void asyncNotify(Observer observer, Event event) { executor.submit(() -> { try { observer.update(event); } catch (Exception e) { logger.error("Notification failed", e); } }); }5. 模式变体与替代方案
5.1 中介者模式结合
当观察者间需要复杂交互时,可以引入中介者:
public class NotificationMediator { private Map<EventType, List<Observer>> mappings = new HashMap<>(); public void register(EventType type, Observer observer) { mappings.computeIfAbsent(type, k -> new ArrayList<>()).add(observer); } public void notify(Event event) { mappings.getOrDefault(event.getType(), Collections.emptyList()) .forEach(o -> o.update(event)); } }5.2 事件总线实现
更松耦合的实现方式:
public class EventBus { private static final Map<Class<?>, List<Consumer<?>>> handlers = new ConcurrentHashMap<>(); public static <T> void subscribe(Class<T> eventType, Consumer<T> handler) { handlers.computeIfAbsent(eventType, k -> new CopyOnWriteArrayList<>()) .add(handler); } public static <T> void publish(T event) { List<Consumer<?>> consumers = handlers.get(event.getClass()); if (consumers != null) { consumers.forEach(c -> ((Consumer<T>)c).accept(event)); } } }5.3 与发布-订阅模式对比
虽然观察者模式与发布-订阅模式相似,但存在关键差异:
| 特性 | 观察者模式 | 发布-订阅模式 |
|---|---|---|
| 耦合度 | 相对较高(直接引用) | 低(通过中间件) |
| 通信方式 | 同步 | 通常异步 |
| 关系 | 1对多 | 多对多 |
| 典型实现 | Java Observable | Kafka/RabbitMQ |
| 适用场景 | 单进程内部通信 | 分布式系统 |
在实际项目中,我通常会根据以下原则选择:
- 单体应用内部通信 → 观察者模式
- 微服务间通信 → 发布-订阅模式
- 需要强一致性 → 观察者模式
- 允许最终一致性 → 发布-订阅模式
6. 现代框架中的观察者模式
6.1 Spring事件机制
Spring框架提供了完善的事件发布-监听机制:
// 自定义事件 public class OrderCreatedEvent extends ApplicationEvent { private Order order; public OrderCreatedEvent(Object source, Order order) { super(source); this.order = order; } // getter... } // 发布事件 applicationContext.publishEvent(new OrderCreatedEvent(this, order)); // 监听事件 @Component public class OrderEventListener { @EventListener public void handleOrderCreated(OrderCreatedEvent event) { // 处理订单创建逻辑 } }Spring事件支持:
- 异步事件处理(@Async)
- 事件过滤(condition表达式)
- 事务绑定事件(@TransactionalEventListener)
6.2 Vue.js的响应式系统
前端框架Vue的响应式原理本质也是观察者模式:
// Vue 3 Composition API import { ref, watch } from 'vue' const count = ref(0) watch(count, (newVal, oldVal) => { console.log(`count changed from ${oldVal} to ${newVal}`) }) count.value++ // 触发观察者Vue实现的特点:
- 基于Proxy的自动依赖追踪
- 细粒度的更新调度
- 批量异步更新机制
6.3 React的Hooks机制
React的函数组件通过Hooks实现状态观察:
function Counter() { const [count, setCount] = useState(0); useEffect(() => { console.log(`Count updated: ${count}`); return () => { // 清理函数 }; }, [count]); // 只在count变化时运行 return <button onClick={() => setCount(c => c + 1)}>Increment</button>; }React观察模式特点:
- 显式声明依赖数组
- 自动处理组件卸载时的清理
- 支持效果合并和批处理
7. 设计考量与最佳实践
7.1 接口设计原则
单一职责原则:
- Subject只负责维护观察者列表和通知
- Observer只负责响应更新
接口隔离原则:
- 为不同类型的观察者定义专用接口
- 避免臃肿的Observer接口
依赖倒置原则:
- Observer依赖抽象接口而非具体实现
- 便于替换具体观察者
7.2 生命周期管理
完善的观察者模式应包含:
- 注册/注销机制:
public interface Observable { void addObserver(Observer o); void removeObserver(Observer o); void removeObservers(); }- 状态验证:
public void addObserver(Observer o) { if (o == null) throw new NullPointerException(); if (observers.contains(o)) return; observers.add(o); }- 资源清理:
public void close() { observers.clear(); // 释放其他资源 }7.3 测试策略
观察者模式的测试要点:
- Subject测试:
@Test public void shouldNotifyObserversOnStateChange() { TestObserver observer = new TestObserver(); subject.addObserver(observer); subject.setState("new"); assertTrue(observer.isNotified()); assertEquals("new", observer.getLastState()); }- Observer测试:
@Test public void shouldUpdateStateWhenNotified() { Observer observer = new ConcreteObserver(); observer.update("test"); assertEquals("test", observer.getCurrentState()); }- 集成测试:
@Test public void shouldMaintainConsistencyBetweenMultipleObservers() { // 设置多个观察者 // 改变主题状态 // 验证所有观察者状态一致 }8. 复杂场景处理方案
8.1 循环依赖问题
当观察者反过来修改Subject状态时可能导致无限通知:
解决方案:
- 添加修改标记:
public void setState(String newState) { if (isUpdating) return; isUpdating = true; this.state = newState; notifyObservers(); isUpdating = false; }- 使用命令模式:
public void applyChange(ChangeCommand command) { command.execute(this); if (!command.isSilent()) { notifyObservers(); } }8.2 部分更新通知
当只有部分状态变化时需要通知:
public void partialUpdate(String fieldName, Object value) { if (!shouldNotify(fieldName, value)) return; updateField(fieldName, value); notifyObservers(new FieldChangeEvent(fieldName, value)); } private boolean shouldNotify(String field, Object newValue) { Object oldValue = getFieldValue(field); return !Objects.equals(oldValue, newValue); }8.3 跨进程观察者
分布式系统下的实现方案:
- 基于RPC:
public class RemoteObserver implements Observer { private final ObserverService stub; public void update(StateChange change) { stub.notifyRemote(change); // 远程调用 } }- 基于消息队列:
public class MessageQueueSubject { private final MessageProducer producer; public void stateChanged(State newState) { producer.send(new StateChangeMessage(newState)); } }- 基于Webhook:
public class WebhookObserver implements Observer { private final String callbackUrl; public void update(StateChange change) { HttpClient.post(callbackUrl, serialize(change)); } }9. 经典案例:股票价格通知系统
9.1 系统设计
模拟股票市场监控系统:
// 股票数据主题 public class StockSubject { private Map<String, Double> prices = new HashMap<>(); private List<StockObserver> observers = new ArrayList<>(); public void updatePrice(String symbol, double price) { if (price != prices.getOrDefault(symbol, 0.0)) { prices.put(symbol, price); notifyObservers(symbol, price); } } private void notifyObservers(String symbol, double price) { observers.forEach(o -> o.onPriceChange(symbol, price)); } } // 观察者接口 public interface StockObserver { void onPriceChange(String symbol, double price); } // 具体观察者:价格报警器 public class PriceAlert implements StockObserver { private String symbol; private double threshold; public PriceAlert(String symbol, double threshold) { this.symbol = symbol; this.threshold = threshold; } @Override public void onPriceChange(String s, double price) { if (s.equals(symbol) && price >= threshold) { System.out.println("Alert! " + symbol + " reached " + price); } } }9.2 扩展功能
- 价格变化历史:
public class PriceHistory implements StockObserver { private Map<String, List<Double>> history = new HashMap<>(); @Override public void onPriceChange(String symbol, double price) { history.computeIfAbsent(symbol, k -> new ArrayList<>()) .add(price); } }- 移动平均计算:
public class MovingAverage implements StockObserver { private final int windowSize; private Map<String, Queue<Double>> windows = new HashMap<>(); @Override public void onPriceChange(String symbol, double price) { Queue<Double> window = windows.computeIfAbsent(symbol, k -> new LinkedList<>()); window.offer(price); if (window.size() > windowSize) window.poll(); double avg = window.stream().mapToDouble(d -> d).average().orElse(0); System.out.println(symbol + " MA(" + windowSize + "): " + avg); } }- 组合观察:
public class PortfolioTracker implements StockObserver { private Set<String> symbols = new HashSet<>(); public void addSymbol(String symbol) { symbols.add(symbol); } @Override public void onPriceChange(String symbol, double price) { if (symbols.contains(symbol)) { updatePortfolioValue(symbol, price); } } }10. 性能关键型实现
10.1 高效通知策略
- 差异化通知:
public void notifyObservers(ChangeType type) { for (Observer o : observers) { if (o.getInterestTypes().contains(type)) { o.update(this); } } }- 批量通知:
public void batchNotify(List<Change> changes) { if (changes.isEmpty()) return; ChangeEvent event = mergeChanges(changes); for (Observer o : observers) { o.update(event); } }- 异步通知:
private final ExecutorService executor = Executors.newWorkStealingPool(); public void asyncNotify(Observer o, Event event) { executor.submit(() -> { try { o.update(event); } catch (Exception e) { logger.error("Notification failed", e); } }); }10.2 内存优化技巧
- 观察者弱引用:
private final List<WeakReference<Observer>> observers = new ArrayList<>(); public void addObserver(Observer o) { observers.add(new WeakReference<>(o)); } private void notifyObservers() { Iterator<WeakReference<Observer>> it = observers.iterator(); while (it.hasNext()) { Observer o = it.next().get(); if (o != null) { o.update(this); } else { it.remove(); // 清理被GC的观察者 } } }- 事件对象池:
private final ObjectPool<Event> eventPool = new ObjectPool<>(Event::new); public void notifyObservers() { Event event = eventPool.borrowObject(); try { event.setData(getState()); for (Observer o : observers) { o.update(event); } } finally { eventPool.returnObject(event); } }- 增量更新:
public void notifyObservers(DeltaUpdate delta) { for (Observer o : observers) { if (o.supportsDelta()) { o.applyDelta(delta); // 只发送变化部分 } else { o.update(getFullState()); } } }11. 与其他模式的协作
11.1 与责任链模式结合
实现按优先级处理通知:
public class ChainedObserver implements Observer { private List<Observer> chain; public ChainedObserver(List<Observer> chain) { this.chain = new ArrayList<>(chain); } @Override public void update(Event event) { for (Observer o : chain) { if (event.isHandled()) break; o.update(event); } } }11.2 与策略模式结合
动态切换通知策略:
public class SmartSubject { private NotificationStrategy strategy; public void setStrategy(NotificationStrategy strategy) { this.strategy = strategy; } public void notifyObservers() { strategy.notify(observers, getState()); } } interface NotificationStrategy { void notify(List<Observer> observers, State state); }11.3 与备忘录模式结合
实现状态回滚通知:
public class UndoableSubject { private final Stack<Memento> history = new Stack<>(); public void setState(State newState) { history.push(createMemento()); this.state = newState; notifyObservers(); } public void undo() { if (!history.isEmpty()) { restoreFromMemento(history.pop()); notifyObservers("UNDO"); } } }12. 反模式与常见误区
12.1 典型反模式
过度通知:
- 问题:频繁状态变化导致通知风暴
- 解决:添加节流机制或批量更新
隐式依赖:
- 问题:Observer对Subject的内部实现有假设
- 解决:明确定义通知契约
循环通知:
- 问题:Observer更新又触发Subject修改
- 解决:添加修改标记或使用命令模式
12.2 性能陷阱
同步阻塞通知:
// 错误做法:长耗时Observer阻塞通知链 public void notifyObservers() { for (Observer o : observers) { o.update(this); // 可能阻塞 } }Observer内存泄漏:
// 错误做法:忘记注销Observer public void init() { subject.addObserver(this); // 但未在destroy时移除 }过度细粒度通知:
// 错误做法:每个字段变化都通知 public void setX(int x) { notify(); } public void setY(int y) { notify(); }
12.3 设计原则违反
违反单一职责:
// 错误做法:Subject承担过多责任 public class Subject { public void businessLogic() { /*...*/ } public void saveToDatabase() { /*...*/ } public void notifyObservers() { /*...*/ } }暴露实现细节:
// 错误做法:Observer需要了解Subject内部 public void update(Subject s) { String state = s.getInternalState().getNestedField(); }缺乏生命周期管理:
// 错误做法:没有提供注销机制 public class Subject { private final List<Observer> observers = new ArrayList<>(); // 缺少removeObserver方法 }
13. 测试策略与验证
13.1 单元测试要点
- Subject测试:
@Test public void shouldNotifyAllObservers() { Subject subject = new Subject(); MockObserver obs1 = new MockObserver(); MockObserver obs2 = new MockObserver(); subject.attach(obs1); subject.attach(obs2); subject.setState("new"); assertTrue(obs1.isNotified()); assertTrue(obs2.isNotified()); }- Observer测试:
@Test public void shouldUpdateStateWhenNotified() { TestObserver observer = new TestObserver(); observer.update("test"); assertEquals("test", observer.getState()); }13.2 集成测试方案
- 顺序验证:
@Test public void shouldNotifyInRegistrationOrder() { Subject subject = new Subject(); List<String> notificationOrder = new ArrayList<>(); subject.attach(e -> notificationOrder.add("first")); subject.attach(e -> notificationOrder.add("second")); subject.notifyObservers(); assertEquals(Arrays.asList("first", "second"), notificationOrder); }- 异常处理测试:
@Test public void shouldContinueWhenObserverFails() { Subject subject = new Subject(); subject.attach(e -> { throw new RuntimeException(); }); subject.attach(e -> results.add("success")); subject.notifyObservers(); assertEquals(1, results.size()); }13.3 性能测试
- 吞吐量测试:
@Test public void throughputTest() { Subject subject = new Subject(); for (int i = 0; i < 1000; i++) { subject.attach(new DummyObserver()); } long start = System.nanoTime(); subject.notifyObservers(); long duration = System.nanoTime() - start; assertTrue(duration < TimeUnit.MILLISECONDS.toNanos(100)); }- 内存泄漏测试:
@Test public void memoryLeakTest() { Subject subject = new Subject(); WeakReference<Observer> ref = new WeakReference<>(new DummyObserver()); subject.attach(ref.get()); ref.clear(); System.gc(); subject.notifyObservers(); assertEquals(0, subject.countObservers()); }14. 现代语言特性应用
14.1 Java函数式实现
利用Java 8+特性简化代码:
public class FunctionalSubject { private final List<Consumer<String>> listeners = new ArrayList<>(); public void addListener(Consumer<String> listener) { listeners.add(listener); } public void changeState(String newState) { state = newState; listeners.forEach(l -> l.accept(newState)); } } // 使用 subject.addListener(state -> System.out.println("State: " + state));14.2 Kotlin委托属性
Kotlin提供更优雅的实现:
class ObservableProperty(var value: String) { private val observers = mutableListOf<(String) -> Unit>() fun addObserver(observer: (String) -> Unit) { observers.add(observer) } operator fun setValue(thisRef: Any?, property: KProperty<*>, newValue: String) { if (value != newValue) { value = newValue observers.forEach { it(newValue) } } } } // 使用 var state by ObservableProperty("init") state.addObserver { println("Changed to $it") } state = "new" // 自动通知14.3 C#事件机制
C#原生支持事件语法:
public class Subject { public event EventHandler<string> StateChanged; private string state; public string State { get => state; set { if (state != value) { state = value; StateChanged?.Invoke(this, value); } } } } // 使用 subject.StateChanged += (sender, newState) => Console.WriteLine(newState); subject.State = "new";15. 架构层面的应用
15.1 领域事件模式
在DDD中实现领域事件:
public abstract class DomainEvent { private final Instant occurredOn = Instant.now(); } public class OrderCreatedEvent extends DomainEvent { private final OrderId orderId; // ... } // 发布 public class Order { public static Order create(OrderDetails details) { Order order = new Order(details); DomainEventPublisher.publish(new OrderCreatedEvent(order.getId())); return order; } } // 订阅 public class OrderMetrics { @Subscribe public void onOrderCreated(OrderCreatedEvent event) { // 更新指标 } }15.2 CQRS模式中的使用
在命令查询职责分离架构中:
// 写模型 public class InventoryCommandHandler { public void handle(UpdateInventoryCommand cmd) { inventory.update(cmd.productId(), cmd.quantity()); eventBus.publish(new InventoryUpdatedEvent(...)); } } // 读模型 public class InventoryReadModel { @Subscribe public void onInventoryUpdate(InventoryUpdatedEvent event) { // 更新物化视图 } }15.3 微服务间事件驱动
跨服务事件通知:
// 通过消息中间件 public class OrderService { @Transactional public void completeOrder(OrderId id) { orderRepository.complete(id); kafkaTemplate.send("order-completed", new OrderCompletedEvent(id)); } } @Service public class DeliveryService { @KafkaListener(topics = "order-completed") public void handleOrderCompleted(OrderCompletedEvent event) { deliveryRepository.scheduleFor(event.orderId()); } }16. 可视化调试工具
16.1 通知流程追踪
添加调试支持:
public class DebuggableSubject extends Subject { private final NotificationTracer tracer; @Override protected void notifyObservers() { tracer.traceStart(observers.size()); super.notifyObservers(); tracer.traceEnd(); } } public interface NotificationTracer { void traceStart(int observerCount); void traceObserverCalled(Observer o); void traceEnd(); }16.2 依赖关系可视化
生成观察者关系图:
public class GraphvizExporter { public String export(Subject subject) { StringBuilder dot = new StringBuilder("digraph G {\n"); subject.getObservers().forEach(o -> dot.append(" Subject -> ").append(o.getClass().getSimpleName()).append("\n") ); return dot.append("}").toString(); } }16.3 性能分析工具
监控通知性能:
public class MonitoredSubject extends Subject { private final NotificationMetrics metrics; @Override protected void notifyObservers() { long start = System.nanoTime(); super.notifyObservers(); metrics.recordDuration(System.nanoTime() - start); } } public interface NotificationMetrics { void recordDuration(long nanos); Stats getStats(); }17. 模式演进与替代方案
17.1 响应式流规范
现代替代方案:
// 使用Java Flow API public class FlowSubject { private final SubmissionPublisher<String> publisher = new SubmissionPublisher<>(); public void changeState(String newState) { state = newState; publisher.submit(newState); } public Flow.Publisher<String> asPublisher() { return publisher; } } // 使用 subject.asPublisher().subscribe(new Flow.Subscriber<>() { @Override public void onNext(String item) { System.out.println("Received: " + item); } // 其他方法... });17.2 Actor模型实现
基于消息传递的方案:
// 使用Akka框架 public class SubjectActor extends AbstractActor { private List<ActorRef> observers = new ArrayList<>(); @Override public Receive createReceive() { return receiveBuilder() .match(Subscribe.class, this::onSubscribe) .match(StateChange.class, this::onStateChange) .build(); } private void onSubscribe(Subscribe msg) { observers.add(getSender()); } private void onStateChange(StateChange change) { observers.forEach(o -> o.tell(change, getSelf())); } }17.3 数据绑定框架
声明式替代方案:
// 使用Vue.js const app = new Vue({ data: { message: 'Hello' } }) // 自动建立观察关系 app.$watch('message', (newVal) => { console.log('Message changed:', newVal) }) app.message = 'Updated' // 触发回调18. 行业应用案例分析
18.1 GUI框架中的实现
典型实现:Java Swing的Button监听
JButton button = new JButton("Click"); button.addActionListener(e -> { System.out.println("Button clicked"); });内部机制:
- Button维护ActionListener列表
- 点击事件触发notifyListeners调用
- 每个listener的actionPerformed被调用
18.2 游戏开发中的应用
游戏事件系统示例:
public class GameEventSystem { private Dictionary<Type, List<object>> handlers = new(); public void Subscribe<T>(Action<T> handler) { var type = typeof(T); if (!handlers.ContainsKey(type)) { handlers[type] = new List<object>(); } handlers[type].Add(handler); } public void Publish<T>(T event) { if (handlers.TryGetValue(typeof(T), out var list)) { foreach (Action<T> handler in list.Cast<Action<T>>()) { handler(event); } } } } // 使用 eventSystem.Subscribe<PlayerDiedEvent>(e => GameOver()); eventSystem.Publish(new PlayerDiedEvent());18.3 金融交易系统
价格变动通知实现:
public class PriceFeed { private final Map<String, List<PriceListener>> listeners = new ConcurrentHashMap<>(); public void subscribe(String symbol, PriceListener listener) { listeners.computeIfAbsent(symbol, k -> new CopyOnWriteArrayList<>()) .add(listener); } public void onMarketData(MarketData data) { List<PriceListener> symbolListeners = listeners.get(data.getSymbol()); if (symbolListeners != null) { symbolListeners.forEach(l -> l.onPriceChange(data)); } } } interface PriceListener { void onPriceChange(MarketData data); }19. 反模式与重构方案
19.1 常见反模式
上帝观察者:
- 问题:单个Observer处理所有通知导致臃肿
- 解决:按职责拆分多个专门Observer
过度通知:
- 问题:频繁无关状态变化触发通知
- 解决:细化变更类型或添加过滤条件
隐式耦合:
- 问题:Observer对Subject有非接口依赖
- 解决:明确定义通知契约和数据结构
19.2 重构为发布-订阅
当观察者模式变得复杂时的重构方向:
// 重构前 public class Subject { private List<Observer> observers = new ArrayList<>(); public void notifyAll() { observers.forEach(o -> o.update(getFullState())); } } // 重构后 public class EventBus { private Map<Class<?>, List<Consumer<?>>> handlers = new HashMap<>(); public <T> void publish(T event) { handlers.getOrDefault(event.getClass(), List.of()) .forEach(c -> ((Consumer<T>)c).accept(event)); } }19.3 性能优化重构
从同步到异步的改造:
// 原始版本 public class Subject { public void notifyObservers() { observers.forEach(o -> o.update(state)); // 同步阻塞 } } // 优化版本 public class AsyncSubject { private final Executor executor; public void notifyObservers() { observers.forEach(o -> executor.execute(() -> o.update(state)) ); } }20. 个人实践心得
在实际项目中应用观察者模式时,我总结了以下几点经验:
明确生命周期管理:
- 注册/注销必须成对出现
- 特别注意在容器环境中的销毁处理
- 推荐使用try-with-resources或类似机制
通知内容设计:
- 事件对象应包含足够上下文
- 但避免暴露Subject内部状态
- 考虑使用不可变事件对象
**异常处理策略