推送服务做到一定规模,连接管理往往比业务逻辑本身更让人头疼。最朴素的写法是用一个全局的 ConcurrentHashMap 或者干脆一个普通 Map 加一把 synchronized 锁,连接建立时放进去,断开时移除,广播时遍历。连接数在几千以内时这套代码毫无问题,但一旦冲到五万、十万级别,这把全局锁就成了整个服务的单点瓶颈:每个连接的心跳、读写事件都要去抢同一把锁,CPU 大量时间花在锁竞争上,广播一条消息的耗时从毫秒级恶化到秒级。这篇文章就来拆解这个问题,并给出几种被生产环境验证过的替代方案。

为什么全局锁列表在高并发下必然失效
先看一段典型的反面代码,很多项目里的连接管理器就是这个骨架:
public class ConnectionManager {
private final Map<String, WebSocketSession> sessions = new HashMap<>();
private final Object lock = new Object();
public void add(String userId, WebSocketSession session) {
synchronized (lock) {
sessions.put(userId, session);
}
}
public void remove(String userId) {
synchronized (lock) {
sessions.remove(userId);
}
}
public void broadcast(String message) {
synchronized (lock) {
for (WebSocketSession s : sessions.values()) {
s.sendMessage(message); // 在持锁状态下做 IO!
}
}
}
}这段代码有两个致命问题。第一个是持锁做 IO:broadcast 在遍历时对每个 session 调用发送方法,网络 IO 的耗时(慢消费者可能阻塞几百毫秒)全部被算进临界区,其他线程的连接建立、断开操作全部被卡住。第二个是锁的粒度:五万个连接意味着所有事件都串行通过同一扇门,即使把 HashMap 换成 ConcurrentHashMap,broadcast 里的复合操作仍然需要外部一致性,锁依然逃不掉。
从操作系统层面看,锁竞争激烈时线程会频繁陷入内核做 futex 或者 park/unpark,上下文切换的开销会吞掉大量 CPU。你可以用 perf 或者 JFR 观察到,处于 RUNNABLE 却拿不到锁的线程占比异常高,服务吞吐反而随核数增加而下降——这是典型的锁反扩展性(lock convoy)现象。结论很明确:管理海量连接的第一原则,就是把“一把锁管所有连接”的结构彻底打散。
方案一:分片锁与分段哈希表
最直接的改进是分段(sharding)。把连接按 userId 的哈希值打散到 N 个独立的桶里,每个桶有自己的锁和 Map,不同桶之间的操作完全并行:
public class ShardedConnectionManager {
private static final int SHARD_COUNT = 256; // 建议为 2 的幂
private final Shard[] shards;
public ShardedConnectionManager() {
shards = new Shard[SHARD_COUNT];
for (int i = 0; i < SHARD_COUNT; i++) {
shards[i] = new Shard();
}
}
private static class Shard {
final Map<String, WebSocketSession> map = new HashMap<>();
final ReentrantLock lock = new ReentrantLock();
}
private int shardIndex(String key) {
int h = key.hashCode() ^ (key.hashCode() >>> 16);
return h & (SHARD_COUNT - 1);
}
public void add(String userId, WebSocketSession session) {
Shard s = shards[shardIndex(userId)];
s.lock.lock();
try { s.map.put(userId, session); } finally { s.lock.unlock(); }
}
public void sendTo(String userId, String message) {
Shard s = shards[shardIndex(userId)];
s.lock.lock();
try {
WebSocketSession session = s.map.get(userId);
if (session != null) session.sendMessage(message);
} finally { s.lock.unlock(); }
}
}分片之后,理论上的临界区并行度提升了 256 倍。实践中五万连接配 256 个分片,平均每个分片不到两百个连接,锁竞争概率已经极低。分片数量的选择有讲究:太少退化回全局锁,太多则浪费内存且缓存局部性变差,一般取 CPU 核数的 4 到 16 倍并向上取 2 的幂,方便用位运算取模。
这个方案的短板在广播。广播需要遍历所有分片,如果为了强一致而对所有分片加锁,就回到了老问题。正确的做法是接受最终一致:对每个分片短暂加锁只做“快照拷贝”或者干脆用 ConcurrentHashMap 作为桶内结构配合弱一致迭代器,遍历时允许少量新连接加入、旧连接移除,业务上几乎无感知。记住一个原则:广播场景下,消息发出去比消息严格有序更重要,为了一点一致性把 IO 拖进临界区是最常见的错误。
方案二:连接绑定线程,彻底消灭共享
分片锁只是降低了竞争概率,还有一种更彻底的思路:让每个连接只属于一个线程,读写、生命周期管理都在这个线程内完成,天然不存在共享,也就不需要任何锁。这是 Netty 的线程模型,也是 Go 的 goroutine-per-connection 模型的本质。
Netty 的做法是把连接注册到某个 EventLoop 上,之后该连接的所有 IO 事件都由这一个线程处理。如果连接管理也按 EventLoop 组织——每个 EventLoop 维护自己的本地连接 Map——那么增删查都变成单线程操作,零锁竞争。广播时由发起方往每个 EventLoop 的任务队列里投递一个广播任务,各线程并行发送:
EventLoopGroup group = new NioEventLoopGroup(Runtime.getRuntime().availableProcessors() * 2);
// 每个 EventLoop 持有私有的连接表,广播任务串行投递到各 loop
public void broadcast(ByteBuf msg) {
for (EventExecutor loop : group) {
loop.execute(() -> {
// 在 loop 自身线程内遍历本地连接并发送,无锁
localChannelsOf(loop).forEach(ch -> ch.writeAndRetain(msg.duplicate()));
});
}
}Go 的写法更直观,每个连接一个 goroutine,配合 channel 做消息分发。Go 的调度器本身就在用户态完成协程切换,百万级 goroutine 的内存开销(初始栈仅 2KB)完全可接受。需要注意的是写并发问题:多个 goroutine 同时往同一个连接写会造成帧交错,必须用带缓冲的发送 channel 把写操作收敛到连接所属的 goroutine 内,或者用 mutex 只保护单连接的写操作——注意这把锁是“每连接一把”而不是全局一把,粒度完全不同。
连接绑定线程方案的性能上限非常高,单机几十万连接是常规操作。它的代价是心智负担:任何跨线程访问连接状态的代码都可能引入死锁或竞态,务必通过 eventLoop().inEventLoop() 判断或 channel 通信来规范访问路径。
方案三:按业务维度组织连接,而不是一张大表
回到业务本身:绝大多数推送场景并不需要“给所有人发消息”,而是“给某个房间、某个群、某个用户的粉丝发消息”。如果按业务维度把连接组织成 Room 结构,每个 Room 内部是一个小连接集合,那么全局大表就退化成了多个小集合,锁的问题自然消解。
type Room struct {
mu sync.RWMutex
sessions map[string]*Session
}
type Hub struct {
rooms map[string]*Room
mu sync.RWMutex // 只保护 rooms 表本身
}
func (h *Hub) Join(roomID, userID string, s *Session) {
h.mu.RLock()
room, ok := h.rooms[roomID]
h.mu.RUnlock()
if !ok {
h.mu.Lock()
room, ok = h.rooms[roomID] // 双重检查
if !ok {
room = &Room{sessions: make(map[string]*Session)}
h.rooms[roomID] = room
}
h.mu.Unlock()
}
room.mu.Lock()
room.sessions[userID] = s
room.mu.Unlock()
}
func (r *Room) Broadcast(msg []byte) {
r.mu.RLock()
defer r.mu.RUnlock()
for _, s := range r.sessions {
select {
case s.sendCh <- msg: // 非阻塞投递,避免慢消费者拖垮整个房间
default:
// 队列满则断开或降级,绝不能阻塞广播协程
}
}
}注意上面 Broadcast 里的非阻塞投递,这是海量连接服务保命的设计:如果某个客户端网络差、发送 channel 堆满,直接阻塞会导致广播协程卡死,进而拖垮整个房间甚至整个进程。正确姿势是给每个连接设置有界队列,满则判定为慢消费者,踢掉或丢弃消息。同时要为每个用户考虑多端登录的场景——同一 userId 可能同时有 App 和 Web 两个连接,数据结构应该是 userId -> []*Session 的多值映射。
不可忽视的运维细节
数据结构选对了,还有几个高频踩坑点需要处理。第一是心跳与僵尸连接:五万连接里总有相当比例的客户端已经断网但 TCP 未感知,必须靠应用层心跳(比如 30 秒一次 ping,连续 3 次未收到 pong 判定死亡)来清理。清理动作要温和,先取消注册再关闭 fd,避免在关闭回调里又去抢锁。
第二是文件描述符与内存的系统性配置。单机五万连接至少要把 ulimit -n 调到十万以上,Linux 下还要关注 net.ipv4.ip_local_port_range(作为客户端压测时)和 TCP 内存参数。每个连接对象自身的内存也要精打细算:一个 Session 如果携带 4KB 缓冲区,五万连接光缓冲就是 200MB,尽量用池化(Netty 的 PooledByteBufAllocator)和按需分配。
第三是优雅下线与重连风暴。服务重启时如果瞬间断掉五万连接,客户端同时重连会形成雪崩。标准做法是网关层配合返回一个带抖动的重连退避时间,服务端分批关闭连接,并把连接状态外置到 Redis 一类的地方,让重连的客户端可以路由回原节点以恢复订阅关系。把这些细节和前面的三种数据结构方案组合起来,单机承载五万到五十万连接并不是什么黑魔法,关键就在于从架构上消灭全局锁这一个单点。
WebSocket连接管理连接池C10M并发修改时间:2026-09-12 03:22:46