实时聊天的高吞吐设计:WebSocket + Disruptor 事件环实战拆解

悦目图库

聊天功能上线后的第一个高峰夜,运维的告警比用户的吐槽来得还快:WebSocket 服务的 GC 时间飙升,消息延迟从几十毫秒涨到两秒多,个别用户的连接直接被挤断。日志里能看到同一秒钟里,好几个线程在往同一个 WebSocketSession 上写消息。

这就是最常见的聊天架构翻车现场:每个连接直接发。收到一条消息,立刻拿着 session 广播出去,看着没毛病,但一旦并发上来,问题一串串地冒出来——多个线程同时写同一个 session 会乱序、会抛异常、会把 TCP 缓冲写炸;对象满天飞把年轻代塞满,GC 频繁到服务卡顿。聊天这种 IO 密集加高并发的场景,恰恰是这些问题最集中的地方。

这篇文章拆的是这个项目真实在跑的一套方案:WebSocket 接住连接,Disruptor 消化并发。前端怎么心跳、怎么断线重连、怎么保证消息不丢,后端怎么用无锁环形缓冲把“并发广播”变成“单线程消费”,弹幕又是怎么在 canvas 上跑满 60 帧的。代码都贴真实原文,能讲透的地方绝不糊弄。这也不是一篇讲理论的科普——每一个组件都是真实部署、真实扛过流量之后再回头的复盘,包括那些事后才看明白的选型理由。你会发现,很多"为什么"不是一开始就想清楚的,而是被线上问题逼着改出来的。

一条消息的旅程

在贴任何代码之前,先记住这条链路,后面每一节都是它的一个环节:

  1. 前端拿到用户输入,把消息 JSON 塞进 WebSocket,发到 /api/ws/chat
  2. 后端 ChatWebSocketServer.handleTextMessage 解析消息、校验字段、确定目标(私聊/图片聊天室/空间聊天);
  3. 不直接处理,而是把消息装进一个事件,扔给 Disruptor 的 RingBuffer
  4. Disruptor 的 WorkerPool 里一个工作线程把事件取出来,调 handleXxxChatMessage:落库、填充发送者信息、把消息广播给目标房间里所有还开着的 session;
  5. 前端 onmessage 收到广播,把消息插进列表。

关键在第三步和第四步之间:生产者只负责“发布”,真正的“处理+广播”被串成了单线程。为什么非要绕这一圈?因为 WebSocketSession.sendMessage 不是线程安全的——同一个 session 被多个线程并发写,轻则消息乱序,重则直接抛 IllegalStateException 把连接搞断。Disruptor 的价值,就是把“谁都能往里塞消息”的并发入口,收敛成一个“一次只处理一条”的串行出口。

一句话记住这个架构:入口并发,出口串行

这里说的“高吞吐”,不是追求每秒能发多少条,而是追求在消息尖峰到来时不崩、不乱、不丢。聊天室的流量是突发性的——一场活动、一个热点,能瞬间把消息量顶上十倍。真正考验架构的,是那个十倍的时刻,而不是平时的一倍。Disruptor 的预分配、单线程消费、背压,全都是在为那个时刻兜底。

WebSocket + Disruptor Architecture

聊天室界面 ▲ 图:用户眼里的聊天室——从输入框到一串气泡,架构里的一切最终都落在这里

为什么偏偏是 Disruptor

如果你没接触过 Disruptor,它其实是 LMAX 交易平台开源的一个高性能队列,靠“无锁环形缓冲”出名。要理解它为什么适合聊天,得先回答一个问题:在 Java 里做“生产者-消费者”,不是有 BlockingQueue 吗?

有,但 BlockingQueue 在高并发下有两个硬伤:一是它底层靠锁(或者 ConcurrentLinkedQueue 的 CAS),竞争激烈时锁的切换开销不小;二是它不能复用对象——每个事件都是新 new 出来的,消息量一大,年轻代 GC 就跟着遭殃。聊天场景偏偏是“高并发 + 海量小对象”,两个硬伤全踩中。

Disruptor 换个思路:预先分配一个固定大小的环形数组(这里 bufferSize = 1024 * 256,也就是 26 万多个槽位),所有事件从一开始就存在于这个数组里,生产者要发布事件,不是 new 一个对象,而是从数组里取一个槽位,往里面填数据。处理完,把槽位清空,下一个循环再用它。整个生命周期里,事件对象一个都没新建过——GC 压力被直接摁死。

这是 Disruptor 的第一个关键设计:对象复用靠预分配 + 清空,而不是靠垃圾回收。

