FastAPI事件循环中多端口TCP服务器未启动问题排查求助
问题排查:TCP服务器未启动导致连接失败
我需要实现的功能:创建多实例TCP服务器,分别监听8001、8002、8003端口,将TCP接收的原始数据通过WebSocket转发给客户端,端到端流程为:client1(实例1)-> TCP -> server(实例2)-> WebSocket -> client2(实例2)
执行uvicorn main:app后,FastAPI服务器在本地正常启动,WebSocket连接(端口8000)已确认正常,但start_tcp_server未被触发,客户端尝试建立TCP连接时始终失败。服务器与客户端位于不同EC2实例,已确认实例间可通过服务器IP互通。
核心问题分析
- Lifespan上下文管理器逻辑错误:
startup_event中await asyncio.gather(...)写在yield之后,这部分代码会在FastAPI应用关闭阶段才执行,而非启动阶段,导致TCP服务器完全没在启动时初始化。 - 异步任务阻塞主线程:
start_tcp_server内的await server.serve_forever()是阻塞式异步调用,直接用gather会卡住整个应用主线程,需将TCP服务器任务作为后台异步任务启动。 - 冗余线程锁:当前TCP和WebSocket均为异步操作,无需
threading.Lock,反而引入不必要的复杂度。
修复后的代码
main.py
from fastapi import FastAPI, WebSocket import asyncio import globals from server import start_tcp_server from contextlib import asynccontextmanager @asynccontextmanager async def startup_event(app: FastAPI): # 启动阶段创建TCP服务器后台任务 print("start TCP servers") ports = [8001, 8002, 8003] tasks = [asyncio.create_task(start_tcp_server(port)) for port in ports] yield # 关闭阶段取消所有TCP服务器任务 for task in tasks: task.cancel() try: await task except asyncio.CancelledError: print(f"TCP server task on port {task.get_name()} cancelled") app = FastAPI(lifespan=startup_event) @app.websocket("/ws") async def websocket_endpoint(websocket: WebSocket): print("about to connect to websocket") await globals.websocket_manager.connect(websocket) print("websocket:", websocket) try: while True: # 保持WebSocket长连接(无需接收文本可忽略此步骤) await websocket.receive_text() except Exception as e: print(f"WebSocket Error: {e}") globals.websocket_manager.disconnect(websocket)
server.py
import asyncio import globals async def handle_client(reader: asyncio.StreamReader, writer: asyncio.StreamWriter): try: while True: data = await reader.read(1024) if not data: break # 处理非UTF-8数据,避免解码崩溃 try: text_data = data.decode('utf-8') await globals.websocket_manager.broadcast(text_data) except UnicodeDecodeError: print(f"Received non-UTF8 data: {data}") writer.close() await writer.wait_closed() except Exception as e: print(f"Client handler error: {e}") writer.close() await writer.wait_closed() async def start_tcp_server(port): print(f"starting tcp server on port {port}...") server = await asyncio.start_server(handle_client, '0.0.0.0', port) # 给任务命名,方便关闭时识别 asyncio.current_task().set_name(str(port)) async with server: await server.serve_forever()
globals.py
from websocket_manager import WebSocketManager data_storage = {} websocket_manager = WebSocketManager() # 移除冗余线程锁,异步环境无需线程同步锁
websocket_manager.py
from fastapi import WebSocket from typing import List class WebSocketManager: 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): # 移除连接前做存在性检查,避免KeyError if websocket in self.active_connections: self.active_connections.remove(websocket) async def broadcast(self, data: str): # 遍历连接列表副本,避免遍历中修改列表引发异常 for connection in self.active_connections.copy(): try: await connection.send_text(data) except Exception as e: print(f"Broadcast error: {e}") self.disconnect(connection)
额外验证步骤
- EC2安全组检查:确保服务器EC2的安全组开放8001、8002、8003端口的TCP入站规则,允许客户端EC2的IP段或指定IP访问。
- 端口占用检查:在服务器上执行
netstat -tulpn | grep -E "8001|8002|8003",确认端口未被其他进程占用。 - 启动日志验证:启动FastAPI后,检查控制台是否输出
starting tcp server on port 8001...等日志,确认TCP服务器已正常启动。
内容的提问来源于stack exchange,提问作者emptyfullstack
相关产品推荐
相关产品推荐

