Python asyncio客户端服务端通信异常:Timer命令阻塞无消息返回
问题描述
使用Python asyncio实现客户端/服务器架构,服务器为回显服务器,支持两类命令:
start:启动定时器,每秒向客户端发送消息并在服务器控制台打印时间差stop:停止定时器quit:关闭连接
实际运行问题:启动服务和客户端后,触发start命令,定时器消息无法发送到客户端,且客户端与服务器均出现阻塞。
原服务器代码
import asyncio import time HOST = "127.0.0.1" PORT = 9999 class Timer(object): '''Simple timer class that can be started and stopped.''' def __init__(self, writer: asyncio.StreamWriter, name = None, interval = 1) -> None: self.name = name self.interval = interval self.writer = writer async def _tick(self) -> None: while True: await asyncio.sleep(self.interval) delta = time.time() - self._init_time self.writer.write(f"Timer {delta} ticked\n".encode()) self.writer.drain() print("Delta time: ", delta) async def start(self) -> None: self._init_time = time.time() self.task = asyncio.create_task(self._tick()) async def stop(self) -> None: self.task.cancel() print("Delta time: ", time.time() - self._init_time) async def msg_handler(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: '''Handle the echo protocol.''' # timer task that the client can start: timer_task = False try: while True: data = await reader.read(1024) # Read 256 bytes from the reader. Size of the message msg = data.decode() # Decode the message addr, port = writer.get_extra_info("peername") # Get the address of the client print(f"Received {msg!r} from {addr}:{port!r}") send_message = "Message received: " + msg writer.write(send_message.encode()) # Echo the data back to the client await writer.drain() # This will wait until everything is clear to move to the next thing. if data == b"quit" and timer_task is True: # cancel the timer_task (if any) if timer_task: timer_task.cancel() await timer_task writer.close() # Close the connection await writer.wait_closed() # Wait for the connection to close elif data == b"quit" and timer_task is False: writer.close() # Close the connection await writer.wait_closed() # Wait for the connection to close elif data == b"start" and timer_task is False: print("Starting timer") t = Timer(writer) timer_task = True await t.start() elif data == b"stop" and timer_task is True: print("Stopping timer") await t.stop() timer_task = False except ConnectionResetError: print("Client disconnected") async def run_server() -> None: # Our awaitable callable. # This callable is ran when the server recieves some data server = await asyncio.start_server(msg_handler, HOST, PORT) async with server: await server.serve_forever() if __name__ == "__main__": loop = asyncio.new_event_loop() # new_event_loop() is for python 3.10. For older versions, use get_event_loop() loop.run_until_complete(run_server())
原客户端代码
import asyncio HOST = '127.0.0.1' PORT = 9999 async def run_client() -> None: # It's a coroutine. It will wait until the connection is established reader, writer = await asyncio.open_connection(HOST, PORT) while True: message = input('Enter a message: ') writer.write(message.encode()) await writer.drain() data = await reader.read(1024) if not data: raise Exception('Socket not communicating with the client') print(f"Received {data.decode()!r}") if (message == 'quit'): writer.write(b"quit") writer.close() await writer.wait_closed() exit(2) # break # Don't know if this is necessary if __name__ == '__main__': loop = asyncio.new_event_loop() loop.run_until_complete(run_client())
问题定位与修复
核心问题
- 客户端阻塞:同步
input()卡住事件循环,且客户端仅在发送命令后读取一次响应,无法接收服务器主动推送的定时器消息 - 服务器逻辑缺陷:
- 定时器实例
t为局部变量,stop命令无法访问 timer_task被赋值为布尔值,无法实际控制定时器任务- 命令匹配未处理输入换行,导致
quit等命令判断失效
- 定时器实例
修复后的代码
修复后客户端代码
import asyncio HOST = '127.0.0.1' PORT = 9999 async def read_server_messages(reader): """持续读取服务器推送的消息""" while True: data = await reader.read(1024) if not data: print("服务器连接已关闭") break print(f"收到服务器消息: {data.decode()!r}") async def run_client() -> None: reader, writer = await asyncio.open_connection(HOST, PORT) # 启动独立任务监听服务器消息,避免阻塞输入流程 asyncio.create_task(read_server_messages(reader)) while True: # 用异步线程包装同步input,避免阻塞事件循环 message = await asyncio.to_thread(input, '输入命令(start/stop/quit): ') if not message: continue writer.write(message.encode()) await writer.drain() if message.strip() == 'quit': writer.close() await writer.wait_closed() break if __name__ == '__main__': asyncio.run(run_client())
修复后服务器代码
import asyncio import time HOST = "127.0.0.1" PORT = 9999 class Timer(object): '''Simple timer class that can be started and stopped.''' def __init__(self, writer: asyncio.StreamWriter, name=None, interval=1) -> None: self.name = name self.interval = interval self.writer = writer self.task = None self._init_time = None async def _tick(self) -> None: try: while True: await asyncio.sleep(self.interval) delta = time.time() - self._init_time msg = f"Timer {delta:.2f} ticked\n".encode() self.writer.write(msg) await self.writer.drain() print(f"Delta time: {delta:.2f}") except asyncio.CancelledError: print(f"Timer stopped, total duration: {time.time() - self._init_time:.2f}") async def start(self) -> None: if self.task is None or self.task.done(): self._init_time = time.time() self.task = asyncio.create_task(self._tick()) async def stop(self) -> None: if self.task and not self.task.done(): self.task.cancel() await self.task self.task = None async def msg_handler(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: '''Handle client commands and echo messages.''' timer = None addr, port = writer.get_extra_info("peername") print(f"客户端 {addr}:{port} 已连接") try: while True: data = await reader.read(1024) if not data: print(f"客户端 {addr}:{port} 断开连接") break msg = data.decode().strip() # 去除换行和空格,避免命令匹配失败 print(f"收到 {addr}:{port} 的命令: {msg!r}") # 回显命令 send_message = f"命令已接收: {msg}\n" writer.write(send_message.encode()) await writer.drain() if msg == "quit": print(f"客户端 {addr}:{port} 请求断开连接") if timer: await timer.stop() break elif msg == "start": if not timer or not timer.task or timer.task.done(): print(f"为 {addr}:{port} 启动定时器") timer = Timer(writer) await timer.start() else: writer.write(b"定时器已在运行\n") await writer.drain() elif msg == "stop": if timer and timer.task and not timer.task.done(): print(f"为 {addr}:{port} 停止定时器") await timer.stop() else: writer.write(b"定时器未运行\n") await writer.drain() else: writer.write(b"未知命令,支持的命令: start/stop/quit\n") await writer.drain() except ConnectionResetError: print(f"客户端 {addr}:{port} 意外断开") finally: # 退出时清理定时器资源 if timer: await timer.stop() writer.close() await writer.wait_closed() print(f"与 {addr}:{port} 的连接已关闭") async def run_server() -> None: server = await asyncio.start_server(msg_handler, HOST, PORT) addr = server.sockets[0].getsockname() print(f"服务器启动,监听 {addr}") async with server: await server.serve_forever() if __name__ == "__main__": asyncio.run(run_server())
修复说明
客户端:
- 使用
asyncio.to_thread()包装input(),避免同步输入阻塞事件循环 - 新增独立异步任务持续读取服务器消息,处理定时器的主动推送
- 简化退出逻辑,移除重复发送
quit的冗余代码
- 使用
服务器:
- 保存Timer实例而非布尔值,解决
stop命令无法访问定时器的问题 - 处理消息时使用
strip()去除换行和空格,确保命令匹配准确 - 在Timer的
_tick方法中捕获CancelledError,优雅处理任务取消 - 完善连接断开时的资源清理逻辑,增加未知命令提示
- 保存Timer实例而非布尔值,解决
内容的提问来源于stack exchange,提问作者nunodsousa
相关产品推荐
相关产品推荐

