设计模式详解-观察者模式

设计模式详解:观察者模式

一、模式概述

观察者模式(Observer Pattern)是行为型设计模式中最具解耦价值的模式,其核心意图在于定义对象间的一种一对多依赖关系,使得每当一个对象状态发生改变时,所有依赖于它的对象都得到通知并被自动更新。这一模式将发布者与订阅者解耦,使它们无需相互知晓,却能协同响应状态变化。

观察者模式的命名源自其隐喻——"观察者"如同订阅报纸的读者,无需主动询问,每当新闻发生(状态变化),报纸(发布者)自动投递(通知)到所有订阅者。在软件系统中,这一模式是事件驱动架构、响应式编程、消息总线等现代基础设施的理论基石。

观察者模式的深层价值在于依赖倒置的极致体现。传统设计中,对象A需要知道对象B的存在,直接调用B的方法;观察者模式中,对象A仅维护一个观察者列表,通过抽象接口通知它们,B只需实现该接口即可接收通知。这种间接性使系统具备极强的扩展能力——新增观察者无需修改发布者,发布者的变化也不影响既有观察者。

二、模式结构

观察者模式包含两个核心角色,形成松散的发布-订阅关系:

主题(Subject):也称为发布者或目标,知道它的观察者,提供注册和删除观察者的接口,当状态变化时通知所有已注册的观察者。

观察者(Observer):定义一个更新接口,用于接收主题的通知。具体观察者维护一个指向主题对象的引用,以便获取更新后的状态。

经典实现中,主题维护观察者列表,状态变化时遍历调用update()方法。现代演进中,观察者模式发展为事件总线、消息队列、响应式流等更复杂的形态,但核心思想不变。

三、深度案例:企业级实时风控事件系统

以下展示一个真实场景下的观察者模式应用——金融交易平台的实时风控事件系统,处理交易、账户、设备、行为等多维度事件,驱动规则计算、告警通知、决策响应、数据归档等异构处理逻辑。

3.1 问题域分析:紧耦合的灾难

java
// 反模式:交易服务直接硬编码所有后续处理 public class NaiveTransactionService { private final RuleEngine ruleEngine; private final AlertService alertService; private final NotificationService notificationService; private final AuditService auditService; private final StatisticsService statisticsService; private final DataWarehouseService dataWarehouseService; private final AntiFraudService antiFraudService; private final ComplianceService complianceService; public TransactionResult processTransaction(TransactionRequest request) { // 执行交易... Transaction transaction = execute(request); // 硬编码调用所有后续处理 ruleEngine.evaluate(transaction); alertService.checkThresholds(transaction); notificationService.sendNotification(transaction); auditService.recordTransaction(transaction); statisticsService.updateMetrics(transaction); dataWarehouseService.syncTransaction(transaction); antiFraudService.analyze(transaction); complianceService.checkRegulatoryRequirements(transaction); // 新增处理?修改此代码,重新部署,全量回归 // 某个服务故障?阻塞交易完成 return TransactionResult.success(transaction); } }

上述代码每新增一种后置处理,需修改核心交易服务,且任一服务故障或延迟直接影响交易响应时间。观察者模式将此重构为事件驱动的松耦合架构。

3.2 抽象层:事件与观察者契约

