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

如何将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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 19:10:47