第二个关键设计是无锁。RingBuffer 维护两个序列号:生产者序列和消费者序列。生产者要拿下一个可用槽位,调 ringBuffer.next(),这是一个 CAS 操作,返回一个序列号;拿到后往里填数据,最后调 ringBuffer.publish(sequence) 把数据“发布”出去,消费者才能看到。没有锁,只有 CAS 和一个简单的序列比较。

顺便算一笔账。bufferSize = 1024 * 256 是 262144 个槽位,每个 ChatEvent 对象里装着一条消息、一个 session、一个用户和几个 Long,粗算一个对象一两百字节,整个环形缓冲预留的内存也就几十 MB 级别。这个内存是一次性的,启动时就分配好了,之后无论消息多密集,都不再增长。这跟 BlockingQueue 的区别就很直观:BlockingQueue 是"用多少、new 多少",GC 跟着跑;Disruptor 是"一次性买断",后面全是零成本复用。为什么 bufferSize 必须是 2 的幂?因为环形数组定位槽位用位运算 index = sequence & (bufferSize - 1),长度是 2 的幂,才能用掩码替代取模。

配置代码长这样,短短几行,但信息量很大:

int bufferSize = 1024 * 256;
Disruptor<ChatEvent> disruptor = new Disruptor<>(
        ChatEvent::new,          // 事件工厂:预分配时用来创建对象
        bufferSize,              // 环形缓冲大小,必须是 2 的幂
        ThreadFactoryBuilder.create().setNamePrefix("chatEventDisruptor").build()
);
disruptor.handleEventsWithWorkerPool(chatEventWorkHandler); // 工作线程池消费
disruptor.start();

handleEventsWithWorkerPool 有个细节值得单独说:它开的是一个工作线程池,池里的线程彼此分工——每个事件只会被其中一个线程处理,不会重复消费。这里池子实际上只有一个工作线程,所以消息被严格串行处理。这就是“并发入口、串行出口”的落点:无论多少连接同时发消息,真正执行“落库 + 广播”的,同一时刻只有这一个线程,session 永远不会被并发写。

生产者:抢槽位,填数据,发布

生产者 ChatEventProducer.publishEvent 是 Disruptor 用法的标准范式,五步走完:

public void publishEvent(ChatMessage chatMessage, WebSocketSession session,
        User user, Long targetId, Integer targetType) {
    RingBuffer<ChatEvent> ringBuffer = chatEventDisruptor.getRingBuffer();
    long sequence = ringBuffer.next();       // 1. 抢占一个槽位(CAS)
    try {
        ChatEvent event = ringBuffer.get(sequence); // 2. 拿到槽位里的对象
        event.setChatMessage(chatMessage);    // 3. 填数据
        event.setSession(session);
        event.setUser(user);
        event.setTargetId(targetId);
        event.setTargetType(targetType);
    } finally {
        ringBuffer.publish(sequence);        // 4. 发布,消费者可见
    }
}

注意两个细节。第一,ringBuffer.next() 是可能阻塞的——如果环形缓冲满了(消费者处理不过来),生产者会在这里等着。这意味着 Disruptor 天然带背压:消息堆积时,不是无限往里塞导致内存爆掉,而是让生产者(也就是 WebSocket 消息处理线程)停下来等。对一个聊天系统来说,这是很合理的取舍:宁可让消息处理慢一点,也不能把内存打爆。

第二,finallypublish。这是 Disruptor 的铁律:next()publish() 必须成对出现,中间哪怕抛异常,也要保证 publish 被调用,否则槽位被占着不放,环形缓冲会越跑越小,最终把整个事件环卡死。这个坑我踩过一次,后面踩坑那一节会细讲。

事件对象本身平平无奇,就是一组字段加上一个 clear()

public class ChatEvent {
    private ChatMessage chatMessage;
    private WebSocketSession session;
    private User user;
    private Long targetId;
    private Integer targetType; // 1-私聊 2-图片聊天室 3-空间聊天

    public void clear() {
        this.chatMessage = null;
        this.session = null;
        this.user = null;
        this.targetId = null;
        this.targetType = null;
    }
}

clear() 是配合“对象复用”的关键:事件处理完,把字段清空,下次这个槽位被复用的时候,就不会残留上一个事件的数据。忘写 clear,或者 clear 漏字段,轻则数据污染,重则把上一个用户的消息发到下一个用户的房间里——隐私事故。

串行,全在消费者这里

消费者 ChatEventWorkHandler 实现了 Disruptor 的 WorkHandler 接口,是整条链路里真正干活的地方:

