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

asyncio服务器循环向客户端写入数据仅首次成功后续缓存问题

问题根源

  • 你在write_loop里使用了阻塞式的同步time.sleep,并且整个循环是同步执行的,全程占用了asyncio事件循环所在的线程,导致事件循环没有机会执行「将内核缓冲区的数据实际发送到网络」的后台任务,所有transport.write写入的内容都停留在进程缓存中,直到进程退出时才会一次性刷新。
  • 同步的while死循环不会让出CPU给事件循环,asyncio所有IO调度逻辑都无法运行,自然也不会发送后续数据。

修复方案

  1. 将write_loop改为异步函数,用asyncio.sleep替代同步的time.sleep,每次休眠时主动让出CPU给事件循环处理IO
  2. 在data_received回调中通过asyncio.create_task提交异步写任务,不要同步调用死循环函数
  3. 连接关闭时主动取消未完成的写任务,避免出现异常

修复后可直接运行的完整代码如下:

import asyncio
import socket
import sys

class server_protocol(asyncio.Protocol):
    def __init__(self):
        self.connection_closed = True
        self.write_task = None

    def connection_made(self, transport):
        peername = transport.get_extra_info('peername')
        print('Connection from {}'.format(peername))
        self.transport = transport
        self.write_counter = 0
        self.connection_closed = False
        
    async def write_loop(self):
        while not self.connection_closed:
            self.write_counter += 1
            print(self.write_counter, 'Writing bytes, Buffer Size:', self.transport.get_write_buffer_size())
            self.transport.write(b'AhXTMkHrJHdExaKBLmkvTnvRduENcusjnRnrBAHtjnUMtdjxsnKgDRtpDMjncFczrqwjrSrVNwxtBSmmLJnFAfkgbDwEuBAcdVMCLVeMuSXfxyYRdaNvvEhEFnGWNNtk')
            # 异步休眠让出事件循环调度权限,触发缓存数据的网络发送逻辑
            await asyncio.sleep(1)
        self.transport.close()
        
    def data_received(self, data):
        # 提交异步任务,不阻塞当前回调线程
        self.write_task = asyncio.create_task(self.write_loop())
        
    def connection_lost(self, exc):
        print("Connection lost")
        self.connection_closed = True
        if self.write_task:
            self.write_task.cancel()
        
    def eof_received(self):
        print("EOF received")
        self.connection_closed = True
        if self.write_task:
            self.write_task.cancel()
        
def client(port):
    sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_address = ('localhost', port)
    print('connecting to {}:{}'.format(*server_address))
    sock.connect(server_address)
    try:
        message = b'start'
        print('sending {!r}'.format(message))
        sock.sendall(message)
        while True:
            print('Receiving')
            data = sock.recv(256)
            print('Received {!r}'.format(data))
            if len(data) == 0:
                break
    finally:
        print('closing socket')
        sock.close()
        
async def main(type, port):
    if type == 'server':
        loop = asyncio.get_running_loop()
        server = await loop.create_server(lambda: server_protocol(), 'localhost', port)
        async with server:
            await server.serve_forever()
    elif type == 'client':
        client(port)
        
if __name__ == '__main__':
    if len(sys.argv) == 3:
        asyncio.run(main(sys.argv[1], int(sys.argv[2])))
    else:
        print('provide argument \'server\' or \'client\' and port')

修改后运行即可看到客户端每秒都能收到服务端发送的字节数据,不会再出现缓存累积的问题。核心原则是asyncio的事件循环线程中绝对不能运行长时间的同步阻塞代码,所有休眠、IO操作都要使用异步版本,主动让出CPU给事件循环调度。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 06:15:10