You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 03:31:39