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

如何实现基于缓冲区的asyncio StreamReader与StreamWriter?

实现异步版BytesIO用于asyncio流测试的正确方式

直接继承asyncio.StreamReader和asyncio.StreamWriter不可行——这两个类与底层传输层(Transport)深度耦合,依赖_transport等未公开的内部属性,并非设计用于子类化。以下是两种可靠的实现方案:

方案一:独立实现兼容接口的AsyncBytesIO类

不依赖StreamReader/Writer的继承,直接实现符合异步流读写接口的类,完全控制内部逻辑:

import asyncio
from io import BytesIO

class AsyncBytesIO:
    def __init__(self, initial_data: bytes = b""):
        self._buffer = BytesIO(initial_data)
        self._closed = False
        # 模拟transport避免类似StreamWriter的__del__报错
        self._transport = None

    # 异步读接口
    async def read(self, n: int = -1) -> bytes:
        if self._closed:
            raise ValueError("Read from closed stream")
        return self._buffer.read(n)

    async def readline(self) -> bytes:
        if self._closed:
            raise ValueError("Read from closed stream")
        return self._buffer.readline()

    async def readexactly(self, n: int) -> bytes:
        if self._closed:
            raise ValueError("Read from closed stream")
        data = self._buffer.read(n)
        if len(data) < n:
            raise asyncio.IncompleteReadError(data, n)
        return data

    def at_eof(self) -> bool:
        return self._buffer.tell() >= len(self._buffer.getvalue())

    # 写接口
    def write(self, data: bytes) -> None:
        if self._closed:
            raise ValueError("Write to closed stream")
        self._buffer.write(data)

    async def drain(self) -> None:
        # 内存缓冲区无需等待,直接返回
        pass

    def close(self) -> None:
        self._closed = True

    def is_closing(self) -> bool:
        return self._closed

    # 工具方法
    def getvalue(self) -> bytes:
        pos = self._buffer.tell()
        self._buffer.seek(0)
        data = self._buffer.read()
        self._buffer.seek(pos)
        return data

    def seek(self, pos: int) -> int:
        return self._buffer.seek(pos)

    def tell(self) -> int:
        return self._buffer.tell()

    def reset(self) -> None:
        self._buffer.seek(0)
        self._closed = False

方案二:结合asyncio.StreamReader和自定义写入端

如果需要严格匹配asyncio.StreamReader的内部缓冲逻辑,可以直接使用StreamReader,搭配自定义的写入缓冲区:

import asyncio
from io import BytesIO

class AsyncBytesWriter:
    def __init__(self):
        self._buffer = BytesIO()
        self._closed = False
        self._transport = None  # 避免__del__报错

    def write(self, data: bytes) -> None:
        if self._closed:
            raise ValueError("Write to closed stream")
        self._buffer.write(data)

    async def drain(self) -> None:
        pass

    def close(self) -> None:
        self._closed = True

    def is_closing(self) -> bool:
        return self._closed

    def getvalue(self) -> bytes:
        return self._buffer.getvalue()

# 使用示例
async def test_stream_reader():
    reader = asyncio.StreamReader()
    writer = AsyncBytesWriter()

    # 给reader喂数据
    reader.feed_data(b"hello async\n")
    reader.feed_eof()

    # 测试读操作
    assert await reader.readline() == b"hello async\n"
    assert await reader.read() == b""

    # 测试写操作
    writer.write(b"test write")
    assert writer.getvalue() == b"test write"

原始代码的临时修复方案

如果一定要基于原始继承思路临时修复,至少需要补充_transport、_protocol等内部依赖属性:

import asyncio
from io import BytesIO

class AsyncBytesIO(asyncio.StreamReader, asyncio.StreamWriter):
    def __init__(self, data: bytes = b""):
        # 先调用父类初始化
        asyncio.StreamReader.__init__(self)
        # StreamWriter需要的必填内部属性
        self._transport = None
        self._protocol = None
        self._loop = asyncio.get_event_loop()
        # 自定义缓冲区
        self._buffer = BytesIO(data)
        self._closed = False

    # 保留你原来实现的read、write等方法...

async def test():
    x = AsyncBytesIO()
    assert False

asyncio.run(test())

但这种方式仍可能触发其他未文档化的内部依赖问题,不推荐长期使用。

内容的提问来源于stack exchange,提问作者Joseph Garvin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 20:35:21