public class ChatEventWorkHandler implements WorkHandler<ChatEvent> {

    @Resource
    @Lazy
    private ChatWebSocketServer chatWebSocketServer;

    @Override
    public void onEvent(ChatEvent event) {
        try {
            ChatMessage chatMessage = event.getChatMessage();
            switch (event.getTargetType()) {
                case 1: // 私聊
                    chatMessage.setPrivateChatId(event.getTargetId());
                    chatWebSocketServer.handlePrivateChatMessage(chatMessage, event.getSession());
                    break;
                case 2: // 图片聊天室
                    chatMessage.setPictureId(event.getTargetId());
                    chatWebSocketServer.handlePictureChatMessage(chatMessage, event.getSession());
                    break;
                case 3: // 空间聊天
                    chatMessage.setSpaceId(event.getTargetId());
                    chatWebSocketServer.handleSpaceChatMessage(chatMessage, event.getSession());
                    break;
                default:
                    log.error("Unknown target type: {}", event.getTargetType());
            }
        } catch (Exception e) {
            log.error("处理聊天消息失败", e);
        } finally {
            event.clear(); // 清空事件数据,供环形缓冲复用
        }
    }
}

@Lazy 注解值得注意。ChatWebSocketServerChatEventWorkHandler 互相依赖,如果不用 @Lazy 懒加载,Spring 在创建 Bean 的时候会直接循环依赖报错。加 @Lazy 的意思是:别急着在启动时把对方注入进来,等到真正要用的时候再说。这是 Spring 处理循环依赖的一个常用手段,代价是第一次调用会慢一点(要现场解析),但对聊天这种低频启动的场景无所谓。

真正的广播逻辑在 handleXxxChatMessage 里。以空间聊天为例:

private void sendToSpaceRoom(ChatMessage chatMessage) throws IOException {
    Set<WebSocketSession> sessions = spaceSessions.get(chatMessage.getSpaceId());
    if (sessions != null) {
        String messageStr = webSocketObjectMapper.writeValueAsString(chatMessage);
        for (WebSocketSession session : sessions) {
            if (session.isOpen()) {
                session.sendMessage(new TextMessage(messageStr));
            }
        }
    }
}

流程是:保存消息 → fillMessageInfo 填充发送者头像昵称等 → 把整条消息序列化成 JSON → 遍历这个房间的 session 集合逐个发送。因为整个 onEvent 跑在同一个 Disruptor 工作线程里,所以这些 sendMessage 天然是串行的,不会有两个线程同时写同一个 session。

有个地方很容易被忽略:为什么广播之前要落库? 因为聊天记录必须可回溯——刷新页面、重新进房间,都要把历史消息拉出来。所以每条消息都是“先存库,再广播”,顺序不能反:先落库失败,就不该广播出去让人看到一条数据库里不存在、刷新就消失的消息。

四张表,管住所有连接

聊天不止一个房间,这个系统里同时存在四种聊天场景:私聊、图片聊天室、空间聊天,外加一个“用户全局在线”的概念。对应的,后端维护了四张 ConcurrentHashMap

private static final Map<Long, WebSocketSession> userSessions = new ConcurrentHashMap<>();
private static final Map<Long, Set<WebSocketSession>> pictureSessions = new ConcurrentHashMap<>();
private static final Map<Long, Set<WebSocketSession>> spaceSessions = new ConcurrentHashMap<>();
private static final Map<Long, Set<WebSocketSession>> privateChatSessions = new ConcurrentHashMap<>();

userSessions 是“用户 ID → 他的 session”,用来判断一个用户整体在线;后面三张是“房间 ID → 房间里所有 session 的集合”,广播就是遍历这三张表里对应的集合。集合用的是 ConcurrentHashMap.newKeySet(),是并发安全的,加 session、删 session 都不会炸。

连接建立的时候,afterConnectionEstablished 根据握手时带进来的参数决定把 session 塞进哪张表,并且顺手做几件事:给私聊用户发历史消息、给房间广播一次“在线用户列表”。连接关闭的时候,afterConnectionClosed 把 session 从对应集合里移除,如果集合空了就把整个房间的条目删掉,再广播一次在线用户更新——这样别人那里显示的在线人数才不会一直挂着已离开的人。

这些表都是 static 的,也就是挂在 JVM 上的。这就引出一个聊天架构的经典话题:多实例部署时怎么办。这个项目目前是单实例,所以静态表没问题;一旦要水平扩展,就得把 session 注册表搬进 Redis 之类的共享存储,或者用 STOMP + 消息代理。这是一套方案的边界,得心里有数,别等上了多实例才发现静态表不通。

