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

Spring WebSocket聊天系统DefaultSubscriptionRegistry高并发线程阻塞问题咨询

解决Spring WebSocket STOMP中DefaultSubscriptionRegistry导致的高并发线程阻塞问题

我之前也碰到过一模一样的场景——用Spring WebSocket+STOMP开发的聊天系统,高并发下大量线程阻塞在DefaultSubscriptionRegistry$DestinationCache$1的getSubscriptions方法,最多的时候线程数也破过400。这个问题的核心原因和解决方案我梳理得很清楚,分享给你:

问题根源分析

DefaultSubscriptionRegistry是Spring默认用来管理STOMP订阅关系的组件,它内部的DestinationCache使用了全局synchronized锁来保护缓存的读写操作。在高并发场景下(比如大量用户同时订阅/取消订阅同一个主题),所有请求都会争抢这把锁,导致线程排队阻塞,甚至引发线程池耗尽的问题。另外,当出现连接异常断开、批量订阅请求涌入这类异常场景时,锁的持有时间会被拉长,进一步加剧阻塞情况。

可行解决方案

1. 替换为RedisSubscriptionRegistry(分布式场景首选)

如果你的系统是分布式部署的,直接改用RedisSubscriptionRegistry是最省心的方案。它把订阅信息存储在Redis中,避免了单节点的锁竞争,同时天然支持多节点的订阅同步。

配置示例:

@Configuration
@EnableWebSocketMessageBroker
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {

    @Autowired
    private RedisConnectionFactory redisConnectionFactory;

    @Override
    public void configureMessageBroker(MessageBrokerRegistry config) {
        config.enableStompBrokerRelay("/topic", "/queue")
              .setRelayHost("localhost")
              .setRelayPort(61613);
        config.setApplicationDestinationPrefixes("/app");
    }

    @Bean
    public SubscriptionRegistry subscriptionRegistry() {
        return new RedisSubscriptionRegistry(redisConnectionFactory);
    }
}

2. 自定义细粒度锁的SubscriptionRegistry(单节点场景)

如果是单节点部署,不想依赖Redis,可以自己实现一个基于细粒度锁的SubscriptionRegistry,替换默认的实现。核心思路是用ConcurrentHashMap存储缓存项,每个缓存项对应一把独立的锁,避免全局锁的竞争。

示例代码:

public class FineGrainedSubscriptionRegistry extends DefaultSubscriptionRegistry {

    @Override
    protected DestinationCache createDestinationCache() {
        return new FineGrainedDestinationCache();
    }

    private class FineGrainedDestinationCache extends DestinationCache {
        private final ConcurrentHashMap<String, ReentrantLock> lockMap = new ConcurrentHashMap<>();

        @Override
        public Set<Subscription> getSubscriptions(String destination) {
            // 为每个目的地获取独立的锁
            ReentrantLock lock = lockMap.computeIfAbsent(destination, k -> new ReentrantLock());
            lock.lock();
            try {
                return super.getSubscriptions(destination);
            } finally {
                lock.unlock();
            }
        }

        // 重写其他修改缓存的方法,同样使用细粒度锁
        @Override
        public void addSubscription(Subscription subscription) {
            String destination = subscription.getDestination();
            ReentrantLock lock = lockMap.computeIfAbsent(destination, k -> new ReentrantLock());
            lock.lock();
            try {
                super.addSubscription(subscription);
            } finally {
                lock.unlock();
            }
        }

        @Override
        public void removeSubscription(Subscription subscription) {
            String destination = subscription.getDestination();
            ReentrantLock lock = lockMap.computeIfAbsent(destination, k -> new ReentrantLock());
            lock.lock();
            try {
                super.removeSubscription(subscription);
                // 如果该目的地没有订阅了,移除锁节省内存
                if (super.getSubscriptions(destination).isEmpty()) {
                    lockMap.remove(destination);
                }
            } finally {
                lock.unlock();
            }
        }
    }
}

然后在配置类中注册这个自定义Bean:

@Bean
public SubscriptionRegistry subscriptionRegistry() {
    return new FineGrainedSubscriptionRegistry();
}

3. 优化订阅目的地设计

从业务层面减少同一目的地的订阅量,比如:

  • 按用户分组拆分主题,比如把/topic/group拆分为/topic/group/{groupId},让不同组的订阅请求分散到不同的缓存项,降低锁竞争。
  • 对私有聊天使用用户专属的目的地,比如/queue/user/{userId},避免大量用户抢占同一个主题的锁。

4. 调整线程池参数(辅助优化)

适当调大Tomcat的NIO线程池(比如server.tomcat.max-threads)和Spring WebSocket的消息处理线程池,避免因为线程阻塞导致请求无法及时处理,但这只是缓解手段,核心还是要解决锁竞争问题。


内容的提问来源于stack exchange,提问作者tkec

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:35:33