如何仅当有用户订阅WebSocket主题时才执行指定逻辑?
Spring WebSocket/STOMP 按需处理并推送数据方案
方法一:直接利用Spring内置订阅注册表查询
Spring确实内置了订阅管理机制,你可以注入SubscriptionRegistry接口(默认实现为DefaultSubscriptionRegistry),直接检查目标目的地是否存在活跃订阅。
修改你的代码如下:
@Component class SendAccessLogMetrics( private val simpMessageSendingOperations: SimpMessageSendingOperations, private val accessLogStats: AccessLogStats, private val subscriptionRegistry: SubscriptionRegistry ) { @Scheduled(fixedDelay = 5, timeUnit = TimeUnit.SECONDS) fun sendAccessLogMetrics() { // 检查目标目的地是否有订阅者 val hasSubscribers = subscriptionRegistry.findSubscriptions( DestinationMatchers.patternEquals("/message/test") ).hasNext() if (hasSubscribers) { val chartData = createChartData(accessLogStats) simpMessageSendingOperations.convertAndSend("/message/test", chartData) } } }
Spring Boot环境下,DefaultSubscriptionRegistry会自动注册为Bean,直接注入即可使用。
方法二:监听订阅事件维护本地计数(性能优化)
如果定时任务频率较高,频繁查询订阅注册表可能存在性能损耗,可通过监听STOMP订阅/取消订阅事件,维护本地计数:
- 实现订阅计数器组件:
@Component class SubscriptionCounter { private val destinationSubscribers = mutableMapOf<String, AtomicInteger>() fun increment(destination: String) { destinationSubscribers.computeIfAbsent(destination) { AtomicInteger(0) }.incrementAndGet() } fun decrement(destination: String) { destinationSubscribers[destination]?.decrementAndGet() } fun hasSubscribers(destination: String): Boolean { return destinationSubscribers[destination]?.get() ?: 0 > 0 } }
- 监听STOMP会话事件:
@Component class StompEventListener(private val subscriptionCounter: SubscriptionCounter) : ApplicationListener<AbstractSubProtocolEvent> { override fun onApplicationEvent(event: AbstractSubProtocolEvent) { when (event) { is SessionSubscribeEvent -> { val destination = event.message.headers[StompHeaders.DESTINATION] as String subscriptionCounter.increment(destination) } is SessionUnsubscribeEvent -> { val destination = event.message.headers[StompHeaders.DESTINATION] as String subscriptionCounter.decrement(destination) } } } }
- 在定时任务中使用计数器:
@Component class SendAccessLogMetrics( private val simpMessageSendingOperations: SimpMessageSendingOperations, private val accessLogStats: AccessLogStats, private val subscriptionCounter: SubscriptionCounter ) { @Scheduled(fixedDelay = 5, timeUnit = TimeUnit.SECONDS) fun sendAccessLogMetrics() { if (subscriptionCounter.hasSubscribers("/message/test")) { val chartData = createChartData(accessLogStats) simpMessageSendingOperations.convertAndSend("/message/test", chartData) } } }
架构优化建议
- 若
createChartData逻辑计算量大,优先选择第二种方法,避免定时任务频繁查询注册表,减少开销。 - 多实例部署场景下,本地计数器会存在节点不一致问题,可改用分布式缓存(如Redis)维护订阅计数,或直接用第一种方法(每个节点独立查询本地订阅,推送自身节点数据)。
- 更高效的方式是事件驱动启停任务:第一个用户订阅时启动数据处理任务,最后一个用户取消订阅时停止任务,替代定时轮询。可通过
ScheduledExecutorService实现动态启停。
内容的提问来源于stack exchange,提问作者Andy
相关产品推荐
相关产品推荐

