WebSocket 协议与实时通信

深入 WebSocket 协议握手机制、帧格式与心跳保活,掌握 Spring 与 Netty 实现及大规模实时推送架构

WebSocket 是一种在单个 TCP 连接上进行全双工通信的协议,它使得浏览器和服务器之间可以建立持久连接,实现低延迟、高效率的实时双向数据传输。

一、为什么需要 WebSocket

1.1 传统方案对比

方案原理延迟缺点
轮询定时 HTTP 请求高(取决于间隔)无效请求多,浪费带宽
长轮询挂起请求直到有数据服务器压力大,消息可能乱序
SSE服务端单向推送仅服务端→客户端,不支持二进制
WebSocket全双工长连接极低协议复杂度高,需处理断线重连

1.2 适用场景

  • 即时通讯(IM、聊天室)
  • 在线协作(文档编辑、白板)
  • 实时数据(股票行情、监控大盘)
  • 在线游戏(操作同步、位置更新)
  • IoT 设备数据上报

二、WebSocket 握手

2.1 HTTP Upgrade 过程

客户端请求(HTTP/1.1):

GET /chat HTTP/1.1
Host: server.example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Sec-WebSocket-Version: 13
Origin: http://example.com

服务端响应(101 Switching Protocols):

HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=

Sec-WebSocket-Accept 计算:
  base64(sha1(Sec-WebSocket-Key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"))

2.2 Spring Boot WebSocket

@Configuration
@EnableWebSocketMessageBroker
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {
    
    @Override
    public void configureMessageBroker(MessageBrokerRegistry registry) {
        registry.enableSimpleBroker("/topic", "/queue");
        registry.setApplicationDestinationPrefixes("/app");
        registry.setUserDestinationPrefix("/user");
    }
    
    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
        registry.addEndpoint("/ws")
            .setAllowedOriginPatterns("*")
            .withSockJS();
    }
}

@Controller
public class ChatController {
    
    @Autowired
    private SimpMessagingTemplate messagingTemplate;
    
    @MessageMapping("/chat")
    @SendTo("/topic/messages")
    public ChatMessage handleChat(ChatMessage message) {
        message.setTimestamp(System.currentTimeMillis());
        return message;
    }
    
    public void sendToUser(String userId, ChatMessage message) {
        messagingTemplate.convertAndSendToUser(userId, "/queue/notify", message);
    }
}

2.3 Netty WebSocket 服务端

public class NettyWebSocketServer {
    
    public static void main(String[] args) throws Exception {
        EventLoopGroup boss = new NioEventLoopGroup(1);
        EventLoopGroup worker = new NioEventLoopGroup();
        
        try {
            ServerBootstrap bootstrap = new ServerBootstrap();
            bootstrap.group(boss, worker)
                .channel(NioServerSocketChannel.class)
                .childHandler(new ChannelInitializer<SocketChannel>() {
                    @Override
                    protected void initChannel(SocketChannel ch) {
                        ch.pipeline()
                            .addLast(new HttpServerCodec())
                            .addLast(new HttpObjectAggregator(65536))
                            .addLast(new WebSocketServerProtocolHandler("/ws"))
                            .addLast(new WebSocketFrameHandler());
                    }
                });
            
            ChannelFuture future = bootstrap.bind(8080).sync();
            future.channel().closeFuture().sync();
        } finally {
            boss.shutdownGracefully();
            worker.shutdownGracefully();
        }
    }
}

public class WebSocketFrameHandler extends SimpleChannelInboundHandler<WebSocketFrame> {
    
    private static final ChannelGroup channels = new DefaultChannelGroup(GlobalEventExecutor.INSTANCE);
    
    @Override
    public void handlerAdded(ChannelHandlerContext ctx) {
        channels.add(ctx.channel());
    }
    
    @Override
    protected void channelRead0(ChannelHandlerContext ctx, WebSocketFrame frame) {
        if (frame instanceof TextWebSocketFrame) {
            String text = ((TextWebSocketFrame) frame).text();
            System.out.println("Received: " + text);
            
            // 广播给所有客户端
            channels.writeAndFlush(new TextWebSocketFrame(
                "[" + ctx.channel().id() + "] " + text
            ));
        }
    }
    
    @Override
    public void handlerRemoved(ChannelHandlerContext ctx) {
        channels.remove(ctx.channel());
    }
}

三、WebSocket 帧格式

WebSocket 数据帧:

 0                   1                   2                   3
 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-------+-+-------------+-------------------------------+
|F|R|R|R| opcode|M| Payload len |    Extended payload length    |
|I|S|S|S|  (4)  |A|     (7)     |             (16/64)           |
|N|V|V|V|       |S|             |   (if payload len==126/127)   |
| |1|2|3|       |K|             |                               |
+-+-+-+-+-------+-+-------------+ - - - - - - - - - - - - - - - +
|     Extended payload length continued, if payload len == 127  |
+ - - - - - - - - - - - - - - - +-------------------------------+
|                               |Masking-key, if MASK set to 1  |
+-------------------------------+-------------------------------+
| Masking-key (continued)       |          Payload Data         |
+-------------------------------- - - - - - - - - - - - - - - -+
:                     Payload Data continued ...                :
+ - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - +
|                     Payload Data continued ...                |
+---------------------------------------------------------------+

FIN: 是否为最后一帧
RSV: 保留位
Opcode: 0=继续, 1=文本, 2=二进制, 8=关闭, 9=Ping, 10=Pong
MASK: 客户端必须掩码,服务端不应掩码
Payload len: 数据长度(7/16/64 bits)

四、心跳与保活

4.1 Ping/Pong 机制

// 服务端定期发送 Ping
public class HeartbeatHandler extends ChannelInboundHandlerAdapter {
    
    private static final int PING_INTERVAL = 30; // 秒
    
    @Override
    public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
        if (evt instanceof IdleStateEvent) {
            IdleStateEvent event = (IdleStateEvent) evt;
            if (event.state() == IdleState.READER_IDLE) {
                // 客户端长时间未发送消息,关闭连接
                ctx.close();
            } else if (event.state() == IdleState.WRITER_IDLE) {
                // 发送 Ping
                ctx.writeAndFlush(new PingWebSocketFrame());
            }
        }
    }
}

