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

基于asyncio的TCP服务器如何实现向客户端按需发送数据?

修复Asyncio TCP服务器主动向客户端发送数据的问题

我用asyncio实现了一台TCP服务器,它能定期(比如每30秒)接收客户端的数据。现在想实现向客户端按需发送数据的功能,用来给客户端设置特定参数,所以需要服务器主动给客户端发消息。自己尝试实现但没成功,原代码如下:

import asyncio
from protocol import Message  # Import your Message class from protocol.py

class TCPServer:

    def __init__(self, host, port):
        self.host = host
        self.port = port
        self.server = None
        self.clients = {}  # Dictionary to store connected clients
        self.input_queue = asyncio.Queue()  # Queue for console input

    async def handle_client(self, reader, writer):
       message = Message()  # Create an instance of your Message class
       addr = writer.get_extra_info('peername')

       # Convert the client address to a string for consistency
       client_address = f"{addr[0]}:{addr[1]}"

       try:
           self.clients[client_address] = writer  # Store the writer for this client

           while True:
               data = await reader.read(1024)
               if not data:
                   break

               # Assuming your Message class has a `raw_data` method
               response = message.raw_data(data)

               print(f"Received {data!r} from {addr!r}")

               print(f"Send: {response!r}")
               writer.write(response.encode())
               await writer.drain()

    except ConnectionResetError as e:
        print(f"ConnectionResetError: {e}")
    except Exception as e:
        print(f"Error: {e}")
    finally:
        print(f"Client {addr!r} closed the connection")
        writer.close()
        await writer.wait_closed()
        del self.clients[client_address]  # Remove the client from the list when they disconnect

    async def send_to_client(self, client_address, message):
          if client_address in self.clients:
             writer = self.clients[client_address]
             writer.write(message.encode())
             await writer.drain()
         else:
             print("Client not found or not connected.")

   async def start(self):
        self.server = await asyncio.start_server(
        self.handle_client, self.host, self.port)

        addr = self.server.sockets[0].getsockname()
        print(f'Serving on {addr}')

        async with self.server:
            await self.server.serve_forever()

  async def stop(self):
        if self.server:
           self.server.close()
           await self.server.wait_closed()
           print("Server stopped")

  async def console_input(self):
       while True:
           client_address = await self.input_queue.get()
           hex_data = input("Enter the hex data to send: ")
           try:
               byte_data = bytes.fromhex(hex_data)
               await self.send_to_client(client_address, byte_data)
          except ValueError:
               print("Invalid hex data format.")

 if __name__ == '__main__':
     host = '127.0.0.1'  # Change this to the desired host
     port = 8888  # Change this to the desired port

     server = TCPServer(host, port)

    async def main():
        try:
          server_task = asyncio.create_task(server.start())
          console_task = asyncio.create_task(server.console_input())

          # Add client address to the input queue for testing
          await asyncio.sleep(2)  # Wait for the server to start (adjust as needed)
          client_address = '127.0.0.1:8888'
          server.input_queue.put_nowait(client_address)

          await asyncio.gather(server_task, console_task)
      except KeyboardInterrupt:
          print("KeyboardInterrupt: Server shutting down...")
          await server.stop()

  asyncio.run(main())

问题分析

原代码存在几个关键问题导致功能失效:

  • 缩进错误:多个方法(handle_client的try块、send_to_client、start等)缩进不一致,直接触发语法错误。
  • 客户端地址错误:测试时硬编码的127.0.0.1:8888是服务器监听地址,不是客户端的实际连接地址(客户端连接后会用随机端口,格式为127.0.0.1:xxxx),导致找不到目标客户端。
  • 同步input阻塞事件循环:console_input里用了同步input(),会卡住整个asyncio事件循环,服务器无法处理客户端连接和数据。
  • 消息编码错误:send_to_client中对传入的bytes类型消息再次调用encode(),会抛出类型错误。

修复后的代码

