Python asyncio websocket传输大文件时接收端文件体积异常过大问题求助
1. 接收端重复写入全量数据(导致文件超大的核心原因)
FileManager.py的chunk_receiver方法存在致命逻辑错误:每次收到分片都会将累计的整个received_file写入文件,而不是仅写入新收到的分片。比如收到第1个分片写1份数据,收到第2个分片写2份,收到第N个分片写N份,最终文件大小是原文件的N(N+1)/2倍,大文件场景下会直接膨胀到数GB。
修复代码:
async def chunk_receiver(self, binary): async with self.lock: self.received_file += binary # 仅写入本次新收到的分片,不是全量 self.file.write(binary) perc = ((len(self.received_file) * 100)/self.filesize) print("\rDownloading file: " + colored(str(round(perc, 2)) + "%", "magenta"), end='', flush=True)
2. 发送端分片偏移计算+二进制转码错误
chunk_sender方法不判断实际读取的字节数,一律按BUFFER_SIZE累加偏移量,最后一次分片读取长度不足BUFFER_SIZE时,偏移量会超出文件实际大小,后续读取会返回空字节,导致发送多余空数据。另外直接将二进制数据用UTF-8解码会导致非文本文件转码错误,数据损坏。
修复代码:
async def chunk_sender(self): async with self.lock: self.file.seek(self.sent) bytes_read = self.file.read(self.BUFFER_SIZE) if not bytes_read: return # 二进制数据需用base64编码才能安全放入JSON,避免utf-8解码错误 import base64 await self.ws.send(json.dumps({ "cmd": "send", "key": self.secret, "data": base64.b64encode(bytes_read).decode("utf-8") })) # 按实际读取长度累加偏移量 self.sent += len(bytes_read)
对应接收端解码逻辑修改(ClientManager.py):
elif cmd == "send": if "data" in deserialized: import base64 binary_chunk = base64.b64decode(deserialized["data"]) await self.FileManager.chunk_receiver(binary_chunk)
3. 类属性导致的数据污染
FileManager中的file、filesize、received_file、sent等可变变量定义为类属性,所有实例共用同一套变量,多次传输或者多用户传输时数据会互相叠加污染。需要把这些变量移到实例初始化方法中:
修复代码(FileManager.py):
class Manager(): BUFFER_SIZE = 8092 lock = asyncio.Lock() def __init__(self, secret=None, ws: websockets.WebSocketServerProtocol = None): self.ws = ws self.secret = secret # 所有可变变量移到实例属性 self.file = None self.filesize = None self.filename = None self.received_file = bytearray() self.sent = 0
4. 客户端FileManager实例污染
ClientManager中的FileManager定义为类属性,所有客户端实例共用同一个文件管理器,多任务传输时数据会乱套,修改为实例属性:
修复代码(ClientManager.py):
class Manager: SERVER_URL = None filename = None filesize = 0 secret = None # 删掉类级别的FileManager定义 def __init__(self, SERVER_URL, filename, filesize, secret): self.SERVER_URL = SERVER_URL self.filename = filename self.filesize = filesize self.secret = secret # 初始化实例专属的FileManager from classes import File self.FileManager = File.Manager() self.FileManager.secret = self.secret self.FileManager.filesize = self.filesize self.FileManager.filename = self.filename
5. 服务端连接注销逻辑错误
Server.py的unregister方法中del client仅删除临时变量,不会从clients列表中移除断开的连接,会导致重复转发消息给已断开的客户端,也可能出现重复转发给多个匹配的客户端。
修复代码:
async def unregister(self, ws: websockets.WebSocketServerProtocol) -> None: with self.clients_lock: self.clients = [client for client in self.clients if client["ws"] != ws] logging.info(f'{ws.remote_address[0]} 断开连接')
6. 可选优化:避免协程乱序问题
当前用asyncio.as_completed并发发送分片,虽然加了锁保证读取顺序,但如果后续去掉锁很容易出现分片乱序的问题,建议改为串行发送分片,避免接收端需要额外做排序逻辑:
async def start_sending(self): self.open_file(self.filename, "rb") spawn = math.ceil(self.filesize / self.BUFFER_SIZE) pbar = tqdm.tqdm(total=spawn, leave=True, mininterval=0) for _ in range(spawn): await self.chunk_sender() pbar.update() pbar.close() self.close_file()
内容的提问来源于stack exchange,提问作者Carlos Molero

