观察者模式实战——从消息订阅看一对多通知
·
观察者模式实战:从事件总线到Android广播的完整落地
本文从零构建一个事件总线,进而对接 Spring ApplicationEvent、Android BroadcastReceiver,最后展示自研框架中"强制下线"的完整实现。所有代码可直接运行,核心思想可迁移至任何需要发布-订阅的业务系统。
文章目录
一、场景与目标
观察者模式解决的核心问题:一个对象的状态变化,需要通知多个接收方,且双方互不感知。
本文目标:从最简版事件总线 → Spring标准实现 → Android系统广播 → 框架实战,逐级展示观察者模式的不同落地形态。
二、角色定义
/**
* 观察者模式角色枚举
*/
public enum ObserverRole {
/** 被观察者 / 事件发布方 */
SUBJECT("主题"),
/** 观察者 / 事件接收方 */
OBSERVER("观察者"),
/** 事件对象 */
EVENT("事件"),
/** 事件总线 */
EVENT_BUS("事件总线");
private final String desc;
ObserverRole(String desc) { this.desc = desc; }
public String getDesc() { return desc; }
}
Publisher EventBus Listeners
┌──────────┐ post() ┌──────────┐ 逐个通知 ┌──────────────┐
│ 状态变化 │ ────────→│listeners │──────────→│ A: 发短信 │
│ 发布事件 │ │ [A,B,C] │──────────→│ B: 减库存 │
└──────────┘ └──────────┘ │ C: 写日志 │
└──────────────┘
一对多通知:Publisher不知道谁在监听,Listeners之间互不感知
三、从零构建事件总线
import java.util.*;
import java.util.concurrent.CopyOnWriteArrayList;
/**
* 自定义事件对象
*/
class LoginEvent {
private final String userId;
private final Date loginTime;
public LoginEvent(String userId) {
this.userId = userId;
this.loginTime = new Date();
}
public String getUserId() { return userId; }
public Date getLoginTime() { return loginTime; }
}
/**
* 监听器接口——观察者契约
*/
interface LoginListener {
void onLogin(LoginEvent event);
}
/**
* 事件总线——发布-订阅核心
*/
class EventBus {
// 使用线程安全的 CopyOnWriteArrayList
private static final List<LoginListener> listeners = new CopyOnWriteArrayList<>();
public static void register(LoginListener listener) {
Objects.requireNonNull(listener, "监听器不能为空");
listeners.add(listener);
}
public static void unregister(LoginListener listener) {
listeners.remove(listener);
}
public static void post(LoginEvent event) {
Objects.requireNonNull(event, "事件不能为空");
for (LoginListener listener : listeners) {
try {
listener.onLogin(event);
} catch (Exception e) {
System.err.println("监听器处理异常: " + e.getMessage());
}
}
}
public static int getListenerCount() {
return listeners.size();
}
}
四、Spring ApplicationEvent 标准实现
import org.springframework.context.ApplicationEvent;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
/**
* 自定义 Spring 事件
*/
public class OrderPaidEvent extends ApplicationEvent {
private final String orderId;
public OrderPaidEvent(Object source, String orderId) {
super(source);
this.orderId = orderId;
}
public String getOrderId() { return orderId; }
}
/**
* 监听器一:发送短信
*/
@Component
class SmsNotificationListener {
@EventListener
public void handleOrderPaid(OrderPaidEvent event) {
System.out.println("[SMS] 订单" + event.getOrderId() + "已支付,发送短信通知");
}
}
/**
* 监听器二:更新库存
*/
@Component
class InventoryUpdateListener {
@EventListener
public void handleOrderPaid(OrderPaidEvent event) {
System.out.println("[库存] 订单" + event.getOrderId() + "支付完成,扣减库存");
}
}
/**
* 监听器三:记录审计日志
*/
@Component
class AuditLogListener {
@EventListener
public void handleOrderPaid(OrderPaidEvent event) {
System.out.println("[审计] 订单" + event.getOrderId() + "支付记录已写入审计日志");
}
}
/**
* 支付服务——事件发布方
*/
@Component
class PaymentService {
private final ApplicationEventPublisher publisher;
public PaymentService(ApplicationEventPublisher publisher) {
this.publisher = publisher;
}
public void pay(String orderId) {
System.out.println("[支付] 订单" + orderId + "支付成功");
// 发布事件——不知道谁在监听
publisher.publishEvent(new OrderPaidEvent(this, orderId));
}
}
五、框架实战——强制下线广播
import android.content.BroadcastReceiver;
import android.content.Context;
import android.content.Intent;
import java.util.ArrayList;
import java.util.List;
/**
* 广播接收器——收到"强制下线"指令后清空所有Activity
*/
public class ForceLogoutReceiver extends BroadcastReceiver {
@Override
public void onReceive(Context context, Intent intent) {
// 关闭所有Activity
ActivityContainer.finishAll();
// 跳转登录页
Intent loginIntent = new Intent(context, LoginActivity.class);
loginIntent.addFlags(Intent.FLAG_ACTIVITY_NEW_TASK);
context.startActivity(loginIntent);
}
}
/**
* 全局 Activity 容器——观察者模式中的"监听器管理"
*/
class ActivityContainer {
private static final List<BaseActivity> activities = new ArrayList<>();
public static synchronized void addActivity(BaseActivity activity) {
activities.add(activity);
}
public static synchronized void removeActivity(BaseActivity activity) {
activities.remove(activity);
}
public static synchronized void finishAll() {
// 倒序遍历避免ConcurrentModification
for (int i = activities.size() - 1; i >= 0; i--) {
BaseActivity activity = activities.get(i);
if (activity != null && !activity.isFinishing()) {
activity.finish();
}
}
activities.clear();
}
}
/**
* 基础Activity——自动注册到全局容器
*/
class BaseActivity {
private boolean finishing = false;
protected void onCreate() {
ActivityContainer.addActivity(this);
}
protected void onDestroy() {
ActivityContainer.removeActivity(this);
}
public boolean isFinishing() { return finishing; }
public void finish() { finishing = true; }
}
六、测试用例
import java.util.concurrent.atomic.AtomicInteger;
public class ObserverTest {
public static void main(String[] args) {
System.out.println("=== 自定义事件总线测试 ===");
// 注册监听器
AtomicInteger smsCount = new AtomicInteger(0);
AtomicInteger auditCount = new AtomicInteger(0);
EventBus.register(event -> {
smsCount.incrementAndGet();
System.out.println("[SMS] 用户 " + event.getUserId() + " 登录");
});
EventBus.register(event -> {
auditCount.incrementAndGet();
System.out.println("[审计] 记录登录: " + event.getUserId());
});
// 发布事件
EventBus.post(new LoginEvent("user001"));
// 验证
System.out.println("\n=== 结果验证 ===");
System.out.println("监听器数量: " + EventBus.getListenerCount());
System.out.println("短信通知数: " + smsCount.get()); // 1
System.out.println("审计日志数: " + auditCount.get()); // 1
// 注销并验证
EventBus.unregister(EventBus.listeners.get(0));
System.out.println("注销后监听器数量: " + EventBus.getListenerCount()); // 1
}
}
七、观察者 vs 责任链
| 维度 | 观察者模式 | 责任链模式 |
|---|---|---|
| 通知方式 | 广播,所有观察者都收到 | 链式传递,节点决定是否继续 |
| 处理方式 | 各自独立处理 | 按顺序处理 |
| 终止条件 | 全部处理完毕 | 任一节点可截断 |
| 调用方感知 | 不知道谁在监听 | 不知道链有多长 |
| 典型场景 | 消息订阅、事件通知 | 权限校验、过滤器链 |
八、亮点总结
✅ 三种形态逐级递进:自研总线 → Spring标准 → Android广播
✅ 线程安全的 CopyOnWriteArrayList 实现
✅ 异常隔离——一个监听器失败不影响其他监听器
✅ Activity 容器 + BroadcastReceiver 的完整强制下线方案
✅ 观察者与责任链的对比,清晰区分两种模式
✅ 完整可运行测试代码
九、适用场景
- 事件驱动的微服务架构
- UI 组件间通信(MVP/MVVM)
- 消息推送系统
- 系统广播与自定义广播
- 插件系统的事件通知
- 数据变更的多方同步
十、扩展方向
- 异步事件总线——使用线程池处理事件
- 事件优先级——监听器排序执行
- 事件持久化——未处理事件重启后补偿
- 分布式事件总线——跨进程/跨服务
- 集成MQ——Kafka/RabbitMQ作为事件传输层
结语
观察者模式是所有事件驱动架构的基石。从15行的事件总线到Spring的ApplicationEvent,从Android系统广播到框架的强制下线——实现方式各不相同,但注册→通知→处理的三步走从未改变。理解了这个三步走,任何发布-订阅系统本质上都是观察者模式。
更多推荐
所有评论(0)