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

Python-Websockets异步发送正常但无法接收消息问题排查

问题描述

使用Python的websockets包实现全双工WebSocket客户端,服务器仅做回声返回(收到消息后原样返回)。目前客户端能正常发送消息,但完全无法接收消息,服务器已确认完成消息收发,排除服务器问题。

该代码用于缓冲外部系统音频并发送至其他服务,同时需随时接收会话相关消息。环境:Python 3.9.15、websockets==10.4。

客户端代码

import asyncio

import websockets

sent = []
received = []

URL = "ws://localhost:8001"


async def update_sent(message):
    with open("sent.txt", "a+") as f:
        print(message, file=f)
    sent.append(message)
    return 0


async def update_received(message):
    with open("recv.txt", "a+") as f:
        print(message, file=f)
        received.append(message)
    return 0


async def sending_handler(websocket):
    while True:
        try:
            message = input("send message:")
            await websocket.send(message)
            await update_sent(message)
        except Exception as e:
            print("Sender: connection closed due to Exception", e)
            break


async def receive_handler(websocket):
    while True:
        try:
            message = await websocket.recv()
            await update_received(message)
        except Exception as e:
            print("Receiver: connection closed due to Exception", e)
            break


async def full_duplex_handler(websocket):
    receiving_task = asyncio.create_task(receive_handler(websocket))
    sending_task = asyncio.create_task(sending_handler(websocket))

    done, pending = await asyncio.wait([receiving_task, sending_task],
                                       return_when=asyncio.FIRST_COMPLETED)
                                       # return_when=asyncio.FIRST_EXCEPTION)
    for task in pending:
        print(task)
        task.cancel()


async def gather_handler(websocket):
    await asyncio.gather(
        sending_handler(websocket),
        receive_handler(websocket),
    )


# using asyncio.wait
async def main_1(url=URL):
    async with websockets.connect(url) as websocket:
        await full_duplex_handler(websocket)


# using asyncio.gather
# async def main_2(url=URL):
#     async with websockets.connect(url) as websocket:
#         await gather_handler(websocket)


if __name__ == "__main__":
    asyncio.run(main_1())
    # asyncio.run(main_2())

服务器代码

import asyncio

import websockets

msgs = []
sent = []


async def handle_send(websocket, message):
    await websocket.send(message)
    msgs.append(message)


async def handle_recv(websocket):
    message = await websocket.recv()
    sent.append(message)
    return f"echo {message}"


async def handler(websocket):
    while True:
        try:
            message = await handle_recv(websocket)
            await handle_send(websocket, message)
        except Exception as e:
            print(e)
            print(msgs)
            print(sent)
            break


async def main():
    async with websockets.serve(handler, "localhost", 8001):
        await asyncio.Future()


if __name__ == "__main__":
    print("starting the server now")
    asyncio.run(main())

预期:发送和接收的消息均写入对应文件,但目前仅发送消息被正常处理。


问题原因与解决办法

核心问题

客户端无法接收消息的根本原因是**input()是同步阻塞函数**,它会卡住整个asyncio事件循环。当sending_handler执行到input("send message:")时,事件循环被完全阻塞,无法切换到receive_handler任务处理WebSocket的接收操作,导致服务器返回的回声消息无法被客户端处理。

解决方案

需要替换同步的input()为异步输入方式,或用线程处理同步输入,避免阻塞事件循环。以下提供两种可行方案:

方案1:使用异步输入(推荐)

利用asyncio的StreamReader实现异步读取控制台输入,不会阻塞事件循环。修改后的sending_handler如下:

import sys  # 需要导入sys模块

async def sending_handler(websocket):
    # 创建异步输入流
    loop = asyncio.get_running_loop()
    reader = asyncio.StreamReader()
    protocol = asyncio.StreamReaderProtocol(reader)
    await loop.connect_read_pipe(lambda: protocol, sys.stdin)
    
    while True:
        try:
            # 异步读取输入
            message = await reader.readline()
            message = message.decode().strip()  # 转字符串并去除换行
            if not message:
                continue
            await websocket.send(message)
            await update_sent(message)
        except Exception as e:
            print("Sender: connection closed due to Exception", e)
            break

方案2:用线程处理同步输入

把input()放到单独线程中执行,通过队列传递输入内容到异步任务,避免阻塞事件循环:

import asyncio
import websockets
import sys
from threading import Thread
from queue import Queue

sent = []
received = []
URL = "ws://localhost:8001"
input_queue = Queue()

# 线程函数:处理同步输入
def input_thread():
    while True:
        msg = input("send message:")
        input_queue.put(msg)

async def update_sent(message):
    with open("sent.txt", "a+") as f:
        print(message, file=f)
    sent.append(message)
    return 0

async def update_received(message):
    with open("recv.txt", "a+") as f:
        print(message, file=f)
        received.append(message)
    return 0

async def sending_handler(websocket):
    # 启动输入线程
    Thread(target=input_thread, daemon=True).start()
    
    while True:
        try:
            # 异步等待队列中的输入
            message = await asyncio.to_thread(input_queue.get)
            await websocket.send(message)
            await update_sent(message)
        except Exception as e:
            print("Sender: connection closed due to Exception", e)
            break

# 其余receive_handler、full_duplex_handler、main_1等函数保持不变

额外优化

full_duplex_handler中使用asyncio.wait并设置return_when=asyncio.FIRST_COMPLETED,会在任意一个任务完成后取消另一个任务。如果希望两个任务一直运行到连接关闭,推荐改用asyncio.gather(即你的main_2函数),它会等待所有任务完成并更好地处理异常。修改主函数如下:

if __name__ == "__main__":
    # asyncio.run(main_1())
    asyncio.run(main_2())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 11:40:23