在 session 进入这些表之前,还有一道关:握手。WsHandshakeInterceptor 是连接建立前的最后一道闸,它干三件事。第一,登录校验——userService.isLogin(httpRequest) 拿到当前登录用户,拿不到就直接拒绝握手,返回 false,连接根本不建立。注意这里用的是 isLogin 而不是 getLoginUser,因为后者在未登录时会抛异常,而握手流程里抛异常会把这次握手搞成 500;isLogin 返回 null 就能优雅地拒绝。

第二,解析房间参数。握手 URL 上带着 pictureIdspaceIdprivateChatId,拦截器把它们解析成 Long 塞进 session 的 attributes 里,后面 afterConnectionEstablished 就靠这些属性决定把 session 放进哪张表。参数格式错了(比如传了非数字),直接拒绝握手——这是第一道防线,防止有人乱塞参数。

第三,解析真实 IP。getClientIpAddress 走了一条代理头链:X-Forwarded-ForProxy-Client-IPWL-Proxy-Client-IPHTTP_CLIENT_IPHTTP_X_FORWARDED_FOR → 最后才回落到 getRemoteAddr。因为服务器通常挂在 Nginx 后面,getRemoteAddr 拿到的永远是 Nginx 的 IP,必须从代理头里抠出真实客户端 IP。还有个细节:IPv6 的回环地址 0:0:0:0:0:0:0:1 会被规范成 127.0.0.1,不然日志和 IP 定位都会踩坑。这个 IP 存进 attributes,后面记录登录地、做风控都用它。

顺便说一句,并不是所有 WebSocket 都走这个拦截器——批量上传进度和 AI 令牌用量这两个端点就绕过了它,因为它们的握手只需要 URL 上的 userId 参数,不需要登录态校验。这说明握手策略要按端点灵活设计,不是一刀切。

13 秒心跳,10 分钟判死

WebSocket 连接看起来是长连接,但 TCP 层的长连接不保证“活着”——手机切网、代理掐线、路由清理,都可能让连接变成“半死”状态:两端都不知道对方已经没了。所以必须自己用心跳探测。

Heartbeat & Timeout

后端用了一个后台线程做心跳检查:

private static final long ACTIVITY_TIMEOUT = 10 * 60 * 1000; // 10分钟超时
private static final long HEARTBEAT_INTERVAL = 13 * 1000;    // 13秒心跳间隔

private void checkHeartbeats() {
    while (true) {
        try {
            Thread.sleep(HEARTBEAT_INTERVAL);
            long now = System.currentTimeMillis();
            lastActivityTime.entrySet().removeIf(entry -> {
                WebSocketSession session = entry.getKey();
                long inactiveTime = now - entry.getValue();
                if (inactiveTime > ACTIVITY_TIMEOUT) {
                    session.close();      // 超时,关掉
                    return true;
                }
                return false;
            });
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            break;
        }
    }
}

逻辑朴素但有效:每个 session 记录“最后活跃时间”,后台线程每 13 秒扫一遍,谁超过 10 分钟没动静就关掉。前端发来的任何消息(包括心跳)都会刷新这个时间,所以只要前端还活着,这个 session 就不会被判死。

微妙之处在于:为什么是 10 分钟而不是更短?因为后端把“心跳响应”和“活跃”做了区分——真正断开连接的判断,是前端在收不到心跳响应后主动发起重连(下一节讲),后端的 10 分钟超时只是一个兜底,防止那些“前端已经死了但连接没断”的僵尸连接一直占着资源。两个阈值各管各的:前端 24 秒判活,后端 10 分钟兜底,谁也不会误杀谁。

心跳线程还有个容易被忽略的价值:它顺带做了可观测性。扫描的时候,lastActivityTime 会先把非活跃超过 5 分钟(超时阈值的一半)的 session 单独记一条日志——"用户 xx 已无活动 x 分钟",方便运维提前发现异常连接;真到 10 分钟才真正关闭。这就像心跳除了保活,还顺便给了一张"谁在挂机"的体检表,排查问题的时候不用再临时加日志。主动观测这一步,很多系统都省了,但真出问题的时候,它就是救命的。

断了?退避着重连

前端的 ChatWebSocket 类把心跳和重连做成了一套完整的机制,值得拆开看。先看心跳:

