观察者模式实战:从事件总线到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)
  • 消息推送系统
  • 系统广播与自定义广播
  • 插件系统的事件通知
  • 数据变更的多方同步

十、扩展方向

  1. 异步事件总线——使用线程池处理事件
  2. 事件优先级——监听器排序执行
  3. 事件持久化——未处理事件重启后补偿
  4. 分布式事件总线——跨进程/跨服务
  5. 集成MQ——Kafka/RabbitMQ作为事件传输层

结语

观察者模式是所有事件驱动架构的基石。从15行的事件总线到Spring的ApplicationEvent,从Android系统广播到框架的强制下线——实现方式各不相同,但注册→通知→处理的三步走从未改变。理解了这个三步走,任何发布-订阅系统本质上都是观察者模式。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