在Java多线程程序中,线程间事件通知是指一个工作线程在状态改变或任务完成后,主动告知其他等待中的线程。观察者模式天然适合这种场景,但在并发环境中必须处理好可见性、原子性与死锁问题。

为什么需要线程安全的观察者模式
普通观察者模式在单线程下没有问题,但多个线程同时注册、注销或触发通知时,可能导致通知遗漏或抛出ConcurrentModificationException。并发场景下的核心诉求是:通知过程不阻塞发布线程过久,且所有订阅者都能收到一致事件。
基于并发容器的简单实现
我们可以用ConcurrentHashMap保存观察者,用队列解耦通知行为。
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
// 事件对象
class Event {
public final String msg;
public Event(String msg) {
this.msg = msg;
}
}
// 观察者接口
interface Observer {
void onEvent(Event e);
}
// 线程安全事件总线
class EventBus {
private final ConcurrentMap<String, Observer> observers = new ConcurrentHashMap<>();
private final ExecutorService pool = Executors.newCachedThreadPool();
public void register(String name, Observer o) {
observers.put(name, o);
}
public void unregister(String name) {
observers.remove(name);
}
// 异步通知,避免阻塞发布线程
public void publish(Event e) {
for (Observer o : observers.values()) {
pool.submit(() -> o.onEvent(e));
}
}
}
使用阻塞队列解耦
如果观察者处理逻辑很慢,可以改为让观察者自己从BlockingQueue中取事件,发布线程只负责入队。
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
public class QueueObserver implements Runnable {
private final BlockingQueue<Event> queue = new ArrayBlockingQueue<>(100);
public void accept(Event e) throws InterruptedException {
queue.put(e);
}
@Override
public void run() {
try {
while (true) {
Event e = queue.take();
System.out.println("处理事件:" + e.msg);
}
} catch (InterruptedException ex) {
Thread.currentThread().interrupt();
}
}
}
与内置API的对比
Java旧版提供了java.util.Observable和Observer,但它们的notifyObservers方法是同步的,且已被标记为废弃。在新代码中更推荐上述并发工具方案。
| 方案 | 线程安全 | 通知方式 |
|---|---|---|
| Observable | 否,需自行同步 | 同步调用 |
| ConcurrentHashMap+线程池 | 是 | 异步执行 |
| BlockingQueue | 是 | 生产消费解耦 |
实践建议
- 发布线程不要直接调用耗时观察者逻辑,优先异步化。
- 注册和注销方法必须使用线程安全容器。
- 事件对象应为不可变类,防止共享状态被修改。
观察者模式在并发中的价值在于解耦生产者和消费者,但只有配合正确的并发工具,才能真正实现安全的线程间事件通知。