this.heartbeatTimer = window.setInterval(() => {
  if (this.destroyed) { clearInterval(this.heartbeatTimer); return }

  if (!this.socket || this.socket.readyState !== WebSocket.OPEN) return

  // 超过心跳间隔两倍没收到响应,说明连接可能已经死了
  if (this.lastHeartbeatResponse > 0 &&
      Date.now() - this.lastHeartbeatResponse > this.heartbeatInterval * 2 &&
      !this.connecting) {
    this.reconnect()
    return
  }

  this.socket.send(JSON.stringify({ type: 'HEARTBEAT', time }))
}, this.heartbeatInterval)

前端每 24 秒发一次心跳,lastHeartbeatResponse 记录最后一次收到心跳响应的时间。如果 48 秒(两倍间隔)没收到响应,就判定连接已死,主动触发重连。这个“两倍间隔”的判断很关键——单次丢包不该触发重连,只有连续一个周期都没响应才算真断了,避免在网络抖动时疯狂重连。

重连用的是指数退避:

if (this.reconnectAttempts < this.maxReconnectAttempts) {
  setTimeout(() => {
    this.reconnectAttempts++
    this.reconnectTimeout *= 2 // 指数退避:1s → 2s → 4s → 8s → 16s
    this.connect()
  }, this.reconnectTimeout)
}

最多重试 5 次,等待时间从 1 秒开始每次翻倍。指数退避的意义在于:网络恢复通常需要时间,如果每次失败都立刻重试,只会把服务端的握手请求打爆,而且重试本身也很可能失败。先等 1 秒,不行等 2 秒,再不行等 4 秒——给它一个恢复的时间窗口,也给自己一个喘息的空间。

消息丢不了:队列与乐观回显

心跳只解决“连接活着”的问题,还有一个更棘手的问题:如果连接恰好断开,用户发出去的消息怎么办? 直接丢掉是最糟的——用户明明点了发送,消息却消失,这是聊天产品最忌讳的事。

Message Queue & Optimistic UI

前端的解法是维护一个消息队列,发送前先入队,连上了才真正发出去:

// 将消息添加到队列
this.messageQueue.push({ message: processedMessage, retryCount: 0, timestamp: Date.now() })
this.processMessageQueue()

processMessageQueue 是一个带重试的循环:连接没就绪就等重连,重试 3 次(每次间隔 2 秒)还发不出去才丢弃;连接就绪就立刻发送。更妙的是——乐观回显

// 发送成功,立即触发本地消息事件,不等待服务器响应
if (item.message.type === 1 && item.message.content) {
  this.triggerEvent('message', {
    type: 'message',
    message: {
      ...item.message,
      id: Date.now().toString(),  // 临时ID
      sender: useLoginUserStore().loginUser
    }
  })
}

也就是说,用户点发送,消息立刻出现在自己的对话列表里,用的是本地临时 ID;服务器广播回来后,再用服务器的真实消息替换。这个体验优化非常关键——如果等服务器确认才显示,网络稍慢一点,用户就会觉得“卡了”“没发出去”,而乐观回显让发送成功这件事零延迟反馈给用户。代价是临时 ID 和服务器 ID 的切换,以及极端情况下(消息最终发送失败)要把乐观回显的那条删掉。

这套“队列 + 重试 + 乐观回显”组合起来,就是在尽力保证:只要用户发了,消息就大概率会到;即使网络波动,也只是晚到,不是丢失。

再往前端消息的接收端看一眼。ChatWebSocket 维护一张事件订阅表 eventHandlerson(type, handler) 注册监听,triggerEvent(type, data) 把消息分发给所有订阅者。onmessage 里先做一件事:判断是不是心跳响应,是就更新 lastHeartbeatResponse 并直接返回,不再往下分发;否则把消息抛给 message 事件。这种“订阅-发布”的事件模型好处是解耦:聊天页面、未读角标、声音提醒各订阅各的,互不干扰;新增一种消息类型,不需要改连接类,只需要加一个 handler。

这个类还有个很细的设计:发送时把消息的 idmessageIdpictureId 统一转成字符串。因为 JSON 里数字和字符串经常混着来,后端解析时要猜类型(Integer 还是 Long 还是 String),前端先统一成字符串,能少踩一堆类型推断的坑。

第二个连接:聊天列表与未读数

聊天不止聊天室,还有一个“聊天列表”——左边是会话列表,右边是聊天气泡。列表页需要实时更新未读数和会话排序,所以单独开了一个 /api/ws/chat-list 连接,用的 ChatListWebSocket单例模式

private static instance: ChatListWebSocket | null = null
private constructor() {}
public static getInstance(): ChatListWebSocket {
  if (!ChatListWebSocket.instance) {
    ChatListWebSocket.instance = new ChatListWebSocket()
  }
  return ChatListWebSocket.instance
}

