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

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())

问题定位与修复

核心问题

  1. 客户端阻塞:同步input()卡住事件循环,且客户端仅在发送命令后读取一次响应,无法接收服务器主动推送的定时器消息
  2. 服务器逻辑缺陷:
    • 定时器实例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())

修复说明

  1. 客户端:

    • 使用asyncio.to_thread()包装input(),避免同步输入阻塞事件循环
    • 新增独立异步任务持续读取服务器消息,处理定时器的主动推送
    • 简化退出逻辑,移除重复发送quit的冗余代码
  2. 服务器:

    • 保存Timer实例而非布尔值,解决stop命令无法访问定时器的问题
    • 处理消息时使用strip()去除换行和空格,确保命令匹配准确
    • 在Timer的_tick方法中捕获CancelledError,优雅处理任务取消
    • 完善连接断开时的资源清理逻辑,增加未知命令提示

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 20:35:29