← 返回首页目录
# FastAPI WebSockets 完整指南
作者:吉祥法师
## 核心概念
WebSocket是一种在单个TCP连接上进行全双工通信的协议,允许客户端和服务器之间进行实时、双向的数据传输。与传统的HTTP请求-响应模式不同,WebSocket连接一旦建立,双方都可以随时发送数据,无需等待对方的请求。FastAPI作为现代Python Web框架,原生支持WebSocket协议,为开发者提供了构建实时应用(如聊天室、实时通知、在线游戏等)的便捷工具。
FastAPI的WebSocket支持建立在Starlette框架之上,提供了简洁的装饰器语法和异步支持。通过`@app.websocket`装饰器,开发者可以轻松创建WebSocket端点,并使用类型注解、依赖注入等FastAPI特性。在高并发场景下,WebSocket比轮询(Polling)更高效,因为它避免了HTTP请求的头开销和延迟。
## 环境安装与配置
### 安装WebSocket库
在使用FastAPI的WebSocket功能前,需要确保已安装必要的依赖包。`websockets`是Python生态中广泛使用的WebSocket实现库,FastAPI依赖它来处理底层的WebSocket协议通信。
```bash
$ uv add websockets
```
该命令会将`websockets`库添加到项目依赖中,确保FastAPI能够正常处理WebSocket连接。安装完成后,开发者就可以在FastAPI应用中使用WebSocket相关功能。
## 创建第一个WebSocket端点
### 基础WebSocket服务器实现
创建WebSocket端点的基本流程包括:定义路径、接收连接、循环处理消息。以下展示一个完整的示例,包含HTTP路由返回HTML页面和WebSocket端点处理消息。
```python
from fastapi import FastAPI, WebSocket
from fastapi.responses import HTMLResponse
app = FastAPI()
@app.get("/")
async def get():
return HTMLResponse(html)
@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
await websocket.accept()
while True:
data = await websocket.receive_text()
await websocket.send_text(f"Message text was: {data}")
```
在这个示例中,`@app.websocket("/ws")`装饰器定义了一个WebSocket路由。当客户端连接到`/ws`路径时,FastAPI会调用`websocket_endpoint`函数。函数首先调用`websocket.accept()`接受连接,然后进入无限循环,等待接收客户端消息并发送响应。
### 网页客户端实现
为了测试WebSocket功能,可以创建一个简单的HTML页面,使用JavaScript的WebSocket API与服务器通信。
```python
html = """
Chat
WebSocket Chat
"""
```
这个HTML页面创建了一个WebSocket连接,当用户输入消息并点击发送按钮时,消息通过WebSocket发送到服务器。服务器返回的响应通过`onmessage`回调函数显示在页面的列表中。
## 生产环境考量
在实际生产环境中,前端通常使用现代JavaScript框架(如React、Vue.js或Angular)构建,这些框架都提供了WebSocket客户端库。对于原生移动应用,可以直接使用系统提供的WebSocket API。上述HTML示例仅用于演示和测试WebSocket服务器端功能。
## WebSocket消息类型处理
### 文本、二进制和JSON数据
FastAPI的WebSocket支持接收和发送多种数据类型。`receive_text()`方法接收文本消息,`send_text()`发送文本消息。除了文本,还可以处理二进制数据:
```python
# 接收二进制数据
data = await websocket.receive_bytes()
# 发送二进制数据
await websocket.send_bytes(b"Binary data")
# 使用JSON格式
import json
data = json.loads(await websocket.receive_text())
await websocket.send_text(json.dumps({"message": "response"}))
```
### 多格式消息处理
在复杂应用场景中,可能需要同时处理不同类型的消息。可以使用`receive_json()`和`send_json()`方法实现JSON消息的自动序列化和反序列化。
## 依赖注入与参数处理
### 在WebSocket中使用FastAPI特性
FastAPI的WebSocket端点支持与普通HTTP端点相同的依赖注入机制,包括`Depends`、`Security`、`Cookie`、`Header`、`Path`和`Query`等。这使得开发者可以在WebSocket连接中复用已有的验证和业务逻辑。
```python
from typing import Annotated
from fastapi import (Cookie, Depends, FastAPI, Query, WebSocket, WebSocketException, status)
async def get_cookie_or_token(
websocket: WebSocket,
session: Annotated[str | None, Cookie()] = None,
token: Annotated[str | None, Query()] = None,
):
if session is None and token is None:
raise WebSocketException(code=status.WS_1008_POLICY_VIOLATION)
return session or token
@app.websocket("/items/{item_id}/ws")
async def websocket_endpoint(
*,
websocket: WebSocket,
item_id: str,
q: int | None = None,
cookie_or_token: Annotated[str, Depends(get_cookie_or_token)],
):
await websocket.accept()
while True:
data = await websocket.receive_text()
await websocket.send_text(f"Session cookie or query token value is: {cookie_or_token}")
if q is not None:
await websocket.send_text(f"Query parameter q is: {q}")
await websocket.send_text(f"Message text was: {data}, for item ID: {item_id}")
```
### WebSocketException错误处理
WebSocket连接中发生错误时,不适合抛出`HTTPException`,因为HTTP状态码在WebSocket上下文中没有意义。FastAPI提供了`WebSocketException`类,用于处理WebSocket特定的异常情况。规范定义了多个WebSocket关闭代码,如`WS_1008_POLICY_VIOLATION`表示策略违规。
## 高级功能:多客户端管理
### 连接管理器模式
在真实应用中,经常需要同时管理多个WebSocket客户端连接。可以创建一个`ConnectionManager`类来统一管理连接的生命周期和消息广播。
```python
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
class ConnectionManager:
def __init__(self):
self.active_connections: list[WebSocket] = []
async def connect(self, websocket: WebSocket):
await websocket.accept()
self.active_connections.append(websocket)
def disconnect(self, websocket: WebSocket):
self.active_connections.remove(websocket)
async def send_personal_message(self, message: str, websocket: WebSocket):
await websocket.send_text(message)
async def broadcast(self, message: str):
for connection in self.active_connections:
await connection.send_text(message)
manager = ConnectionManager()
@app.websocket("/ws/{client_id}")
async def websocket_endpoint(websocket: WebSocket, client_id: int):
await manager.connect(websocket)
try:
while True:
data = await websocket.receive_text()
await manager.send_personal_message(f"You wrote: {data}", websocket)
await manager.broadcast(f"Client #{client_id} says: {data}")
except WebSocketDisconnect:
manager.disconnect(websocket)
await manager.broadcast(f"Client #{client_id} left the chat")
```
### WebSocketDisconnect异常处理
当客户端断开连接时,`receive_text()`会抛出`WebSocketDisconnect`异常。通过捕获这个异常,可以执行清理操作,如从活跃连接列表中移除断开的连接。使用`try-except`块可以优雅地处理客户端意外断线的情况。
### 消息广播机制
`broadcast`方法遍历所有活跃的WebSocket连接,向每个客户端发送消息。这在聊天室应用中尤为重要,能够确保所有在线用户都能收到其他用户发送的消息。
## 中断与连接断开处理
### 完整的断开连接示例
以下是一个完整的聊天应用示例,演示如何处理多个客户端的连接和断开:
```python
html = """
Chat
WebSocket Chat
Your ID:
"""
```
### 测试断开连接
要测试这个示例,需要在多个浏览器标签页中打开应用。当一个标签页关闭时,服务器端`WebSocketDisconnect`异常被触发,其他所有客户端都会收到类似"客户端已离开聊天"的通知。这种机制在构建实时协作工具时尤为重要。
## 分布式部署考虑
上述`ConnectionManager`示例使用内存列表存储连接信息,这仅在单进程应用且进程持续运行时有效。对于分布式或多进程部署,可以使用Redis、PostgreSQL等外部服务来管理WebSocket连接状态。`encode/broadcaster`等库提供了与FastAPI集成且支持多进程的广播解决方案,它们利用Redis或其他消息代理实现跨进程通信。
## 最佳实践与性能优化
### 连接生命周期管理
合理管理WebSocket连接的生命周期是构建稳健应用的关键。在连接建立时执行必要的认证和初始化,在连接断开时及时清理资源。使用`WebSocketDisconnect`异常处理机制可以确保资源被正确释放。
### 心跳机制与超时处理
在长时间运行的WebSocket连接中,网络中断或客户端异常可能导致连接僵死。实现心跳机制(Heartbeat)可以定期发送ping消息检测连接是否仍然活跃,并在超时后主动关闭无效连接。
### 消息大小限制
WebSocket协议允许发送任意大小的消息,但过大的消息可能阻塞事件循环或占用过多内存。建议设置合理的消息大小限制,并使用`max_size`参数控制传入消息的最大字节数。
### 错误处理与日志记录
WebSocket应用中,错误处理和日志记录同样重要。适当记录连接建立、关闭和异常事件,有助于监控系统的运行状态。同时,在处理业务逻辑时捕获并处理可能出现的异常,避免未处理异常导致进程崩溃。
## 扩展与集成
### 与其他库的集成
WebSocket可以与其他FastAPI功能无缝集成,如依赖注入、安全验证、数据库操作等。可以利用`Depends`机制在WebSocket连接建立时执行身份验证,使用数据库连接池存储和检索消息记录。
### 与消息队列集成
在生产级应用中,WebSocket通常与消息队列(如Redis Pub/Sub、RabbitMQ)集成,实现事件的异步分发。这样可以处理大量并发连接,并确保消息在多个服务器实例之间的传递。
### 集群部署与负载均衡
当应用需要水平扩展时,需要考虑WebSocket连接在多个服务器实例之间的负载均衡。这通常需要配置负载均衡器支持WebSocket的粘性会话(Sticky Session),确保同一客户端的请求始终指向同一服务器实例。
## 调试与测试
可以使用FastAPI的`TestClient`类对WebSocket端点进行测试。需要注意,TestClient对WebSocket测试有异步事件循环的限制,需要正确管理事件循环。可以编写单元测试和集成测试来验证WebSocket功能的正确性。
## 相关资源
- Starlette WebSocket文档提供了完整的`WebSocket`类参考
- 基于类的WebSocket处理器需要查看Starlette官方文档
- FastAPI官方文档提供更多高级特性和配置选项
WebSocket作为实时通信的重要组成部分,在FastAPI中得到了充分的支持和优化。通过掌握上述核心概念和实践技巧,开发者可以快速构建高性能的实时应用。随着WebSocket技术的不断发展,FastAPI提供了持续更新的功能和最佳实践,建议保持关注最新的技术趋势和框架更新。