为什么要单例?因为聊天列表是全局的——不管你在哪个页面,未读数都得实时更新。如果每个页面都 new 一个连接,就会有一堆重复的 WebSocket 在空转,服务端还要维护一堆重复 session。单例保证全站只有一个聊天列表连接,谁要数据都从它这里过。

信息中心界面 ▲ 图:信息中心——未读数、通知、会话,都靠那一条全局 WebSocket 实时刷新

未读数的更新有一个处理竞态的细节,代码注释里写得很直白:

public requestUnreadCounts(): void {
  // 增加 500ms 延迟,以避开后端数据库或缓存同步的异步事务写入空窗期(竞态条件)
  setTimeout(() => {
    this.sendMessage({ type: 'REQUEST_UNREAD_COUNTS' })
  }, 500)
}

刚发完一条消息,立刻去请求未读数,后端那条消息可能还没提交事务,查出来就少一条。这个 500ms 延迟是血泪换来的:不加它,未读数会时不时地“差一条”,看着就像 bug,但又不是每次都复现——最折磨人的那种问题。

弹幕在 canvas 上飘

弹幕和聊天消息是两套完全不同的渲染逻辑。聊天消息是 DOM 列表,弹幕是 canvas 上的飘屏——因为弹幕要同时飘几十上百条,用 DOM 会卡成幻灯片,canvas 才能撑住 60 帧。

弹幕的移动逻辑核心就两行:

// 速度 = 视口宽度 / 播放时长,保证不同屏幕下弹幕穿越时间一致
speed: (viewportWidth.value + 400) / (barrageSpeed.value / 1000)

// 每一帧:位置按时间增量移动
barrage.currentX -= barrage.speed * deltaTime
animationFrameId = requestAnimationFrame(updateBarragePositions)

第一行很关键:弹幕的“速度”不是拍脑袋定的,而是视口宽度除以时长。这样同一个弹幕,在手机窄屏和电脑宽屏上,穿越整屏的时间是相同的——视觉节奏一致,不会出现手机上飞得飞快、电脑上慢吞吞的情况。+400 是给弹幕本身的宽度留的余量。

第二行是动画的标准做法:requestAnimationFrame 每一帧回调一次,位置减去“速度 × 帧间隔”,再用 deltaTime(真实帧间隔)修正。为什么用 deltaTime 而不是直接 x -= speed?因为帧率不是恒定的——掉到 30 帧的时候,如果不修正,弹幕速度就会变成一半,看起来一卡一卡的。用 deltaTime 修正后,无论帧率多少,每帧移动的距离都能换算成相同的时间位移,动画就是匀速的。

canvas 弹幕还有几个细节活。先说按行分配:屏幕纵向按弹幕高度切成几行,新弹幕优先塞进最空的那一行,避免弹幕叠成墙。再看命中检测:判断一条新弹幕会不会撞上同一条轨道上正在飘的弹幕,会撞就换一行或者微调速度。最后是回收:弹幕 x 坐标飘出屏幕左边之后,把它从活跃列表里移除,下一帧就不再绘制——不回收的话,列表越积越长,每帧绘制的内容越来越多,帧率就掉下去了。

还有一个所有 canvas 弹幕都要注意的点:文本绘制成本。每条弹幕每帧都要 fillText 一次,几十条弹幕就是几十次绘制,所以弹幕字体要尽量用同一套配置,避免每次绘制都重新解析字体;把字体字符串提成常量,能省下肉眼可见的绘制开销。骨架就是上面那两行——速度公式和帧循环,剩下的都是在这个骨架上做文章。

AI 助手走的是另一条道

这个聊天系统里还藏着一个 AI 助手,用户在公共空间里 @悦目小助手 提问,AI 异步回复。AI 回复是个耗时操作(要调外部模型接口),而且同一时刻只能有一个 AI 在回复,否则会乱。所以它单独走了一条通道:

AI Queue

private final LinkedBlockingQueue<AIMessageTask> aiMessageQueue = new LinkedBlockingQueue<>(1000);
private final ThreadPoolExecutor aiMessageExecutor = new ThreadPoolExecutor(
        1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>(), ...);

核心线程数 1 的线程池 + 一个容量 1000 的阻塞队列:AI 提问先入队,单线程逐个处理,保证“同一时刻只有一个 AI 在思考”;队列满了,offer 返回 false,直接回一条“抱歉,我现在有点忙”给用户,而不是无限堆积把内存打爆。这个设计跟 Disruptor 的背压思路如出一辙——处理不过来的请求,宁可拒掉,也不能让系统崩溃

