← 返回题目列表

asyncio 中流式 I/O 和背压怎么理解?StreamReader、StreamWriter 怎么用?

高频 困难 第 12 / 27 题 更新于 2026/07/31
asyncioStreamReaderStreamWriter背压

简化版

asyncioStreamReaderStreamWriter 提供基于流的异步网络 I/O。读用 await reader.read(),写用 writer.write()await writer.drain()drain() 很关键,它让写入方在缓冲区过大时等待,形成背压,避免内存被无限写爆。

详细版

流式 I/O 适合 TCP 这类连续字节流。StreamReader 负责异步读取字节,StreamWriter 负责写入字节。因为 TCP 没有天然消息边界,应用层要自己设计分隔符、长度前缀或固定长度协议。

reader, writer = await asyncio.open_connection("127.0.0.1", 8888)
writer.write(b"ping\n")
await writer.drain()
data = await reader.readline()
writer.close()
await writer.wait_closed()

背压指下游处理不过来时,上游要减速。writer.write() 只是把数据放入缓冲,await writer.drain() 才会在缓冲高水位时暂停生产者。面试里要强调:不 drain 的高速写入可能导致内存膨胀。

完整版教学

一、为什么 TCP 流需要应用层协议

TCP 提供的是可靠字节流,不保留你每次 write() 的消息边界。你写两次,对方可能一次读到;你写一次,对方也可能分几次读到。应用层必须定义如何拆包。

发送端:
write("hello")
write("world")

接收端可能读到:
"helloworld"
或 "hel" + "loworld"

常见做法有换行分隔、固定长度、长度前缀。面试如果能补这一点,就说明你不只是会调用 read()

二、StreamReader 怎么读

StreamReader 提供 read(n)readline()readexactly(n) 等方法。选择取决于协议。按行协议用 readline(),长度前缀协议先 readexactly(4) 读长度,再读对应字节数。

line = await reader.readline()
body = await reader.readexactly(1024)
chunk = await reader.read(4096)
方法含义适合
read(n)最多读 n 字节流式块处理
readline()读到换行行协议
readexactly(n)必须读满 n 字节长度前缀

如果对方断开,读取可能返回空 bytes 或抛出异常。网络代码必须处理断开、超时和半包。

三、StreamWriter 写入为什么要 drain

writer.write(data) 通常不会立刻把数据全部发到网络,它先写入传输缓冲区。缓冲区如果持续增长,内存会越来越大。await writer.drain() 会在缓冲区达到高水位时等待,直到降到低水位。

for chunk in chunks:
    writer.write(chunk)
    await writer.drain()

数字化理解:生产者每秒生成 100MB,网络只能发送 10MB。如果不等待,90MB/s 会堆在内存缓冲里;10 秒就是约 900MB。drain() 让生产者感知下游速度。

背压就是“下游慢,上游要跟着慢”,否则异步程序只是更快地把内存撑爆。

四、背压和队列有什么关系

背压不只出现在网络写入,也出现在生产者消费者队列里。asyncio.Queue(maxsize=N) 满了以后,await queue.put() 会等待,这也是背压。没有 maxsize 的队列在高峰期可能无限增长。

queue = asyncio.Queue(maxsize=100)
await queue.put(item)  # 满了会等待
位置背压表现
StreamWriterawait drain() 等缓冲下降
Queueawait put() 等队列有空位
HTTP 客户端连接池满时等待
数据库池连接耗尽时等待

背压的目标不是让系统永远不慢,而是让慢可控、内存可控、故障可恢复。

五、连接关闭要怎么做

写完数据后要 writer.close(),并 await writer.wait_closed() 等待关闭完成。只 close 不 wait,程序退出时可能还有缓冲未清理。对于长连接,还要处理对端断开和取消。

writer.close()
await writer.wait_closed()

如果协议需要半关闭、优雅结束或 TLS,关闭流程还会更复杂。普通面试至少要说出 close + wait_closed 和异常处理。

六、服务端处理要注意什么

asyncio.start_server 可以启动异步 TCP 服务。每个连接通常由一个协程处理。连接处理协程里不能做阻塞 CPU 或同步 I/O,否则会卡住事件循环。

async def handle(reader, writer):
    while line := await reader.readline():
        writer.write(line.upper())
        await writer.drain()

server = await asyncio.start_server(handle, "127.0.0.1", 8888)

如果有 1000 个连接,每个连接都能在等待网络时让出控制权,这是 asyncio 的优势。但如果某个 handler 里跑 500ms CPU 计算,所有连接都会受影响。

七、常见误区与追问

  • 误区:write() 后数据已经发到对端。 它通常只是写入缓冲,发送由底层传输完成。
  • 误区:异步写入不需要考虑内存。 不 drain 或无界队列会让内存快速膨胀。
  • 误区:TCP 一次 write 对应一次 read。 TCP 是字节流,没有消息边界。
  • 追问:drain() 的作用是什么? 等待写缓冲从高水位降下来,形成背压。
  • 追问:如何处理粘包拆包? 设计应用层协议,如长度前缀、换行分隔或固定长度。
  • 追问:服务端 handler 能否调用阻塞函数? 不应直接调用,会阻塞事件循环;用异步库或线程池。
  • 追问:队列背压怎么做? 设置 maxsize,让生产者在队列满时 await。

八、加强记忆

asyncio 流式 I/O 记住“读协议、写 drain、关 wait_closed、慢要背压”。TCP 是字节流,不帮你切消息;write() 不等于发完,drain() 才是控制缓冲的关键。高并发异步程序真正难的是让快生产者尊重慢消费者,否则不阻塞线程也会把内存打穿。