You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

部署在GCP App Engine的WebSocket服务器全局连接字典异常排查

WebSocket服务器问题解答与优化方案

1. 全局字典跟踪连接是否为最佳实践?

用全局字典跟踪单实例内的活跃WebSocket连接是可行的临时方案,但存在明确限制:

  • 仅适用于单进程/单实例部署:如果GCP App Engine启用自动扩缩容(多实例),全局字典无法跨实例共享,此时需改用Redis等分布式存储统一管理连接映射。
  • 必须保证协程安全:所有对字典的读写操作都要通过异步锁(asyncio.Lock)包裹,避免并发修改导致的数据结构损坏。

2. asyncio.Task与ensure_future替换后连接字典为空的原因?

两者在Python 3.7+功能等价,问题不在任务创建方式,而是服务器代码存在并发安全漏洞:
服务器的ping分支直接修改全局字典connections[sno] = websocket,未通过async with lock保护。当客户端频繁发送ping(比如你的ping任务每10秒触发一次),多个协程并发修改字典会导致数据结构损坏,极端情况下可能直接清空字典。

客户端任务创建方式的差异可能间接提升了消息发送频率,进一步加剧了并发冲突的概率。

3. 连接顺序导致字典失效的问题根源?

核心问题是全局字典的并发修改未完全加锁,以及日志打印时机的误导:

  • 移动端发送消息时,print(connections)在锁外执行,此时字典可能正被其他协程修改,打印的是旧状态;锁内的检查才是准确的,但你的代码未输出锁内的字典状态。
  • 客户端连接后,ping消息触发的无锁字典修改可能覆盖或损坏之前的连接条目,导致移动端后续请求找不到对应记录。

另外需检查移动端发送的serialno是否与客户端完全一致(如大小写、空格),避免因字符串不匹配导致的查找失败。

4. 代码不良实践与优化方案

不良实践清单

  • 部分字典操作未加异步锁(ping分支),存在并发安全风险
  • finally块仅在最后一条消息是client类型时才删除连接,若最后一条是ping,会导致无效连接残留
  • 用print代替logging输出调试信息,生产环境无法统一管理日志
  • 捕获所有Exception,掩盖了WebSocketDisconnect等特定异常的处理逻辑
  • 自定义ping/pong逻辑,未利用WebSocket协议原生的心跳机制
  • 客户端未在连接关闭时取消ping任务,存在资源泄漏风险

优化后的服务器代码

import asyncio
import logging
import uvicorn
import json
from fastapi import FastAPI, WebSocket, WebSocketDisconnect

# Configure logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

app = FastAPI()

# Dictionary to store connections
connections = {}
lock = asyncio.Lock()  # Async lock for thread-safe access

@app.websocket("/ask")
async def websocket_endpoint(websocket: WebSocket):
    await websocket.accept()
    serial_number = None
    connection_type = None  # Track if this is a client/mobile connection
    try:
        while True:
            data = await websocket.receive_text()
            logger.info(f"Raw data received: {data}")
            data_dict = json.loads(data)
            logger.info(f"Parsed data: {data_dict}")
            
            # Validate required fields
            if "type" not in data_dict or "serialno" not in data_dict:
                logger.warning("Missing required fields in message")
                continue
            
            msg_type = data_dict["type"]
            sno = data_dict["serialno"]
            
            async with lock:
                current_connections = connections.copy()
            logger.info(f"Current connections: {current_connections}")
            
            if msg_type == "client":
                serial_number = sno
                connection_type = "client"
                async with lock:
                    connections[serial_number] = websocket
                logger.info(f"Client {serial_number} connected")
            elif msg_type == "mobile":
                async with lock:
                    target_ws = connections.get(sno)
                if target_ws and isinstance(target_ws, WebSocket):
                    try:
                        await target_ws.send_text("guava")
                        logger.info(f"Message sent to client {sno}")
                    except Exception as e:
                        logger.error(f"Failed to send message to {sno}: {str(e)}")
                        # Remove invalid connection
                        async with lock:
                            if sno in connections:
                                del connections[sno]
                else:
                    logger.warning(f"No active client connection for serial number: {sno}")
            elif msg_type == "ping":
                # Update connection entry (with lock) and send pong
                async with lock:
                    connections[sno] = websocket
                await websocket.send_text("pong")
                logger.info(f"Pong sent to {sno}")
            else:
                logger.warning(f"Unknown message type: {msg_type}")

    except WebSocketDisconnect:
        logger.info(f"Connection closed for {serial_number} ({connection_type})")
    except Exception as e:
        logger.error(f"Error handling connection {serial_number}: {str(e)}", exc_info=True)
    
    finally:
        # Clean up connection regardless of last message type
        if serial_number and connection_type == "client":
            async with lock:
                if serial_number in connections:
                    del connections[serial_number]
                    logger.info(f"Client {serial_number} removed from connections")

if __name__ == "__main__":
    uvicorn.run(app, host="0.0.0.0", port=8000)

优化后的客户端代码

import asyncio
import websockets
import json

async def websocket_client():
    uri = "ws://rasd-8bd8c.uc.r.appspot.com/ask"
    
    # Use WebSocket native ping mechanism instead of custom task
    async with websockets.connect(uri, ping_interval=60, ping_timeout=120) as websocket:
        # Initial client registration
        register_msg = json.dumps({
            "type": "client",
            "serialno": "123",
            "msg": "wow"
        })
        await websocket.send(register_msg)
        print(f"Sent registration: {register_msg}")

        # Continuously listen for server messages
        try:
            while True:
                response = await websocket.recv()
                print(f"Received from server: {response}")
        except websockets.ConnectionClosed:
            print("Connection closed by server")

# Run client
asyncio.run(websocket_client())

关键优化点说明

  1. 全链路协程安全:所有对connections的读写操作都通过async with lock包裹,彻底避免并发冲突
  2. 连接清理逻辑修复:用connection_type独立跟踪连接类型,不再依赖最后一条消息的类型,确保客户端断开时必清理无效条目
  3. 原生心跳机制:客户端使用websockets库内置的ping_interval和ping_timeout替代自定义ping任务,减少代码复杂度
  4. 日志规范:用logging替代print,便于生产环境统一管理和排查问题
  5. 异常精细化处理:单独捕获WebSocketDisconnect,区分正常断开与异常错误,便于定位问题

内容的提问来源于stack exchange,提问作者Fawwaz Ali

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 18:30:58