java
/** * 抽象观察者:事件处理器 * 支持同步/异步、有序/并行、过滤/转换等高级特性 */ public interface EventObserver<E extends DomainEvent> { /** * 观察者标识 */ String getObserverId(); /** * 感兴趣的事件类型 */ Set<Class<? extends E>> getInterestedEventTypes(); /** * 事件优先级(数值越小优先级越高) */ int getPriority(); /** * 执行模式:同步或异步 */ ExecutionMode getExecutionMode(); /** * 是否支持批量处理 */ boolean supportsBatching(); /** * 批量大小(支持批量时有效) */ int getBatchSize(); /** * 处理事件 * @param event 事件对象 * @param context 事件上下文(包含元数据、追踪信息等) */ void onEvent(E event, EventContext context); /** * 批量处理事件 */ default void onBatch(List<E> events, EventContext context) { events.forEach(e -> onEvent(e, context)); } /** * 错误处理策略 */ default ErrorPolicy getErrorPolicy() { return ErrorPolicy.RETRY_THEN_DEAD_LETTER; } } /** * 事件上下文:传递追踪、租户、时间等元数据 */ public class EventContext { private final String traceId; private final String tenantId; private final Instant eventTime; private final Instant receiveTime; private final String sourceService; private final Map<String, Object> metadata; // 传播上下文到子线程 public EventContext fork() { return new EventContext(traceId, tenantId, Instant.now(), receiveTime, sourceService, new HashMap<>(metadata)); } } /** * 抽象主题:事件总线 * 定义发布-订阅的核心接口 */ public interface EventBus { /** * 注册观察者 */ void register(EventObserver<? extends DomainEvent> observer); /** * 注销观察者 */ void unregister(String observerId); /** * 发布事件 */ <E extends DomainEvent> void publish(E event); /** * 发布事件(带上下文) */ <E extends DomainEvent> void publish(E event, EventContext context); /** * 发布事件并等待所有同步观察者处理完成 */ <E extends DomainEvent> void publishAndWait(E event, Duration timeout); }

3.3 具体观察者:异构处理逻辑

java
/** * 具体观察者:实时规则计算 * 同步执行,高优先级,直接影响交易结果 */ @Component @Order(10) public class RealtimeRuleObserver implements EventObserver<TransactionEvent> { private final RuleEngine ruleEngine; private final DecisionCache decisionCache; @Override public String getObserverId() { return "realtime-rule"; } @Override public Set<Class<? extends TransactionEvent>> getInterestedEventTypes() { return Set.of(TransactionCreatedEvent.class, TransactionAmendedEvent.class); } @Override public int getPriority() { return 10; } @Override public ExecutionMode getExecutionMode() { return ExecutionMode.SYNC; } @Override public boolean supportsBatching() { return false; } // 需逐笔实时判定 @Override public void onEvent(TransactionEvent event, EventContext context) { // 缓存检查:相同特征的交易复用近期判定结果 String cacheKey = generateCacheKey(event); RuleDecision cached = decisionCache.get(cacheKey); if (cached != null && cached.isValid()) { applyDecision(event, cached); return; } // 执行规则计算 RuleDecision decision = ruleEngine.evaluate(event, context); // 缓存结果 decisionCache.put(cacheKey, decision); // 应用决策 applyDecision(event, decision); } private void applyDecision(TransactionEvent event, RuleDecision decision) { if (decision.getAction() == RuleAction.BLOCK) { throw new TransactionBlockedException(decision.getReasonCodes()); } if (decision.getAction() == RuleAction.CHALLENGE) { event.setChallengeRequired(true); event.setChallengeMethod(decision.getChallengeMethod()); } // ALLOW则无需处理,继续流程 } } /** * 具体观察者:反欺诈分析 * 异步执行,可批量,不阻塞主交易流程 */ @Component @Order(50) public class AntiFraudObserver implements EventObserver<TransactionEvent> { private final FraudDetectionEngine fraudEngine; private final KafkaTemplate<String, FraudEvent> kafkaTemplate; @Override public String getObserverId() { return "anti-fraud"; } @Override public Set<Class<? extends TransactionEvent>> getInterestedEventTypes() { return Set.of(TransactionEvent.class); // 所有交易事件 } @Override public int getPriority() { return 50; } @Override public ExecutionMode getExecutionMode() { return ExecutionMode.ASYNC; } @Override public boolean supportsBatching() { return true; } @Override public int getBatchSize() { return 100; } @Override public void onBatch(List<TransactionEvent> events, EventContext context) { // 批量特征提取 List<FraudFeatureVector> features = events.parallelStream() .map(e -> extractFeatures(e, context)) .collect(Collectors.toList()); // 批量模型预测 List<FraudScore> scores = fraudEngine.scoreBatch(features); // 处理高风险交易 for (int i = 0; i < events.size(); i++) { if (scores.get(i).isHighRisk()) { kafkaTemplate.send("fraud-alerts", new FraudAlertEvent( events.get(i), scores.get(i))); } } } @Override public void onEvent(TransactionEvent event, EventContext context) { // 单条处理(批量未凑齐时) FraudScore score = fraudEngine.score(extractFeatures(event, context)); if (score.isHighRisk()) { kafkaTemplate.send("fraud-alerts", new FraudAlertEvent(event, score)); } } private FraudFeatureVector extractFeatures(TransactionEvent event, EventContext context) { return FraudFeatureVector.builder() .transactionId(event.getTransactionId()) .amount(event.getAmount()) .deviceFingerprint(event.getDeviceFingerprint()) .location(event.getLocation()) .merchantCategory(event.getMerchantCategory()) .timeOfDay(event.getTransactionTime().getHour()) .userHistoricalPattern(loadUserPattern(event.getUserId())) .build(); } } /** * 具体观察者:审计日志 * 同步执行,最高优先级,确保事件被记录 */ @Component @Order(1) public class AuditLogObserver implements EventObserver<DomainEvent> { private final AuditLogRepository auditRepository; private final AsyncAuditBuffer asyncBuffer; @Override public String getObserverId() { return "audit-log"; } @Override public Set<Class<? extends DomainEvent>> getInterestedEventTypes() { return Set.of(DomainEvent.class); // 所有领域事件 } @Override public int getPriority() { return 1; } // 最高优先级 @Override public ExecutionMode getExecutionMode() { return ExecutionMode.SYNC; } @Override public ErrorPolicy getErrorPolicy() { return ErrorPolicy.FAIL_FAST; // 审计失败则交易失败 } @Override public void onEvent(DomainEvent event, EventContext context) { AuditEntry entry = AuditEntry.builder() .eventId(event.getEventId()) .eventType(event.getClass().getSimpleName()) .traceId(context.getTraceId()) .tenantId(context.getTenantId()) .eventTime(context.getEventTime()) .payload(maskSensitiveData(event)) .sourceService(context.getSourceService()) .build(); // 同步写入主库,确保事务一致性 auditRepository.insert(entry); // 异步复制到日志平台 asyncBuffer.buffer(entry); } private String maskSensitiveData(DomainEvent event) { // 脱敏处理 String json = JsonUtils.toJson(event); return DataMasker.mask(json, MaskingPolicy.AUDIT); } } /** * 具体观察者:数据仓库同步 * 异步执行,容忍延迟,失败可重试 */ @Component @Order(100) public class DataWarehouseSyncObserver implements EventObserver<TransactionEvent> { private final DataWarehouseClient warehouseClient; private final RetryTemplate retryTemplate; @Override public String getObserverId() { return "data-warehouse"; } @Override public ExecutionMode getExecutionMode() { return ExecutionMode.ASYNC; } @Override public ErrorPolicy getErrorPolicy() { return ErrorPolicy.RETRY_WITH_BACKOFF; } @Override public void onEvent(TransactionEvent event, EventContext context) { WarehouseRecord record = convertToWarehouseFormat(event); retryTemplate.execute(context -> { warehouseClient.append(record); return null; }); } @Override public void onBatch(List<TransactionEvent> events, EventContext context) { // 批量加载优化 List<WarehouseRecord> records = events.stream() .map(this::convertToWarehouseFormat) .collect(Collectors.toList()); warehouseClient.appendBatch(records); } } /** * 具体观察者:实时指标推送 * 异步,极低延迟,内存聚合 */ @Component @Order(30) public class MetricsObserver implements EventObserver<DomainEvent> { private final MeterRegistry meterRegistry; private final CountersAggregator aggregator; @Override public String getObserverId() { return "metrics"; } @Override public ExecutionMode getExecutionMode() { return ExecutionMode.ASYNC; } @Override public boolean supportsBatching() { return true; } @Override public int getBatchSize() { return 1000; } @Override public void onEvent(DomainEvent event, EventContext context) { // 内存计数器,无阻塞 aggregator.increment(event.getClass().getSimpleName()); // 细粒度指标 if (event instanceof TransactionEvent) { TransactionEvent te = (TransactionEvent) event; meterRegistry.counter("transactions", "type", te.getTransactionType(), "status", te.getStatus(), "channel", te.getChannel() ).increment(); meterRegistry.distributionSummary("transaction.amount") .record(te.getAmount().doubleValue()); } } @Override public void onBatch(List<DomainEvent> events, EventContext context) { // 批量聚合后刷新 Map<String, Long> counts = events.stream() .collect(Collectors.groupingBy( e -> e.getClass().getSimpleName(), Collectors.counting() )); counts.forEach((type, count) -> aggregator.increment(type, count)); } }

3.4 主题实现:高性能事件总线

java
/** * 主题实现:分层事件总线 * 支持同步/异步通道、优先级调度、背压控制 */ @Component public class TieredEventBus implements EventBus { private final ObserverRegistry registry; private final SyncEventExecutor syncExecutor; private final AsyncEventExecutor asyncExecutor; private final EventDeadLetterQueue deadLetterQueue; private final MeterRegistry meterRegistry; // 同步通道:直接内存调用 private final Map<Class<?>, List<EventObserver>> syncObservers = new ConcurrentHashMap<>(); // 异步通道:按优先级分队列 private final PriorityBlockingQueue<AsyncEventTask> asyncQueue = new PriorityBlockingQueue<>(1000, Comparator.comparingInt(t -> t.priority)); @Autowired public TieredEventBus(ObserverRegistry registry, SyncEventExecutor syncExecutor, AsyncEventExecutor asyncExecutor, EventDeadLetterQueue deadLetterQueue, MeterRegistry meterRegistry) { this.registry = registry; this.syncExecutor = syncExecutor; this.asyncExecutor = asyncExecutor; this.deadLetterQueue = deadLetterQueue; this.meterRegistry = meterRegistry; // 启动异步消费者 startAsyncConsumers(); } @Override public void register(EventObserver<? extends DomainEvent> observer) { registry.register(observer); // 按执行模式分类索引 for (Class<? extends DomainEvent> eventType : observer.getInterestedEventTypes()) { if (observer.getExecutionMode() == ExecutionMode.SYNC) { syncObservers.computeIfAbsent(eventType, k -> new CopyOnWriteArrayList<>()) .add(observer); } } meterRegistry.counter("eventbus.observer.registered").increment(); } @Override public void unregister(String observerId) { registry.unregister(observerId); // 从索引中移除... } @Override public <E extends DomainEvent> void publish(E event) { publish(event, EventContext.create()); } @Override public <E extends DomainEvent> void publish(E event, EventContext context) { long startTime = System.currentTimeMillis(); // 1. 获取所有感兴趣的观察者 List<EventObserver<E>> observers = registry.findObservers(event.getClass()); // 2. 按优先级排序 observers.sort(Comparator.comparingInt(EventObserver::getPriority)); // 3. 分离同步与异步 List<EventObserver<E>> syncList = new ArrayList<>(); List<EventObserver<E>> asyncList = new ArrayList<>(); for (EventObserver<E> observer : observers) { if (observer.getExecutionMode() == ExecutionMode.SYNC) { syncList.add(observer); } else { asyncList.add(observer); } } // 4. 执行同步观察者(阻塞,影响交易响应) for (EventObserver<E> observer : syncList) { executeSync(observer, event, context); } // 5. 提交异步观察者(非阻塞) for (EventObserver<E> observer : asyncList) { submitAsync(observer, event, context); } meterRegistry.timer("eventbus.publish").record( System.currentTimeMillis() - startTime, TimeUnit.MILLISECONDS); } @Override public <E extends DomainEvent> void publishAndWait(E event, Duration timeout) { publish(event); // 等待所有相关异步任务完成 CountDownLatch latch = new CountDownLatch(countAsyncObservers(event.getClass())); // 注册临时监听器... try { if (!latch.await(timeout.toMillis(), TimeUnit.MILLISECONDS)) { throw new EventTimeoutException("异步处理超时"); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new EventInterruptedException(e); } } private <E extends DomainEvent> void executeSync(EventObserver<E> observer, E event, EventContext context) { try { observer.onEvent(event, context); meterRegistry.counter("eventbus.sync.success", "observer", observer.getObserverId()).increment(); } catch (Exception e) { meterRegistry.counter("eventbus.sync.failure", "observer", observer.getObserverId()).increment(); switch (observer.getErrorPolicy()) { case FAIL_FAST: throw new SyncObserverException(observer.getObserverId(), e); case SKIP: log.warn("同步观察者失败,跳过: {}", observer.getObserverId(), e); break; case RETRY_THEN_DEAD_LETTER: // 同步重试一次 retrySync(observer, event, context); break; } } } private <E extends DomainEvent> void submitAsync(EventObserver<E> observer, E event, EventContext context) { AsyncEventTask task = new AsyncEventTask( observer.getPriority(), observer.getObserverId(), () -> { try { if (observer.supportsBatching()) { // 批量收集后处理 batchingExecutor.collect(observer, event, context); } else { observer.onEvent(event, context); } } catch (Exception e) { handleAsyncError(observer, event, context, e); } } ); asyncQueue.offer(task); } private void startAsyncConsumers() { int consumerCount = Runtime.getRuntime().availableProcessors(); for (int i = 0; i < consumerCount; i++) { Thread consumer = new Thread(this::asyncConsumeLoop, "eventbus-async-" + i); consumer.setDaemon(true); consumer.start(); } } private void asyncConsumeLoop() { while (!Thread.interrupted()) { try { AsyncEventTask task = asyncQueue.take(); task.runnable.run(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } private <E extends DomainEvent> void handleAsyncError(EventObserver<E> observer, E event, EventContext context, Exception e) { switch (observer.getErrorPolicy()) { case RETTRY_WITH_BACKOFF: scheduleRetry(observer, event, context, e, 1); break; case DEAD_LETTER: deadLetterQueue.enqueue(new DeadLetterEvent(event, context, observer.getObserverId(), e)); break; case IGNORE: log.error("异步观察者失败,忽略: {}", observer.getObserverId(), e); break; } } }

3.5 响应式演进:Reactive Streams

java
/** * 响应式观察者:基于Project Reactor的背压感知事件流 */ public class ReactiveEventBus implements EventBus { private final Sinks.Many<DomainEvent> sink = Sinks.many().multicast().onBackpressureBuffer(); private final Flux<DomainEvent> sharedFlux = sink.asFlux().publish().autoConnect(); @Override public <E extends DomainEvent> void publish(E event) { sink.tryEmitNext(event); } /** * 订阅响应式流 */ public <E extends DomainEvent> Disposable subscribe( Class<E> eventType, Function<E, Mono<Void>> handler, int prefetch) { return sharedFlux .ofType(eventType) .onBackpressureBuffer(1000) // 背压缓冲 .flatMap(event -> handler.apply(event) .onErrorResume(e -> { log.error("处理失败", e); return Mono.empty(); }), prefetch) // 并发度控制 .subscribe(); } } // 使用示例 @Component public class ReactiveAntiFraudAdapter { @PostConstruct public void subscribe() { eventBus.subscribe( TransactionEvent.class, event -> fraudEngine.analyze(event) .flatMap(score -> score.isHighRisk() ? alertService.sendAlert(event, score).then() : Mono.empty()), 10 // 最多10个并发分析 ); } }

四、观察者模式的高级主题

4.1 事件溯源与CQRS

java
/** * 事件存储:持久化所有领域事件,作为系统真相来源 */ public class EventStore { private final EventRepository eventRepository; private final EventBus eventBus; /** * 追加事件:存储并发布 */ public <E extends DomainEvent> void append(E event, long expectedVersion) { // 乐观并发控制 long currentVersion = eventRepository.getVersion(event.getAggregateId()); if (currentVersion != expectedVersion) { throw new ConcurrencyException(expectedVersion, currentVersion); } // 持久化 StoredEvent stored = StoredEvent.builder() .eventId(event.getEventId()) .aggregateId(event.getAggregateId()) .version(expectedVersion + 1) .eventType(event.getClass().getName()) .payload(serialize(event)) .occurredAt(event.getOccurredAt()) .build(); eventRepository.save(stored); // 发布到总线 eventBus.publish(event); } /** * 重放事件:重建聚合状态 */ public <A extends AggregateRoot> A replay(String aggregateId, Class<A> aggregateClass) { List<DomainEvent> events = eventRepository .findByAggregateId(aggregateId) .stream() .map(this::deserialize) .collect(Collectors.toList()); A aggregate = instantiateAggregate(aggregateClass, aggregateId); events.forEach(aggregate::apply); return aggregate; } }

五、观察者模式与相关模式的辨析

观察者 vs 中介者:观察者模式中发布者与订阅者直接通信(通过总线);中介者模式中所有通信通过中介者协调。观察者更松散,中介者更集中。

观察者 vs 发布-订阅:经典观察者中发布者维护订阅者列表;发布-订阅模式中引入事件通道,发布者与订阅者完全解耦。现代消息队列(Kafka、RabbitMQ)是发布-订阅的工业实现。

观察者 vs 责任链:观察者是一对多广播,责任链是一对一传递。观察者所有接收者并行处理,责任链只有一个处理者最终响应。

六、设计陷阱与规避策略

陷阱一:内存泄漏

观察者未注销导致主题持有过期引用。解决方案:使用WeakReference、自动注销机制(如Spring的DisposableBean),或基于存活时间的自动清理。

陷阱二:循环通知

A通知B,B又触发事件通知A,形成无限循环。解决方案:事件去重标记、循环检测、或严格的事件类型隔离。

陷阱三:观察者雪崩

某个观察者处理缓慢阻塞整体。解决方案:异步隔离、超时熔断、背压控制、独立线程池。

七、结语

观察者模式是事件驱动架构的理论基石,它将状态变化的通知机制从业务逻辑中彻底解耦,使系统具备极强的响应能力与扩展弹性。在实时风控、交易处理、物联网、微服务通信等场景中,观察者模式是不可或缺的组织原则。理解其同步/异步分层、优先级调度、背压控制、错误隔离等高级机制,掌握与响应式编程、事件溯源、CQRS等现代架构的融合,警惕内存泄漏、循环通知、雪崩效应等工程陷阱,是构建高可靠事件系统的能力核心。观察者模式的精髓在于承认世界的异步本质,以松耦合的方式响应变化,在不确定的环境中构建确定的响应能力

返回知识中心