设计模式详解-观察者模式
设计模式详解:观察者模式
一、模式概述
观察者模式(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等现代架构的融合,警惕内存泄漏、循环通知、雪崩效应等工程陷阱,是构建高可靠事件系统的能力核心。观察者模式的精髓在于承认世界的异步本质,以松耦合的方式响应变化,在不确定的环境中构建确定的响应能力。