有意思的是,AI 消息不走 Disruptor,而是单独一个队列。为什么?因为 AI 回复的“单位成本”高(一次外部 API 调用可能几百毫秒到几秒),而且和普通聊天消息的实时性要求完全不同。把它单独拎出来,是为了不让一次慢腾腾的 AI 调用堵住普通消息的广播通道——普通消息要毫秒级,AI 可以等,两者不该抢同一个线程。

单线程还有另一层深意:保证顺序。AI 的回复是基于上下文对话的,如果两个问题同时触发两次回复,回复顺序一乱,用户看到的对话就驴唇不对马嘴。单线程队列天然保证“先问的先答”,不需要额外的锁或排序。调 deepSeekService.generateAssistantResponse 的时候,整个线程被占住,后面的问题排队等——对聊天产品来说,用户等 AI 回复本来就比等普通消息更有耐心,所以这种串行是可接受的。

队列满时的处理也值得一提:aiMessageQueue.offer 返回 false,就直接回一条“抱歉,我现在有点忙,请稍后再问我问题”给用户。这不是偷懒,是防止“请求堆积 → 内存爆炸 → 整个聊天服务跟着挂”的雪崩。宁可拒绝新请求,也不能让一个 AI 功能拖垮整站。

踩坑记录

这条链路踩过的坑,我回头分成了三组——一组是 Disruptor 的规矩,一组是"本地和服务器两份状态"的纠缠,一组是连接和帧的小细节。分组比一条条列着清楚,因为每组的坑,根源都是同一个。

Disruptor 的规矩,破一次就卡死。 next()publish() 必须成对——最早我在 try 块中间抛了异常,publish 没执行,槽位占死,环形缓冲固定大小,一个槽被占可用就少一个,消息多了生产者全卡在 next() 上。症状是"系统没崩,但消息完全不动",排查了很久。event.clear() 漏字段也一样毒——事件是复用的,漏清一个字段,下一个复用到这个槽位就带着上一个的数据,我们踩到过一次"上一个人的消息被带到下一个人头上"。所以 publish 必须放 finallyclear() 有哪些字段、set 就得覆盖哪些,写代码时逐字段对着看。

"本地一份、服务器一份"的纠缠。 乐观回显的临时 ID 在消息回得特别快时会翻车——先加了临时那条,又加了服务器那条,同一条消息显示两次;断线重连后重取历史,本地消息还在列表里,又叠出一遍。这俩的修法都是"先按 ID 找,找到替换、找不到才新增"。这类 ID 映射和幂等的问题,只有在高并发或者断线重连下才露头,单机测试永远测不出来。

连接和帧的小细节。 心跳最初设 5 秒,网络一抖就重连,重连又握手失败,形成"重连风暴"打满握手线程——后来拉长到 24 秒、加"两倍间隔没响应才重连"的判断、用指数退避限速才压住。弹幕早期 x -= speed 直接减,低帧率设备越飘越慢,改成 speed * deltaTime 才好。还有在线人数:早期把 User 实体整个序列化发出去,手机号邮箱跟着在线列表漏光——后来改 safeUser 只挑展示字段。这几跤提醒我:心跳别太勤、canvas 动效用时间增量、往外发的对象永远单独做"对外视图"

60 秒的后悔药:撤回

聊天里最常用也最容易被忽略的功能是撤回。这个项目的撤回规则是:只能撤回自己的消息、只能在 60 秒内。前端发一条 RECALL 类型消息,带上 messageId,后端 handleRecallMessage 处理。

核心校验三连:消息存在吗、是你发的吗、超过 60 秒了吗——三个条件缺一个就拒绝。尤其是第三个,60 秒的窗口是产品定的规矩:太短,误发的消息救不回来;太长,对面都看完了再撤回,反而诡异。

通过校验后,后端不是把消息删掉,而是把内容改成“消息已撤回”再更新——为什么保留而不是删除?因为聊天记录要有审计性,删掉就死无对证了。改完之后,把撤回通知广播给整个房间,所有在线的 session 收到 RECALL 后,前端把对应的气泡替换成“消息已撤回”。

注意,撤回是不走 Disruptor 的——它直接在 WebSocket 消息处理线程里同步执行了。因为撤回是低频操作,不需要高吞吐,而且它要读库、改库、再广播,走不走 Disruptor 都一样。这也说明:Disruptor 是为“高并发广播”准备的,不是所有消息都要往里面塞,得按场景选。

历史消息,一次 20 条