// Pipeline 配置
ch.pipeline()
    .addLast(new IdleStateHandler(60, 30, 0))  // 读空闲60s,写空闲30s
    .addLast(new HeartbeatHandler());

五、大规模推送架构

大规模 WebSocket 推送架构:

客户端 ←──→ CDN / WAF
              │
              ▼
┌──────────────────────────────┐
│      L7 负载均衡器            │
│   (Nginx / HAProxy / Envoy)  │
│   ws://example.com/ws        │
└──────────────┬───────────────┘
               │
    ┌──────────┼──────────┐
    ▼          ▼          ▼
┌───────┐  ┌───────┐  ┌───────┐
│ WS-1  │  │ WS-2  │  │ WS-3  │   WebSocket 服务集群
│ (10k) │  │ (10k) │  │ (10k) │   每台维护 1-10 万连接
└───┬───┘  └───┬───┘  └───┬───┘
    │          │          │
    └──────────┼──────────┘
               ▼
┌──────────────────────────────┐
│        Redis Pub/Sub         │   消息总线
│   /  RabbitMQ / Kafka        │   广播消息到所有 WS 节点
└──────────────────────────────┘
               ▲
               │
        业务服务发送消息

关键设计:
1. 有状态服务 → 客户端需保持连接到同一节点
2. 消息总线 → 所有节点共享消息
3. 心跳检测 → 及时清理死连接
4. 连接数限制 → 单机不超过 ulimit

六、连接管理

@Component
public class WebSocketSessionManager {
    
    private final ConcurrentHashMap<String, WebSocketSession> sessions = new ConcurrentHashMap<>();
    
    public void register(String userId, WebSocketSession session) {
        WebSocketSession old = sessions.put(userId, session);
        if (old != null && old.isOpen()) {
            try {
                old.close();  // 踢掉旧连接
            } catch (IOException e) {
                log.warn("关闭旧连接失败", e);
            }
        }
    }
    
    public void unregister(String userId) {
        sessions.remove(userId);
    }
    
    public void sendToUser(String userId, String message) {
        WebSocketSession session = sessions.get(userId);
        if (session != null && session.isOpen()) {
            try {
                session.sendMessage(new TextMessage(message));
            } catch (IOException e) {
                log.error("发送消息失败: {}", userId, e);
            }
        }
    }
    
    public void broadcast(String message) {
        sessions.values().forEach(session -> {
            if (session.isOpen()) {
                try {
                    session.sendMessage(new TextMessage(message));
                } catch (IOException e) {
                    log.error("广播消息失败", e);
                }
            }
        });
    }
    
    public int getOnlineCount() {
        return sessions.size();
    }
}

七、安全性考量

// Origin 校验
@Override
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response,
        WebSocketHandler wsHandler, Map<String, Object> attributes) {
    String origin = request.getHeaders().getOrigin();
    if (!ALLOWED_ORIGINS.contains(origin)) {
        return false;  // 拒绝非白名单 Origin
    }
    return true;
}

// JWT 认证(通过 URL 参数传递 Token)
@Override
public void afterConnectionEstablished(WebSocketSession session) {
    String token = extractToken(session);
    if (!jwtUtil.validate(token)) {
        session.close(CloseStatus.POLICY_VIOLATION);
        return;
    }
    String userId = jwtUtil.getUserId(token);
    sessionManager.register(userId, session);
}

八、总结

方面建议
协议选择WebSocket 适合双向实时,SSE 适合服务端推送
框架选型Spring STOMP 快速开发,Netty 高性能定制
集群部署Redis Pub/Sub 做消息广播,保证多节点一致
连接管理心跳保活 + 超时清理 + 单用户单连接
安全防护Origin 校验 + JWT 认证 + 连接数限制

WebSocket 是实现实时通信的核心技术。在生产环境中,需要重点关注连接的稳定性(心跳、重连)、扩展性(集群、消息广播)和安全性(认证、防攻击)。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「network」更多文章

  1. 网络安全:TLS/SSL、证书与加密通信
  2. 负载均衡算法与高可用架构
  3. DNS 系统与智能解析