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 是实现实时通信的核心技术。在生产环境中,需要重点关注连接的稳定性(心跳、重连)、扩展性(集群、消息广播)和安全性(认证、防攻击)。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。