asyncio服务器循环向客户端写入数据仅首次成功后续缓存问题
问题根源
- 你在
write_loop里使用了阻塞式的同步time.sleep,并且整个循环是同步执行的,全程占用了asyncio事件循环所在的线程,导致事件循环没有机会执行「将内核缓冲区的数据实际发送到网络」的后台任务,所有transport.write写入的内容都停留在进程缓存中,直到进程退出时才会一次性刷新。 - 同步的
while死循环不会让出CPU给事件循环,asyncio所有IO调度逻辑都无法运行,自然也不会发送后续数据。
修复方案
- 将
write_loop改为异步函数,用asyncio.sleep替代同步的time.sleep,每次休眠时主动让出CPU给事件循环处理IO - 在
data_received回调中通过asyncio.create_task提交异步写任务,不要同步调用死循环函数 - 连接关闭时主动取消未完成的写任务,避免出现异常
修复后可直接运行的完整代码如下:
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
相关产品推荐
相关产品推荐

