WebSocket 从入门到实战:用 FastAPI 打造实时租房系统
—— 零基础也能懂的实时通信完全指南 ——
适用人群:Python 初学者 / Web 开发新手 / 对实时通信感兴趣的开发者
目录
六、前端实现:JavaScript 连接 WebSocket
一、什么是 WebSocket?为什么需要它?
想象一下这个场景:你在租房网站上浏览房源,突然管理员上架了一套新房子,但你却不知道,需要手动刷新页面才能看到。这体验是不是很糟糕?
传统 HTTP 请求是"一问一答"模式:客户端发起请求 → 服务器返回响应 → 连接关闭。如果想获取最新数据,只能不断轮询(每隔几秒请求一次),但这会带来两个问题:
- 浪费资源:大部分请求返回的数据都是一样的
- 延迟高:两次请求之间的时间差导致信息不及时
WebSocket 就是为了解决这个问题而生的!它是一种双向通信协议,建立连接后,客户端和服务器可以随时互相发送消息,就像打电话一样,双方都能主动说话。
使用 WebSocket 的好处:
- ✅ 实时性:消息毫秒级到达
- ✅ 高效:不需要频繁建立/关闭连接
- ✅ 双向:服务器可以主动推送数据
- ✅ 轻量:头部信息小,节省带宽
二、WebSocket vs HTTP:有什么区别?
用一个生活中的例子来理解:
📮 HTTP 就像写信:你寄一封信(请求)→ 对方收到后回信(响应)→ 结束。如果想知道新消息,只能再寄一封信去问。
📞 WebSocket 就像打电话:拨通后(建立连接)→ 双方可以随时说话(双向通信)→ 直到挂断(关闭连接)。
|
特性 |
HTTP |
WebSocket |
|
通信方式 |
单向(客户端发起) |
双向(任意一方发起) |
|
连接状态 |
短连接(用完即关) |
长连接(持续保持) |
|
适用场景 |
获取数据、提交表单 |
实时聊天、消息推送 |
|
性能开销 |
每次都要建立连接 |
一次建立,多次复用 |
结论:如果你的应用需要实时更新(如聊天室、股票行情、在线游戏),WebSocket 是不二之选!
三、FastAPI + WebSocket:环境准备
步骤 1:安装依赖
pip install fastapi uvicorn websockets
步骤 2:创建基础项目结构
project/
├── app/
│ ├── core/
│ │ └── websocket_manager.py # WebSocket 连接管理器
│ ├── api/
│ │ └── websocket.py # WebSocket 路由
│ └── main.py # 主应用入口
└── requirements.txt
四、核心概念:连接管理器设计
在开始写代码之前,我们需要理解一个核心概念:如何管理大量的 WebSocket 连接?
想象你的租房系统有 1000 个用户同时在线,每个人都在不同的频道(比如"北京房源"、"上海房源"),你需要:
- 记录谁连接了哪个频道
- 给特定频道的所有用户发消息
- 给特定用户发消息(比如订单通知)
- 处理用户断开连接的情况
解决方案:设计一个 ConnectionManager 类,用字典高效管理所有连接!
五、后端实现:一步步搭建 WebSocket 服务
5.1 创建连接管理器
这是整个 WebSocket 系统的核心!我们用两个字典来管理连接:
- active_connections:按频道分组:{"houses": [ws1, ws2], "public": [ws3]}
- user_connections:按用户分组:{user_id: {ws1, ws2}}
from typing import Dict, List, Set
from fastapi import WebSocket
import json
class ConnectionManager:
def __init__(self):
# 按频道管理连接
self.active_connections: Dict[str, List[WebSocket]] = {}
# 按用户管理连接
self.user_connections: Dict[int, Set[WebSocket]] = {}
async def connect(self, websocket: WebSocket, channel: str = "public", user_id: int = None):
"""接受 WebSocket 连接并注册"""
await websocket.accept() # 握手
# 添加到频道
if channel not in self.active_connections:
self.active_connections[channel] = []
self.active_connections[channel].append(websocket)
# 关联到用户
if user_id:
if user_id not in self.user_connections:
self.user_connections[user_id] = set()
self.user_connections[user_id].add(websocket)
def disconnect(self, websocket: WebSocket, channel: str, user_id: int = None):
"""断开连接并清理"""
if channel in self.active_connections:
if websocket in self.active_connections[channel]:
self.active_connections[channel].remove(websocket)
if user_id and user_id in self.user_connections:
self.user_connections[user_id].discard(websocket)
if not self.user_connections[user_id]:
del self.user_connections[user_id]
async def broadcast_to_channel(self, message: dict, channel: str):
"""向指定频道广播消息"""
if channel in self.active_connections:
dead_connections = []
for connection in self.active_connections[channel]:
try:
await connection.send_json(message)
except:
dead_connections.append(connection)
# 清理失效连接
for ws in dead_connections:
self.active_connections[channel].remove(ws)
async def send_personal_message(self, message: dict, user_id: int):
"""向特定用户发送消息"""
if user_id in self.user_connections:
for websocket in self.user_connections[user_id]:
try:
await websocket.send_json(message)
except:
pass
# 创建全局实例
manager = ConnectionManager()
💡 关键点解析:
- await websocket.accept():完成 WebSocket 握手,建立连接
- dead_connections:收集发送失败的连接,避免内存泄漏
- Set[WebSocket]:一个用户可以有多个连接(多标签页)
- send_json():自动将字典转换为 JSON 字符串发送
5.2 创建 WebSocket 路由
路由负责处理 WebSocket 连接的整个生命周期:连接 → 接收消息 → 断开
from fastapi import APIRouter, WebSocket, WebSocketDisconnect, Query
from app.core.websocket_manager import manager
router = APIRouter(tags=["websocket"])
@router.websocket("/ws/houses")
async def websocket_houses(
websocket: WebSocket,
token: str = Query(None) # 通过 URL 参数传递 Token
):
"""房源更新频道"""
# 1. 验证用户身份(可选)
user_id = None
if token:
# 解码 JWT Token 获取 user_id
user_id = decode_token(token)
# 2. 建立连接,加入 "houses" 频道
await manager.connect(websocket, channel="houses", user_id=user_id)
# 3. 保持连接,监听消息
try:
while True:
data = await websocket.receive_text()
message = json.loads(data)
# 处理心跳包
if message.get("type") == "ping":
await websocket.send_json({
"type": "pong",
"timestamp": message.get("timestamp")
})
except WebSocketDisconnect:
# 4. 断开连接时清理
manager.disconnect(websocket, channel="houses", user_id=user_id)
连接生命周期:
- 客户端发起连接请求:ws://localhost:8080/ws/houses?token=xxx
- 服务器调用 connect() 接受连接
- 进入 while True 循环,等待客户端消息
- 收到 "ping" 消息,回复 "pong"(心跳检测)
- 客户端断开,捕获 WebSocketDisconnect 异常
- 调用 disconnect() 清理连接
5.3 在主应用中注册
from fastapi import FastAPI
from app.api.websocket import router as websocket_router
app = FastAPI(title="租房系统 API")
# 注册 WebSocket 路由
app.include_router(websocket_router)
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8080)
六、前端实现:JavaScript 连接 WebSocket
前端使用浏览器原生的 WebSocket API,非常简单!
// 1. 建立连接
const token = localStorage.getItem('token');
const ws = new WebSocket(`ws://localhost:8080/ws/houses?token=${token}`);
// 2. 连接成功
ws.onopen = () => {
console.log('WebSocket 连接成功!');
};
// 3. 接收消息
ws.onmessage = (event) => {
const data = JSON.parse(event.data);
if (data.type === 'house_update') {
console.log('收到房源更新:', data);
if (data.action === 'create') {
alert(`新房源上线: ${data.data.title}`);
// 重新加载房源列表
loadHouses();
} else if (data.action === 'update') {
alert(`房源更新: ${data.data.title}`);
} else if (data.action === 'delete') {
alert('房源已下架');
}
}
};
// 4. 连接断开 - 自动重连
ws.onclose = () => {
console.log('WebSocket 断开,5秒后重连...');
setTimeout(() => {
initWebSocket(); // 重新连接
}, 5000);
};
// 5. 错误处理
ws.onerror = (error) => {
console.error('WebSocket 错误:', error);
};
// 6. 发送心跳包(每30秒)
setInterval(() => {
if (ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify({
type: 'ping',
timestamp: Date.now()
}));
}
}, 30000);
📌 WebSocket API 速查:
- new WebSocket(url):创建连接
- ws.send(data):发送消息
- ws.close():主动关闭连接
- ws.readyState:连接状态(0=连接中,1=已打开,2=关闭中,3=已关闭)
七、实战案例:房源实时更新推送
现在我们把前后端结合起来,实现一个完整的功能:当管理员添加新房源时,所有在线用户立即收到通知!
7.1 后端:房源创建时广播
from app.core.websocket_manager import manager
from datetime import datetime
@router.post("/houses")
async def create_house(house_data: HouseCreate, db: Session = Depends(get_db)):
# 创建房源记录
new_house = create_house_in_db(db, house_data)
# 向所有用户推送房源创建通知
await manager.broadcast_house_created({
"id": new_house.id,
"title": new_house.title,
"price": new_house.price,
"city": new_house.city
})
return {"code": 200, "message": "房源创建成功"}
async def broadcast_house_created(self, house_data: dict):
"""广播房源创建消息到所有频道"""
message = {
"type": "house_update",
"action": "create",
"data": house_data,
"timestamp": datetime.now().isoformat()
}
# 同时推送至houses和public频道
await self.broadcast_to_channel(message, "houses")
await self.broadcast_to_channel(message, "public")
7.2 前端:接收并显示通知
// 处理房源更新消息
function handleHouseUpdate(data) {
console.log('收到房源更新:', data);
switch(data.action) {
case 'create':
showToast(`🎉 新房源上线: ${data.data.title}`, 'success');
break;
case 'update':
showToast(`📝 房源更新: ${data.data.title}`, 'info');
break;
case 'delete':
showToast(`🗑️ 房源已下架`, 'warning');
break;
}
loadRecommendHouses();
}
// 显示Toast提示
function showToast(message, type = 'info') {
const toast = document.createElement('div');
toast.className = `toast toast-${type}`;
toast.textContent = message;
document.body.appendChild(toast);
// 3秒后自动移除
setTimeout(() => toast.remove(), 3000);
}
八、进阶技巧:断线重连 + 心跳检测
8.1 断线重连机制
网络不稳定时,WebSocket 可能会断开。我们需要实现自动重连:
let ws = null;
let reconnectAttempts = 0;
const MAX_RECONNECT_ATTEMPTS = 10;
function initWebSocket() {
const token = localStorage.getItem('token');
ws = new WebSocket(`ws://localhost:8080/ws/houses?token=${token}`);
ws.onopen = () => {
console.log('WebSocket连接已建立');
reconnectAttempts = 0; // 重置重连计数器
};
ws.onclose = () => {
console.log('WebSocket连接已断开');
// 采用指数退避算法进行重连:1s, 2s, 4s, 8s...
if (reconnectAttempts < MAX_RECONNECT_ATTEMPTS) {
const delay = Math.min(1000 * Math.pow(2, reconnectAttempts), 30000);
reconnectAttempts++;
console.log(`将在 ${delay}ms 后进行第 ${reconnectAttempts} 次重连...`);
setTimeout(initWebSocket, delay);
} else {
console.error('达到最大重连次数,请刷新页面重试');
}
};
}
// 初始化WebSocket连接
initWebSocket();
💡 指数退避策略:每次重连间隔翻倍,避免频繁重连导致服务器压力过大
8.2 心跳检测(Heartbeat)
长时间没有数据传输的连接可能被防火墙或代理服务器关闭。心跳检测可以保持连接活跃:
# 后端:处理心跳响应
async def websocket_houses(websocket: WebSocket, ...):
await manager.connect(websocket, channel="houses")
try:
while True:
data = await websocket.receive_text()
message = json.loads(data)
if message.get("type") == "ping":
# 返回pong响应,附带原时间戳
await websocket.send_json({
"type": "pong",
"timestamp": message.get("timestamp")
})
except WebSocketDisconnect:
manager.disconnect(websocket, channel="houses")
九、常见问题与解决方案
Q1: WebSocket 和 HTTP 能共存吗?
A: 当然可以!FastAPI 同时支持两种协议,互不影响。
Q2: 如何保证 WebSocket 的安全性?
A: 使用 wss://(WebSocket Secure,类似 HTTPS),并在连接时验证 JWT Token。
Q3: WebSocket 断开后,之前的消息会丢失吗?
A: 会的。如果需要离线消息,需要配合数据库或消息队列(如 Redis)实现消息持久化。
Q4: 一个用户可以有多个 WebSocket 连接吗?
A: 可以!比如用户在多个浏览器标签页打开网站,每个标签页都会建立一个连接。
Q5: WebSocket 支持多少人同时在线?
A: 取决于服务器性能。一台普通服务器可以轻松支持数千个并发连接,优化后可达数万。
Q6: 如何在生产环境部署 WebSocket?
A: 使用 Nginx 作为反向代理,配置 WebSocket 支持;或使用云服务商的 WebSocket 服务。
十、总结与下一步学习
恭喜你完成了 WebSocket 的学习之旅!让我们回顾一下:
- ✅ 理解了 WebSocket 的原理和优势
- ✅ 学会了用 FastAPI 搭建 WebSocket 服务
- ✅ 掌握了连接管理器的设计模式
- ✅ 实现了前后端实时通信
- ✅ 学会了断线重连和心跳检测
下一步学习方向:
- 🚀 学习 Redis Pub/Sub,实现分布式 WebSocket 消息广播
- 🚀 探索 Socket.IO,获得更好的跨平台兼容性
- 🚀 实践更多场景:在线聊天室、协作编辑、实时游戏
- 🚀 研究 WebSocket 性能优化:压缩、二进制传输、连接池
记住:最好的学习方式就是动手实践!现在就试着给你的项目加上 WebSocket 功能吧!💪
— 完 —
更多推荐




所有评论(0)