聊天窗口打开时,只加载最近 20 条历史消息。往上滑想看更早的,前端发一条 loadMore 类型的消息,带上当前页码,后端 handleLoadMoreMessage 再去查下一页。

这个“20 条 + 分页”的设计值得说:聊天消息表是无脑增长的,一次性全拉回来既慢又占内存,前端渲染也会卡。每页 20 条是权衡的结果——够撑起聊天窗口的可视区域,又不至于拉太久。

一个常被忽略的转换细节是:前端传过来的 pagepictureIdspaceId 可能是数字也可能是字符串,后端代码里写了一堆 instanceof Integer / Long / String 的三段式转换。看着啰嗦,但这是真实项目里必然的防御——JSON 跨语言序列化时数字类型经常变,写死一种类型迟早翻车。

分页还返回一个 hasMore 标志,前端据此决定还能不能继续往上翻,翻到底了就显示“没有更早的消息了”。这个标志最容易漏:不返回它,前端永远不知道能不能继续加载,只能无限请求下去。

谁在线,谁离线

聊天室顶部通常会显示“xx 人在线”。这个数据怎么来的?连接建立时、有人进出时,后端广播一条 onlineUsers 消息,前端订阅这个类型,更新在线人数和列表。

对空间聊天来说,逻辑稍微复杂一点:一个空间有成员列表,而在线的人只是其中一部分。所以 broadcastOnlineUsers 要算两笔账——在线用户数(当前连接着的),和离线用户数(成员总数减去在线数)。离线用户是这样算出来的:

List<User> allMembers = chatMessageService.getSpaceMembers(spaceId);
Set<User> offlineUsers = new HashSet<>();
for (User member : allMembers) {
    boolean isOnline = onlineUsers.stream()
            .anyMatch(u -> u.getId().equals(member.getId()));
    if (!isOnline) offlineUsers.add(member);
}

这里还藏着一个性能隐患:onlineUsers.stream().anyMatch 对每个成员都线性扫描一遍在线列表,成员多、在线多的时候是 O(n×m) 的。数据量小没问题,但设计上其实可以先把在线用户转成 Set<Long>contains,复杂度就降成 O(n+m)。真实项目里往往先求能用,再回头优化这种小地方,但也得心里有数。

另外,在线用户发出去之前做了脱敏——只挑 id、昵称、头像这些展示字段,重新组装一个对象,不把整个 User 实体发出去。这就是坑六里说的“对外视图”。

收尾:三个词

把整条链路收起来看,就是三个词:

无锁。 Disruptor 用预分配 + CAS 环形缓冲替代锁,对象复用替代垃圾回收,从底层摁住了并发和 GC 两个大头。

串行。 无论多少连接同时发消息,真正落库和广播的只有 Disruptor 那一个工作线程,session 永远不会被并发写。这是整套设计最核心的一招。

心跳。 前端 24 秒判活、指数退避重连、队列兜底消息不丢;后端 13 秒扫描、10 分钟兜底关僵尸连接。一前一后,把“连接会断”这件必然发生的事,处理成了用户几乎感知不到的事。

这套架构还有一个没说透的点:它把“复杂”关在了笼子里。前端要处理心跳、重连、队列、乐观回显,后端要处理四张 session 表、Disruptor、心跳线程、AI 队列——单看任何一环都不算难,难的是让它们协同不打架。Disruptor 把后端最危险的“并发写 session”关进单线程的笼子,前端把最危险的“消息会丢”用队列兜住,剩下的细节都是在这个笼子外围打补丁。

最后说一句大实话:Disruptor 单独拎出来,是个有点小众、有点炫技的组件。但放在聊天这个场景里,它解决的恰恰是“并发写 session”这个最要命的问题,而且解决得干净——不是靠锁堆出来的“能用”,是靠架构选出来的“不会出事”。这大概就是选型的意义:不是找一个最潮的轮子,而是找一个最贴合问题的答案。

如果哪天你的聊天系统在并发上翻车了,回头翻翻这篇文章,多半能找到对应的那一节。而如果你也在写聊天功能,不妨先问自己一句:我的消息,是让每个线程各发各的,还是先串起来再发?

其实还有一句话没在正文里说:这套架构上线后,最让我踏实的不是 Disruptor 有多快,而是线上真的再没出现过"两个线程同时写一个 session"的报错。技术选型对不对,终究要靠生产环境来回答。写这篇也是想把这些"生产环境替我回答过的问题"传给你——下次再遇到"要不要给消息加个队列"的纠结,你会记得,这里已经有人替你踩过一遍。