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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 08:27:58