import asyncio
from protocol import Message  # 从protocol.py导入Message类

class TCPServer:
    def __init__(self, host, port):
        self.host = host
        self.port = port
        self.server = None
        self.clients = {}  # 存储已连接客户端的字典
        self.input_queue = asyncio.Queue()  # 控制台输入队列

    async def handle_client(self, reader, writer):
        message = Message()  # 创建Message实例
        addr = writer.get_extra_info('peername')
        client_address = f"{addr[0]}:{addr[1]}"
        
        print(f"客户端 {client_address} 已连接")
        self.clients[client_address] = writer  # 保存客户端的writer对象

        try:
            while True:
                data = await reader.read(1024)
                if not data:
                    break

                # 假设Message类有raw_data方法处理收到的数据
                response = message.raw_data(data)
                print(f"从 {client_address} 收到: {data!r}")
                print(f"回复: {response!r}")
                writer.write(response.encode())
                await writer.drain()

        except ConnectionResetError as e:
            print(f"客户端 {client_address} 连接重置: {e}")
        except Exception as e:
            print(f"处理客户端 {client_address} 时出错: {e}")
        finally:
            print(f"客户端 {client_address} 断开连接")
            writer.close()
            await writer.wait_closed()
            if client_address in self.clients:
                del self.clients[client_address]  # 从字典中移除客户端

    async def send_to_client(self, client_address, message):
        if client_address not in self.clients:
            print(f"未找到客户端 {client_address} 或客户端未连接")
            return
        
        writer = self.clients[client_address]
        try:
            # 判断消息类型,bytes直接发送,字符串则编码
            if isinstance(message, str):
                writer.write(message.encode())
            else:
                writer.write(message)
            await writer.drain()
            print(f"已向 {client_address} 发送: {message!r}")
        except Exception as e:
            print(f"向 {client_address} 发送消息失败: {e}")

    async def start(self):
        self.server = await asyncio.start_server(
            self.handle_client, self.host, self.port
        )
        addr = self.server.sockets[0].getsockname()
        print(f"服务器运行在 {addr}")
        
        async with self.server:
            await self.server.serve_forever()

    async def stop(self):
        if self.server:
            self.server.close()
            await self.server.wait_closed()
            print("服务器已停止")

    async def console_input(self):
        loop = asyncio.get_event_loop()
        while True:
            # 使用线程池执行同步input,避免阻塞事件循环
            client_address = await loop.run_in_executor(
                None, input, "请输入要发送的客户端地址(格式: ip:port): "
            )
            hex_data = await loop.run_in_executor(
                None, input, "请输入要发送的十六进制数据: "
            )
            try:
                byte_data = bytes.fromhex(hex_data)
                await self.send_to_client(client_address, byte_data)
            except ValueError:
                print("无效的十六进制格式,请重新输入")

if __name__ == '__main__':
    host = '127.0.0.1'
    port = 8888

    server = TCPServer(host, port)

    async def main():
        try:
            server_task = asyncio.create_task(server.start())
            console_task = asyncio.create_task(server.console_input())
            await asyncio.gather(server_task, console_task)
        except KeyboardInterrupt:
            print("\n收到中断信号,正在关闭服务器...")
            await server.stop()

    asyncio.run(main())

修复说明

  1. 修正缩进:统一所有代码块的缩进,确保符合Python语法规范。
  2. 替换同步input:用loop.run_in_executor在单独线程中处理控制台输入,避免阻塞asyncio事件循环,保证服务器能正常处理客户端连接和数据。
  3. 修复客户端地址问题:不再硬编码客户端地址,改为让用户手动输入;客户端连接时会打印其地址,方便用户获取正确的目标地址。
  4. 优化消息发送逻辑:增加消息类型判断,bytes直接发送,字符串自动编码,避免重复encode导致的错误。
  5. 完善异常处理:在消息发送和客户端处理流程中增加更详细的错误提示,便于排查问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 15:59:52