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

Python websockets客户端如何为await异步接收函数创建独立线程

需求完全可行,你之前用threading尝试失败的核心原因是async异步函数必须绑定对应线程的asyncio事件循环,不能直接在普通线程中运行。另外你的原有代码还存在几处可修复的问题:函数名拼写错误、await语句写在异步函数外部、async with作用域结束后连接会自动断开、全局变量存连接容易出现悬空引用。

以下是两种适配不同场景的实现方案:

方案一:仅用asyncio原生并发(无额外线程,推荐)

asyncio本身就是为单线程下实现IO并发设计的,你的需求完全不需要开线程,用asyncio的后台任务就可以实现:接收消息作为后台任务跑,主逻辑同时执行其他任务。

修正后的代码:

from typing import Dict
import websockets
import asyncio
import json

URL = "ws://localhost:你的服务端口"  # websocket地址前缀为ws/wss,不要填http地址
# 不要用全局变量存连接,避免作用域问题

def get_dict():
    # 补全你自己的初始化参数逻辑
    return {"type": "init"}

async def receive_msg(ws):
    # 直接传入连接对象,不依赖全局变量
    while True:
        msg = await ws.recv()
        print(f"Server: {msg}")

async def listen():
    input("按回车开始连接")
    async with websockets.connect(URL) as ws:
        # 发送初始化消息
        msg_initial: Dict[str, str] = get_dict()
        await ws.send(json.dumps(msg_initial))
        
        # 创建后台任务运行接收逻辑,不会阻塞当前协程
        receive_task = asyncio.create_task(receive_msg(ws))
        
        # 这里写你原本要执行的其他业务逻辑
        print("主逻辑正常运行,接收消息已在后台启动")
        for i in range(10):
            print(f"主逻辑执行任务{i}")
            await asyncio.sleep(1)  # 模拟你的其他异步操作
        
        # 如需长期保持连接,等待接收任务结束即可
        await receive_task

if __name__ == "__main__":
    asyncio.run(listen())

说明:

  • asyncio.create_task会把接收消息的协程注册到当前事件循环作为后台任务,只要事件循环在运行就会自动调度,完全不需要额外开线程
  • 连接的生命周期和async with绑定,出了作用域会自动断开,不会出现资源泄露

方案二:必须用独立线程实现(适配原有同步逻辑场景)

如果你原有主线程的逻辑都是同步代码,没法改成async协程,那可以在独立子线程里单独跑一个asyncio事件循环处理websocket接收:

代码示例:

from typing import Dict
import websockets
import asyncio
import json
import threading
import queue

URL = "ws://localhost:你的服务端口"
# 线程安全的通信队列,主线程需要发消息时直接往队列里塞数据即可
send_queue = queue.Queue()

def get_dict():
    return {"type": "init"}

# 子线程入口函数,内部单独运行事件循环
def websocket_thread():
    async def ws_task():
        async with websockets.connect(URL) as ws:
            msg_initial: Dict[str, str] = get_dict()
            await ws.send(json.dumps(msg_initial))
            
            # 两个后台任务:分别处理收消息、发主线程传来的消息
            async def recv_task():
                while True:
                    msg = await ws.recv()
                    print(f"Server: {msg}")
            
            async def send_task():
                while True:
                    if not send_queue.empty():
                        msg = send_queue.get()
                        await ws.send(json.dumps(msg))
                    await asyncio.sleep(0.1)
            
            # 两个任务并行运行
            await asyncio.gather(recv_task(), send_task())
    
    # 子线程单独创建专属的事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    loop.run_until_complete(ws_task())

if __name__ == "__main__":
    # 启动websocket子线程,daemon=True表示主程序退出时自动销毁子线程
    t = threading.Thread(target=websocket_thread, daemon=True)
    t.start()
    
    # 这里是你主线程的原有同步逻辑,完全不受影响
    input("按回车停止程序\n")
    # 主线程发消息示例:send_queue.put({"content": "来自主线程的消息"})

内容的提问来源于stack exchange,提问作者Gui Reis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 10:24:06