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

如何仅当有用户订阅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订阅/取消订阅事件,维护本地计数:

  1. 实现订阅计数器组件:
@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
    }
}
  1. 监听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)
            }
        }
    }
}
  1. 在定时任务中使用计数器:
@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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 22:54:18