asyncio 中流式 I/O 和背压怎么理解?StreamReader、StreamWriter 怎么用?
简化版
asyncio 的 StreamReader 和 StreamWriter 提供基于流的异步网络 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) # 满了会等待
| 位置 | 背压表现 |
|---|---|
| StreamWriter | await drain() 等缓冲下降 |
| Queue | await 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() 才是控制缓冲的关键。高并发异步程序真正难的是让快生产者尊重慢消费者,否则不阻塞线程也会把内存打穿。