本文目录导读:

我将为您提供一个完整的Java发布订阅(Pub/Sub)模式案例,包含三种实现方式:手写实现、JDK内置实现和Spring实现。
手写发布订阅模式
基础接口定义
// 订阅者接口
public interface Subscriber {
void onMessage(String topic, Object message);
}
// 发布者接口
public interface Publisher {
void publish(String topic, Object message);
}
// 事件对象
public class Event {
private String topic;
private Object data;
private long timestamp;
public Event(String topic, Object data) {
this.topic = topic;
this.data = data;
this.timestamp = System.currentTimeMillis();
}
// getters and setters
public String getTopic() { return topic; }
public Object getData() { return data; }
public long getTimestamp() { return timestamp; }
}
事件总线实现
import java.util.*;
import java.util.concurrent.*;
public class EventBus {
// 使用ConcurrentHashMap存储主题和订阅者列表
private final ConcurrentMap<String, CopyOnWriteArrayList<Subscriber>> subscribers = new ConcurrentHashMap<>();
// 线程池处理异步消息
private final ExecutorService executor = Executors.newCachedThreadPool();
// 订阅
public void subscribe(String topic, Subscriber subscriber) {
subscribers.computeIfAbsent(topic, k -> new CopyOnWriteArrayList<>())
.add(subscriber);
System.out.println("订阅者订阅主题: " + topic);
}
// 取消订阅
public void unsubscribe(String topic, Subscriber subscriber) {
List<Subscriber> topicSubscribers = subscribers.get(topic);
if (topicSubscribers != null) {
topicSubscribers.remove(subscriber);
System.out.println("订阅者取消订阅主题: " + topic);
}
}
// 同步发布
public void publishSync(String topic, Object message) {
System.out.println("发布消息到主题 [" + topic + "]: " + message);
List<Subscriber> topicSubscribers = subscribers.get(topic);
if (topicSubscribers != null) {
Event event = new Event(topic, message);
topicSubscribers.forEach(subscriber ->
subscriber.onMessage(topic, event)
);
}
}
// 异步发布
public void publishAsync(String topic, Object message) {
System.out.println("异步发布消息到主题 [" + topic + "]: " + message);
List<Subscriber> topicSubscribers = subscribers.get(topic);
if (topicSubscribers != null) {
Event event = new Event(topic, message);
topicSubscribers.forEach(subscriber ->
executor.submit(() -> subscriber.onMessage(topic, event))
);
}
}
// 关闭线程池
public void shutdown() {
executor.shutdown();
}
}
具体业务实现
// 新闻订阅者
public class NewsSubscriber implements Subscriber {
private String name;
public NewsSubscriber(String name) {
this.name = name;
}
@Override
public void onMessage(String topic, Object message) {
Event event = (Event) message;
System.out.println("[" + name + "] 收到新闻: " + event.getData()
+ " (时间: " + new java.text.SimpleDateFormat("HH:mm:ss")
.format(new java.util.Date(event.getTimestamp())) + ")");
}
}
// 股票订阅者
public class StockSubscriber implements Subscriber {
private String name;
public StockSubscriber(String name) {
this.name = name;
}
@Override
public void onMessage(String topic, Object message) {
Event event = (Event) message;
System.out.println("[" + name + "] 收到股票消息: " + event.getData());
}
}
测试代码
public class PubSubTest {
public static void main(String[] args) throws InterruptedException {
EventBus eventBus = new EventBus();
// 创建订阅者
NewsSubscriber newsSub1 = new NewsSubscriber("News-订阅者1");
NewsSubscriber newsSub2 = new NewsSubscriber("News-订阅者2");
StockSubscriber stockSub1 = new StockSubscriber("Stock-订阅者1");
// 订阅
eventBus.subscribe("news", newsSub1);
eventBus.subscribe("news", newsSub2);
eventBus.subscribe("stock", stockSub1);
// 发布消息
System.out.println("=== 同步发布新闻 ===");
eventBus.publishSync("news", "Java 21 发布了新特性!");
System.out.println("\n=== 异步发布股票 ===");
eventBus.publishAsync("stock", "股票代码:600001 上涨5%");
// 异步等待
Thread.sleep(1000);
// 取消订阅
System.out.println("\n=== 取消订阅 ===");
eventBus.unsubscribe("news", newsSub2);
eventBus.publishSync("news", "第二条新闻消息");
// 关闭
eventBus.shutdown();
}
}
JDK内置的Observer实现
import java.util.Observable;
import java.util.Observer;
// 新闻发布者(继承Observable)
public class NewsPublisher extends Observable {
private String news;
public void publishNews(String news) {
this.news = news;
setChanged();
notifyObservers(news);
}
}
// 新闻订阅者(实现Observer)
public class NewsObserver implements Observer {
private String name;
public NewsObserver(String name) {
this.name = name;
}
@Override
public void update(Observable o, Object arg) {
System.out.println("[" + name + "] 收到新闻: " + arg);
}
}
// JDK实现测试
public class JdkPubSubTest {
public static void main(String[] args) {
NewsPublisher publisher = new NewsPublisher();
NewsObserver observer1 = new NewsObserver("观察者1");
NewsObserver observer2 = new NewsObserver("观察者2");
publisher.addObserver(observer1);
publisher.addObserver(observer2);
publisher.publishNews("JDK自带的观察者模式");
}
}
Spring事件驱动实现
定义事件
import org.springframework.context.ApplicationEvent;
// 自定义事件
public class MessageEvent extends ApplicationEvent {
private String topic;
private String message;
public MessageEvent(Object source, String topic, String message) {
super(source);
this.topic = topic;
this.message = message;
}
public String getTopic() { return topic; }
public String getMessage() { return message; }
}
事件监听器
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
@Component
public class MessageListener {
@EventListener
public void handleMessageEvent(MessageEvent event) {
System.out.println("监听器收到消息 - 主题: " + event.getTopic()
+ ", 内容: " + event.getMessage());
}
// 条件监听器
@EventListener(condition = "#event.topic == 'important'")
public void handleImportantMessage(MessageEvent event) {
System.out.println("重要消息处理: " + event.getMessage());
}
}
事件发布者
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.ApplicationEventPublisherAware;
import org.springframework.stereotype.Component;
@Component
public class MessagePublisher implements ApplicationEventPublisherAware {
private ApplicationEventPublisher publisher;
@Override
public void setApplicationEventPublisher(ApplicationEventPublisher publisher) {
this.publisher = publisher;
}
public void publishMessage(String topic, String message) {
MessageEvent event = new MessageEvent(this, topic, message);
publisher.publishEvent(event);
}
}
Spring配置和启动
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
@Configuration
@ComponentScan("com.example.pubsub")
public class SpringPubSubApplication {
public static void main(String[] args) {
AnnotationConfigApplicationContext context =
new AnnotationConfigApplicationContext(SpringPubSubApplication.class);
MessagePublisher publisher = context.getBean(MessagePublisher.class);
// 发布普通消息
publisher.publishMessage("normal", "普通消息示例");
// 发布重要消息
publisher.publishMessage("important", "重要消息示例");
context.close();
}
}
高级特性扩展
带过滤器的消息总线
public interface MessageFilter {
boolean accept(String topic, Object message);
}
public class FilteredEventBus extends EventBus {
private final Map<String, List<MessageFilter>> filters = new HashMap<>();
public void addFilter(String topic, MessageFilter filter) {
filters.computeIfAbsent(topic, k -> new CopyOnWriteArrayList<>())
.add(filter);
}
@Override
public void publishSync(String topic, Object message) {
List<MessageFilter> topicFilters = filters.get(topic);
if (topicFilters != null && !topicFilters.isEmpty()) {
boolean shouldPublish = topicFilters.stream()
.allMatch(filter -> filter.accept(topic, message));
if (shouldPublish) {
super.publishSync(topic, message);
}
} else {
super.publishSync(topic, message);
}
}
}
支持通配符订阅
public class WildcardEventBus extends EventBus {
private final Map<Subscriber, List<String>> subscriberPatterns = new ConcurrentHashMap<>();
public void subscribePattern(String pattern, Subscriber subscriber) {
subscriberPatterns.computeIfAbsent(subscriber, k -> new CopyOnWriteArrayList<>())
.add(pattern);
}
@Override
public void publishSync(String topic, Object message) {
super.publishSync(topic, message);
// 检查通配符订阅者
subscriberPatterns.forEach((subscriber, patterns) -> {
boolean match = patterns.stream().anyMatch(pattern ->
topic.matches(pattern.replace("*", ".*")));
if (match) {
subscriber.onMessage(topic, message);
}
});
}
}
完整测试演示
public class CompletePubSubDemo {
public static void main(String[] args) throws InterruptedException {
System.out.println("=== 发布订阅模式完整演示 ===\n");
// 使用带过滤功能的EventBus
FilteredEventBus eventBus = new FilteredEventBus();
// 创建订阅者
Subscriber logSubscriber = (topic, message) -> {
if (message instanceof Event) {
Event event = (Event) message;
System.out.println("[" + new java.text.SimpleDateFormat("HH:mm:ss")
.format(new java.util.Date(event.getTimestamp())) + "] "
+ topic + ": " + event.getData());
}
};
Subscriber criticalSubscriber = (topic, message) -> {
if (message instanceof Event) {
Event event = (Event) message;
System.out.println("🔴 紧急订阅者处理: [" + topic + "] " + event.getData());
}
};
// 订阅
eventBus.subscribe("user.login", logSubscriber);
eventBus.subscribe("user.important", criticalSubscriber);
// 添加过滤器
eventBus.addFilter("user.login", (topic, message) -> {
if (message instanceof Event) {
Event event = (Event) message;
return !String.valueOf(event.getData()).contains("BANNED");
}
return true;
});
// 发布消息
System.out.println("--- 发布普通消息 ---");
eventBus.publishSync("user.login", "用户 admin 登录成功");
System.out.println("\n--- 发布被过滤的消息 ---");
eventBus.publishSync("user.login", "用户 banned_user BANNED");
System.out.println("\n--- 发布紧急消息 ---");
eventBus.publishSync("user.important", "检测到异常访问!");
// 异步演示
System.out.println("\n--- 异步发布演示 ---");
eventBus.publishAsync("user.login", "用户 async_user 异步登录");
Thread.sleep(1000);
// 关闭
eventBus.shutdown();
}
}
Maven依赖(Spring版本)
<dependencies>
<!-- Spring Context -->
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-context</artifactId>
<version>5.3.30</version>
</dependency>
</dependencies>
这个案例展示了:
- 手写实现 - 完整理解发布订阅的原理
- JDK内置实现 - 使用Observer模式
- Spring实现 - 企业级事件驱动开发
- 高级特性 - 过滤器、通配符等扩展
您可以根据实际需求选择合适的实现方式,手写实现适合学习和轻量级场景,Spring实现适合企业级应用。