Python基于aiortc的WebRTC数据通道大文件传输吞吐量被限制在~110 Mb/s的问题咨询
Python基于aiortc的WebRTC数据通道大文件传输吞吐量被限制在~110 Mb/s的问题咨询
我正在用Python结合aiortc和WebRTC数据通道开发一个P2P文件传输工具,主打高效处理大文件——具体做法是从磁盘读取64MB的块,再拆分成256KB的小chunk通过数据通道发送。
我的当前配置
- 磁盘读取块大小:64 MB
- 网络传输chunk大小:256 KB
- 最大缓冲量(MAX_BUFFERED_AMOUNT):32 * 256 KB
- 使用两个数据通道:一个用于传输文件数据,一个用于控制消息
- 用asyncio队列来实现chunk的流水线处理
遇到的问题
哪怕是在同一台机器上(localhost)进行传输,我能达到的最大吞吐量也只有~110 Mb/s。观察发现,数据通道的bufferedAmount始终卡在MAX_BUFFERED_AMOUNT的上限,发送循环一直在等待缓冲区空间,这明显限制了吞吐量。
核心代码实现
import asyncio import aiofiles import os import json # 假设以下是项目自定义模块/实现 class ControlMessage: @staticmethod def create_json(msg_type, data): return json.dumps({"type": msg_type, "data": data}) @staticmethod def from_json(message): obj = json.loads(message) return type('ControlMsg', (), {"msg_type": obj["type"], "data": obj["data"]})() class Progress: def __init__(self, total): self.total = total self.sent = 0 def update(self, chunk): self.sent += len(chunk) def compute_hash(file_path): # 假设有文件哈希计算逻辑 return "dummy_hash" # 定义常量 MAX_BUFFERED_AMOUNT = 32 * 256 * 1024 # 32*256KB DISK_READ_SIZE = 64 * 1024 * 1024 # 64MB NETWORK_CHUNK_SIZE = 256 * 1024 # 256KB class FileSender: def __init__(self, file_path): self._file_path = file_path self._file_channel = None self._control_channel = None self._done = asyncio.Future() self._buffer_event = asyncio.Event() self._buffer_event.set() self._chunk_queue = asyncio.Queue(maxsize=128) # 足够容纳大量chunk以实现流水线 def set_channel(self, channel_type, channel): if channel_type == "file": self._file_channel = channel self._file_channel.on("open", self._on_both_channels_open) self._file_channel.bufferedAmountLowThreshold = MAX_BUFFERED_AMOUNT self._file_channel.on("bufferedamountlow", self._on_buffered_amount_low) elif channel_type == "control": self._control_channel = channel self._control_channel.on("open", self._on_both_channels_open) self._control_channel.on("message", lambda message: asyncio.create_task(self._on_control_message(message))) def _on_both_channels_open(self): if self._file_channel.readyState == "open" and self._control_channel.readyState == "open": asyncio.create_task(self._start_file_transfer()) def _on_buffered_amount_low(self): self._buffer_event.set() async def wait_until_done(self): await self._done async def _file_reader(self): """从磁盘读取大块数据,拆分成chunk后推入队列""" async with aiofiles.open(self._file_path, "rb") as f: while True: block = await f.read(DISK_READ_SIZE) if not block: break # 将块拆分成256KB的chunk,推入队列前保持在内存中 for i in range(0, len(block), NETWORK_CHUNK_SIZE): await self._chunk_queue.put(block[i:i + NETWORK_CHUNK_SIZE]) await self._chunk_queue.put(None) # 标记文件读取完成 async def _start_file_transfer(self): metadata = self._construct_metadata(self._file_path) self._control_channel.send(ControlMessage.create_json("metadata", json.dumps(metadata))) progress = Progress(metadata["file_size"]) reader_task = asyncio.create_task(self._file_reader()) while True: chunk = await self._chunk_queue.get() if chunk is None: break # 仅在需要时等待缓冲区空间 while self._file_channel.bufferedAmount > MAX_BUFFERED_AMOUNT: self._buffer_event.clear() await self._buffer_event.wait() self._file_channel.send(chunk) progress.update(chunk) await reader_task self._control_channel.send(ControlMessage.create_json("eof", metadata["file_name"])) self._done.set_result(None) def _construct_metadata(self, file_path): file_size = os.path.getsize(file_path) file_name = os.path.basename(file_path) file_hash = compute_hash(file_path) return {"file_name": file_name, "file_size": file_size, "hash": file_hash} async def _on_control_message(self, message): # 处理接收端的控制消息(比如传输确认) pass class FileReceiver: def __init__(self, path): self._file_channel = None self._control_channel = None self._metadata = None self._location = None self._progress = None self._path = path self._done = asyncio.Future() self._chunk_queue = asyncio.Queue(maxsize=256) # 更大的队列以支持流水线 self._writer_task = None async def wait_until_done(self): await self._done def set_channel(self, channel_type, channel): if channel_type == "file": self._file_channel = channel self._file_channel.on("message", lambda msg: self._chunk_queue.put_nowait(msg)) elif channel_type == "control": self._control_channel = channel self._control_channel.on("message", lambda msg: asyncio.create_task(self._on_control_message(msg))) async def _process_file(self): self._location = os.path.join(self._path, self._metadata["file_name"]) self._progress = Progress(self._metadata["file_size"]) async with aiofiles.open(self._location, "wb") as f: while True: chunk = await self._chunk_queue.get() if chunk is None: break await f.write(chunk) self._progress.update(chunk) # 校验哈希 if compute_hash(self._location) != self._metadata["hash"]: print("[ERROR] File corrupted") else: print("File received successfully") self._control_channel.send(ControlMessage.create_json("transfer_complete", self._metadata["file_name"])) self._done.set_result(None) self._file_channel.close() self._control_channel.close() async def _on_control_message(self, message): control_message = ControlMessage.from_json(message) if control_message.msg_type == "metadata": self._metadata = json.loads(control_message.data) print(f"Receiving file: {self._metadata['file_name']} ({self._metadata['file_size']} bytes)") self._writer_task = asyncio.create_task(self._process_file()) elif control_message.msg_type == "eof": # 告知文件写入任务所有数据已发送完成 await self._chunk_queue.put(None)
我的疑问
我已经尝试用队列来做流水线、调整队列大小,但发送端还是一直卡在等待缓冲区空间的状态。有没有办法优化缓冲区的使用逻辑,或者调整WebRTC数据通道的参数,来突破这个吞吐量瓶颈?
内容来源于stack exchange
相关产品推荐
相关产品推荐

