如何将Python流服务器中的数据传递回主任务?
解决Asyncio TCP服务器中多客户端数据的跨任务消费问题
问题背景
基于asyncio streams实现的TCP echo服务器代码如下:
import asyncio async def handle_echo(reader, writer): data = await reader.read(100) message = data.decode() addr = writer.get_extra_info('peername') print(f"Received {message!r} from {addr!r}") print(f"Send: {message!r}") writer.write(data) await writer.drain() print("Close the connection") writer.close() async def main(): server = await asyncio.start_server( handle_echo, '127.0.0.1', 8888) addrs = ', '.join(str(sock.getsockname()) for sock in server.sockets) print(f'Serving on {addrs}') async with server: await server.serve_forever() asyncio.run(main())
需要将客户端发送的数据作为生产者供其他任务消费,但直接传递asyncio.Queue到handle_echo失败——该函数仅能接收reader和writer两个参数,且模块级声明的队列不属于主任务创建的事件循环。需找到无需退回原生socket的解决方案。
解决方案
可以通过闭包或类来共享队列,确保队列归属当前事件循环,同时让handle_echo能访问到队列:
方法1:闭包封装队列
在main函数内创建队列,定义嵌套的handle_echo函数直接访问外层队列:
import asyncio async def main(): # 创建当前事件循环所属的队列 data_queue = asyncio.Queue() # 消费者任务:持续从队列取数据处理 async def consumer(): while True: message, addr = await data_queue.get() print(f"Consumer processing: {message!r} from {addr!r}") data_queue.task_done() # 嵌套的handle_echo,直接访问外层data_queue async def handle_echo(reader, writer): while True: data = await reader.read(100) if not data: break message = data.decode() addr = writer.get_extra_info('peername') print(f"Received {message!r} from {addr!r}") # 将数据放入队列,交给消费者处理 await data_queue.put((message, addr)) # 保留echo逻辑(按需保留) writer.write(data) await writer.drain() print(f"Close connection with {addr!r}") writer.close() await writer.wait_closed() # 启动后台消费者任务 asyncio.create_task(consumer()) # 启动TCP服务器 server = await asyncio.start_server(handle_echo, '127.0.0.1', 8888) addrs = ', '.join(str(sock.getsockname()) for sock in server.sockets) print(f'Serving on {addrs}') async with server: await server.serve_forever() asyncio.run(main())
方法2:类封装服务器与队列
用类管理队列、服务器和处理逻辑,结构更清晰:
import asyncio class TCPServerWithQueue: def __init__(self, host, port): self.host = host self.port = port self.data_queue = asyncio.Queue() self.server = None async def consumer(self): while True: message, addr = await self.data_queue.get() print(f"Consumer processing: {message!r} from {addr!r}") self.data_queue.task_done() async def handle_echo(self, reader, writer): while True: data = await reader.read(100) if not data: break message = data.decode() addr = writer.get_extra_info('peername') print(f"Received {message!r} from {addr!r}") await self.data_queue.put((message, addr)) writer.write(data) await writer.drain() print(f"Close connection with {addr!r}") writer.close() await writer.wait_closed() async def start(self): asyncio.create_task(self.consumer()) self.server = await asyncio.start_server(self.handle_echo, self.host, self.port) addrs = ', '.join(str(sock.getsockname()) for sock in self.server.sockets) print(f'Serving on {addrs}') async with self.server: await self.server.serve_forever() async def main(): server = TCPServerWithQueue('127.0.0.1', 8888) await server.start() asyncio.run(main())
关键说明
- 闭包/类的方式确保队列属于
main函数启动的事件循环,避免跨循环问题。 - 修改
handle_echo为循环读取,支持客户端持续发送数据(原代码仅读取一次就关闭连接)。 - 消费者任务通过
asyncio.create_task后台运行,持续处理队列中的数据。
内容的提问来源于stack exchange,提问作者Nick S.
相关产品推荐
相关产品推荐

