如何为REST API实现观察者模式?Spring环境下订阅通知方案咨询
实现Spring REST API的观察者注册与实时数据推送方案
我来帮你梳理下在Spring生态里实现这种订阅-推送REST API的具体方案,结合实际项目经验,分模块给你讲清楚实现思路和最佳实践:
一、先搞定订阅者的注册与管理核心逻辑
首先需要一个线程安全的订阅管理器,用来维护「主题-订阅者」的映射关系,处理注册、取消订阅和批量通知的逻辑。这是观察者模式的核心载体:
@Component public class SubscriptionManager { // 用ConcurrentHashMap保证多线程下的安全操作,key是订阅主题,value是该主题下的订阅者集合 private final Map<String, Set<Subscriber>> subscribers = new ConcurrentHashMap<>(); // 注册订阅者到指定主题 public void subscribe(String topic, Subscriber subscriber) { subscribers.computeIfAbsent(topic, k -> ConcurrentHashMap.newKeySet()).add(subscriber); } // 从指定主题移除订阅者 public void unsubscribe(String topic, Subscriber subscriber) { Set<Subscriber> topicSubscribers = subscribers.get(topic); if (topicSubscribers != null) { topicSubscribers.remove(subscriber); // 主题下无订阅者时清理空集合,避免内存浪费 if (topicSubscribers.isEmpty()) { subscribers.remove(topic); } } } // 向指定主题的所有订阅者推送数据 public void notifySubscribers(String topic, Object data) { Set<Subscriber> topicSubscribers = subscribers.get(topic); if (topicSubscribers != null) { topicSubscribers.forEach(subscriber -> { try { subscriber.onDataUpdate(data); } catch (Exception e) { // 推送失败时自动取消该订阅者,避免后续无效推送 unsubscribe(topic, subscriber); } }); } } } // 定义订阅者接口,统一推送回调逻辑 public interface Subscriber { void onDataUpdate(Object data) throws IOException; }
这个管理器要注册为Spring Bean,这样可以在整个应用中共享,保证订阅数据的一致性。
二、选择适合的实时推送技术(REST场景下的两种主流方案)
传统REST是请求响应模式,要实现主动推送,需要用到以下两种技术,根据你的业务场景选:
1. Server-Sent Events (SSE):适合单向推送场景
如果你的客户端只需要接收数据推送,不需要主动给服务端发消息,SSE是最优解——它基于HTTP协议,不需要额外的端口或协议支持,客户端实现简单。
实现订阅接口
@RestController @RequestMapping("/api/subscribe") public class SseSubscriptionController { private final SubscriptionManager subscriptionManager; // 构造注入Spring Bean public SseSubscriptionController(SubscriptionManager subscriptionManager) { this.subscriptionManager = subscriptionManager; } @GetMapping(value = "/{topic}", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribeToTopic(@PathVariable String topic, HttpServletRequest request) { // 设置SSE连接超时时间(比如30分钟,避免长时间闲置占用资源) SseEmitter emitter = new SseEmitter(1800000L); // 创建当前连接对应的订阅者实例,绑定SseEmitter的推送逻辑 Subscriber subscriber = data -> { // 封装SSE事件,支持自定义事件类型、ID等 emitter.send(SseEmitter.event() .id(UUID.randomUUID().toString()) .data(data, MediaType.APPLICATION_JSON)); }; // 注册订阅者到管理器 subscriptionManager.subscribe(topic, subscriber); // 连接完成/超时/异常时自动取消订阅 emitter.onCompletion(() -> subscriptionManager.unsubscribe(topic, subscriber)); emitter.onTimeout(() -> { subscriptionManager.unsubscribe(topic, subscriber); emitter.complete(); }); emitter.onError(e -> { subscriptionManager.unsubscribe(topic, subscriber); emitter.completeWithError(e); }); return emitter; } }
触发数据推送
当外部事件导致数据变更时,通过事件监听或业务逻辑调用订阅管理器的通知方法即可:
@Component public class DataChangeNotifier { private final SubscriptionManager subscriptionManager; private final DataService dataService; public DataChangeNotifier(SubscriptionManager subscriptionManager, DataService dataService) { this.subscriptionManager = subscriptionManager; this.dataService = dataService; } // 假设这是外部事件触发的方法(比如数据库变更监听、MQ消息监听) @EventListener public void onDataUpdated(DataUpdateEvent event) { String topic = event.getTopic(); // 获取最新数据 Object latestData = dataService.getLatestDataByTopic(topic); // 推送给所有订阅该主题的客户端 subscriptionManager.notifySubscribers(topic, latestData); } }
2. WebSocket:适合双向通信场景
如果你的客户端需要主动和服务端交互(比如动态修改订阅主题、主动取消订阅),WebSocket是更好的选择,Spring提供了成熟的WebSocket+STOMP协议支持。
配置WebSocket
@Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Override public void configureMessageBroker(MessageBrokerRegistry config) { // 启用内置消息代理,处理订阅主题的消息转发 config.enableSimpleBroker("/topic"); // 客户端发送消息的前缀,对应@MessageMapping的路径 config.setApplicationDestinationPrefixes("/app"); } @Override public void registerStompEndpoints(StompEndpointRegistry registry) { // 注册WebSocket端点,客户端通过这个地址建立连接,SockJS用于兼容不支持WebSocket的浏览器 registry.addEndpoint("/ws-subscribe").withSockJS(); } }
处理订阅与推送
@Controller public class WebSocketSubscriptionController { private final SimpMessagingTemplate messagingTemplate; private final DataService dataService; public WebSocketSubscriptionController(SimpMessagingTemplate messagingTemplate, DataService dataService) { this.messagingTemplate = messagingTemplate; this.dataService = dataService; } // 处理客户端的订阅请求(可选,如果你需要自定义订阅逻辑) @MessageMapping("/subscribe/{topic}") public void handleSubscription(@DestinationVariable String topic, StompHeaderAccessor headerAccessor) { String sessionId = headerAccessor.getSessionId(); // 这里可以记录sessionId和topic的映射,用于后续精准推送或统计 } // 监听数据变更事件,推送消息到指定主题 @EventListener public void onDataUpdated(DataUpdateEvent event) { String topic = event.getTopic(); Object latestData = dataService.getLatestDataByTopic(topic); // 向指定主题推送消息 messagingTemplate.convertAndSend("/topic/" + topic, latestData); } }
三、关键最佳实践
- 线程安全优先:订阅管理器必须用线程安全的集合(比如
ConcurrentHashMap),避免多客户端并发订阅/取消时出现数据不一致。 - 连接生命周期管理:一定要处理SSE/WebSocket的超时、断开、异常场景,及时清理订阅者,避免内存泄漏。
- 消息可靠性保障:如果需要确保消息不丢失,可以结合消息队列(比如RabbitMQ),将数据变更事件先存入队列,再由服务端消费后推送;同时给客户端提供重连机制,重连后主动拉取最新数据。
- 主题合理划分:按数据类型、用户分组等维度拆分主题,避免向无关客户端推送数据,减少资源消耗。
- 限流与熔断:如果订阅者数量庞大,要考虑限流(比如限制单主题订阅数),并用Hystrix等组件实现熔断,避免服务被压垮。
- 序列化规范:推送的数据统一用JSON序列化,确保客户端能正常解析,同时避免序列化异常。
四、简单客户端示例(SSE)
用JavaScript实现的SSE客户端:
const targetTopic = "user-profile-updates"; const eventSource = new EventSource(`/api/subscribe/${targetTopic}`); // 接收推送数据 eventSource.onmessage = function(event) { const updatedData = JSON.parse(event.data); console.log("收到更新数据:", updatedData); // 这里写更新UI的逻辑 }; // 处理连接错误 eventSource.onerror = function(error) { console.error("SSE连接出错:", error); eventSource.close(); };
内容的提问来源于stack exchange,提问作者J-me
相关产品推荐
相关产品推荐

