asyncio 的 Transport 和 Protocol 是什么?什么时候要用低层 API?
简化版
asyncio 提供了两层网络 API:高层的 Streams(StreamReader/StreamWriter,用 await 顺序读写,符合直觉)和低层的 Transport/Protocol(基于回调**,性能更高但写法完全不同)。两者的职责划分很清晰:Transport 代表「连接」,负责写数据和关闭连接(transport.write(data)、transport.close());Protocol 代表「协议逻辑」,由你实现,事件循环会在特定时刻回调它的方法——connection_made(transport)(连接建立)、data_received(data)(收到数据,注意是 TCP 字节流的任意分片,不保证消息边界)、eof_received()(对端半关闭)、connection_lost(exc)(连接断开)。关键区别在于控制流方向:Streams 是「你主动去 await 读」,Protocol 是「数据到了框架来调你」——后者省掉了协程调度和中间缓冲的开销,所以吞吐更高**,代价是状态要自己维护(消息拆包、粘包处理、缓冲区管理都得手写)。背压在低层是通过一对回调实现的:写太快时 Transport 会调用 protocol.pause_writing(),缓冲降下来后调 resume_writing()——不理会这对回调就会导致内存无限增长;读方向则用 transport.pause_reading()/resume_reading() 控制。事实上 Streams 本身就是基于 Protocol 实现的(StreamReaderProtocol),所以两者不是并列关系而是分层关系。核心记忆:Transport 管连接和写、Protocol 管协议和回调;Streams 更好写、Protocol 更快;data_received 收到的是字节流分片,必须自己处理粘包。
详细版
两层 API 对比:
| 维度 | Streams(高层) | Transport/Protocol(低层) |
|---|---|---|
| 控制流 | await reader.read()(拉) | 回调 data_received()(推) |
| 心智负担 | 低,顺序写法 | 高,状态机自己管 |
| 性能 | 有协程调度和缓冲开销 | 更快(少一层) |
| 粘包处理 | readexactly/readuntil 帮你做 | 完全自己实现 |
| 背压 | await writer.drain() | pause_writing/resume_writing 回调 |
| UDP | ❌ 不支持 | ✅ 只能用它 |
| 适用 | 绝大多数业务 | 高吞吐、自定义协议、UDP、框架 |
import asyncio
# ① ★低层:实现一个 Protocol★
class EchoProtocol(asyncio.Protocol):
def connection_made(self, transport: asyncio.Transport):
self.transport = transport # ★保存,用于写和关闭★
peer = transport.get_extra_info("peername")
print("连接来自", peer)
def data_received(self, data: bytes):
# ★注意:data 是 TCP 流的任意分片,不是"一条消息"★
self.transport.write(data) # ★写是同步的,不用 await★
def eof_received(self) -> bool | None:
# 对端半关闭;返回 True 表示"保持连接开着"
return False # False/None → 关闭连接
def connection_lost(self, exc: Exception | None):
print("连接断开", exc) # ★清理资源★
async def main():
loop = asyncio.get_running_loop()
server = await loop.create_server(EchoProtocol, "0.0.0.0", 8888)
async with server:
await server.serve_forever()
# ② ★高层:同样的功能用 Streams★
async def handle(reader, writer):
while data := await reader.read(4096):
writer.write(data)
await writer.drain() # ★背压★
writer.close()
await writer.wait_closed()
server = await asyncio.start_server(handle, "0.0.0.0", 8888)
# ③ ★粘包处理:低层必须自己拆包★
class LengthPrefixProtocol(asyncio.Protocol):
"""协议:4 字节大端长度 + 消息体"""
def connection_made(self, transport):
self.transport = transport
self._buf = bytearray() # ★自己维护缓冲★
def data_received(self, data):
self._buf.extend(data)
while True:
if len(self._buf) < 4: # ★头都没收全★
return
length = int.from_bytes(self._buf[:4], "big")
if length > MAX_MSG: # ★必须限长,防内存攻击★
self.transport.abort()
return
if len(self._buf) < 4 + length: # ★消息体没收全★
return
msg = bytes(self._buf[4:4 + length])
del self._buf[:4 + length] # ★从缓冲里移除★
self.handle_message(msg) # 处理完整消息
# ④ ★背压:必须实现这对回调★
class FlowControlProtocol(asyncio.Protocol):
def connection_made(self, transport):
self.transport = transport
self._paused = False
transport.set_write_buffer_limits(high=64*1024, low=16*1024)
def pause_writing(self): # ★写缓冲超过 high★
self._paused = True # ★停止产生数据★
def resume_writing(self): # ★降到 low 以下★
self._paused = False # ★继续★
def data_received(self, data):
if self._paused:
self.transport.pause_reading() # ★读方向也可以暂停★
# ⑤ UDP(★只能用低层★)
class UDPProtocol(asyncio.DatagramProtocol):
def datagram_received(self, data, addr):
self.transport.sendto(b"pong", addr)
transport, protocol = await loop.create_datagram_endpoint(
UDPProtocol, local_addr=("0.0.0.0", 9999))
# ⑥ 客户端也一样
transport, protocol = await loop.create_connection(
lambda: EchoProtocol(), "example.com", 80)
⚠️ 三个必须理解的点:①
data_received(data)收到的绝不是「一条完整消息」——TCP 是字节流,没有消息边界:发送方一次send的数据可能被拆成多次data_received回调(分包),也可能多次send的数据合并成一次回调(粘包)。所以任何基于 TCP 的协议都必须自己定义边界(长度前缀、分隔符、固定长度),并维护一个缓冲区反复尝试拆包。Streams 的readexactly(n)/readuntil(sep)就是把这件事封装好了。② 背压在低层是「一对回调」而不是await:调用transport.write(data)不会阻塞——数据先进 Transport 的写缓冲,如果对端很慢,缓冲会无限增长直到内存耗尽。asyncio 的机制是:缓冲超过high水位时调用你的pause_writing(),降到low以下时调用resume_writing()——你必须在这两个回调里真正停止/恢复产生数据,否则背压形同虚设。Streams 的await writer.drain()就是基于这对回调实现的。③ Streams 不是「另一套实现」,而是基于 Protocol 的封装:asyncio.start_server内部创建的是StreamReaderProtocol,它在data_received里把数据喂给StreamReader的缓冲、在pause_writing/resume_writing里控制drain()的等待——理解了这层关系,就知道两者的性能差异来自哪里(多了一层缓冲和协程调度)。
完整版教学
一、两层 API 的分工
★ asyncio 的网络分层:
┌─────────────────────────────────────────────┐
│ 你的业务代码 │
├─────────────────────────────────────────────┤
│ ★Streams★:StreamReader / StreamWriter │ ← 高层(await 风格)
│ start_server / open_connection │
├─────────────────────────────────────────────┤
│ ★Protocol★(你实现):data_received 等回调 │ ← 低层(回调风格)
│ ★Transport★(框架提供):write / close │
├─────────────────────────────────────────────┤
│ 事件循环 + selector(epoll/kqueue/IOCP) │
└─────────────────────────────────────────────┘
★ Streams ★就是基于★ Protocol 实现的(StreamReaderProtocol)
→ 不是两套并列的实现,是★分层★
★ 职责划分(★面试要能说清★):
★Transport(框架实现,你只调用)★
- 代表"一个连接"
- transport.write(data) ★同步调用,不阻塞(数据进缓冲)★
- transport.writelines(list)
- transport.close() ★优雅关闭(先把缓冲写完)★
- transport.abort() ★立即关闭(丢弃缓冲)★
- transport.get_extra_info("peername"/"socket"/"sslcontext"/...)
- transport.pause_reading() / resume_reading() ★读方向流控★
- transport.set_write_buffer_limits(high, low)
★Protocol(你实现,框架回调你)★
- connection_made(transport) ★连接建立★(保存 transport)
- data_received(data) ★收到数据(字节流分片)★
- eof_received() 对端发了 FIN(半关闭)
- connection_lost(exc) ★连接彻底断开(清理)★
- pause_writing() / resume_writing() ★背压回调★
★ 控制流的方向(★两者的本质区别★):
Streams(拉模型 pull):
data = await reader.read(1024) # ★我主动要数据,没有就挂起★
Protocol(推模型 push):
def data_received(self, data): # ★数据来了框架调我★
...
→ 拉模型符合直觉、状态在协程栈里自然保存
→ 推模型没有协程调度开销,但★状态要自己存在 self 上★
★ 三对 API 的对应关系:
loop.create_server(ProtocolFactory, host, port) ←→ asyncio.start_server(cb)
loop.create_connection(ProtocolFactory, host, port) ←→ asyncio.open_connection()
loop.create_datagram_endpoint(...) ←→ ★没有高层对应(UDP 只能低层)★
loop.subprocess_exec(...) ←→ asyncio.create_subprocess_exec()
asyncio 的网络 API 是分层而不是并列的:Streams 就建立在 Protocol 之上(asyncio.start_server 内部创建的是 StreamReaderProtocol)。职责划分要能说清:Transport 由框架实现、代表「一个连接」,提供 write()(同步调用、不阻塞,数据先进缓冲)、close()(优雅关闭,先把缓冲写完)、abort()(立即关闭、丢弃缓冲)、get_extra_info()(拿 peername、底层 socket、SSL 上下文)和流控方法;Protocol 由你实现、代表「协议逻辑」,框架会在特定时刻回调 connection_made、data_received、eof_received、connection_lost 以及背压相关的 pause_writing/resume_writing。本质区别在控制流方向:Streams 是拉模型(你主动 await 读,状态自然保存在协程栈里),Protocol 是推模型(数据来了框架调你,状态必须自己存在 self 上)。还有一个实用信息:UDP 只有低层 API(create_datagram_endpoint),没有 Streams 版本。
二、data_received:TCP 没有消息边界
★ 最重要的认知:TCP 是★字节流★,不是消息流
发送方:sock.send(b"HELLO");sock.send(b"WORLD")
接收方 data_received 可能被调用成:
① data_received(b"HELLOWORLD") ★粘包★(两次发送合并)
② data_received(b"HEL")
data_received(b"LOWORLD") ★分包★(一次发送被拆开)
③ data_received(b"HELLO")
data_received(b"WORLD") 碰巧一致(★不能依赖★)
→ ★协议必须自己定义消息边界★
★ 三种常见的边界方案:
① ★长度前缀★(最常用、最高效)
[4 字节长度][消息体]
✓ 解析简单、二进制安全、能预知消息大小
★ 必须校验长度上限★(否则一个伪造的 4GB 长度就让你 OOM)
② ★分隔符★(如 \r\n,HTTP、Redis 协议用)
✓ 人类可读、易调试
✗ 消息体里出现分隔符要转义;★必须限制最大行长★
③ ★固定长度★
✓ 最简单
✗ 只适合定长记录
★ 标准的拆包循环(★模板,背下来★):
def data_received(self, data):
self._buf.extend(data) # ★① 追加到缓冲★
while True: # ★② 循环拆包(一次可能收到多条)★
if len(self._buf) < HEADER: # ★③ 头不完整 → 等下次★
return
length = parse_header(self._buf)
if length > MAX_MSG: # ★④ 限长防攻击★
self.transport.abort()
return
if len(self._buf) < HEADER + length: # ★⑤ 体不完整 → 等★
return
msg = bytes(self._buf[HEADER:HEADER+length])
del self._buf[:HEADER+length] # ★⑥ 从缓冲移除★
self.handle(msg) # ★⑦ 处理★
★ 三个易错点:
- ★忘了 while 循环★:一次 data_received 可能包含多条完整消息
- ★忘了限长★:恶意的超大长度字段 → 内存耗尽
- ★缓冲用 bytes 而不是 bytearray★:每次拼接都是 O(n) 复制 → O(n²)
★ Streams 帮你做了什么:
await reader.readexactly(4) # ★读满 4 字节,不够就继续等★
await reader.readuntil(b"\r\n") # ★读到分隔符★
→ 内部就是维护缓冲 + 判断是否满足条件 + 挂起等待
★ readuntil 也有 limit(默认 64KB)→ 超了抛 LimitOverrunError(★同样是防攻击★)
★ 一个真实的协议实现要处理的:
□ 消息边界(拆包)
□ ★最大消息长度★
□ 版本/魔数校验
□ 心跳/keepalive(检测半开连接)
□ ★半包超时★(收到半条消息后对端不发了)
□ 背压(见下节)
□ 优雅关闭(发完缓冲再关)
最重要的认知是「TCP 是字节流,没有消息边界」:发送方两次 send 的数据可能粘包(合并成一次 data_received)也可能分包(一次发送被拆成多次回调),碰巧一致也不能依赖。所以协议必须自己定义边界:长度前缀(最常用最高效,但必须校验长度上限否则一个伪造的 4GB 长度就能让你 OOM)、分隔符(HTTP、Redis 用,但必须限制最大行长)、或固定长度。那个拆包循环模板要背下来,三个易错点是:忘了 while 循环(一次回调可能包含多条完整消息)、忘了限长、以及用 bytes 而不是 bytearray 做缓冲(每次拼接都是 O(n) 复制,整体退化成 O(n²))。Streams 的 readexactly(n) 和 readuntil(sep) 就是把这件事封装好了(注意 readuntil 也有 64KB 的 limit,超了抛 LimitOverrunError,同样是防攻击设计)。真实协议还要处理心跳、半包超时、优雅关闭等。
三、背压:一对回调而不是 await
★ 问题:transport.write() 不阻塞
for chunk in huge_data:
transport.write(chunk) # ★立刻返回,数据进写缓冲★
→ 如果对端读得慢,写缓冲★无限增长★ → ★内存耗尽★
→ 这就是"没有背压"
★ asyncio 的低层背压机制(★水位线 + 回调★):
transport.set_write_buffer_limits(high=64*1024, low=16*1024)
写缓冲 > high → 框架调用 ★protocol.pause_writing()★
写缓冲 < low → 框架调用 ★protocol.resume_writing()★
★ 关键:这只是"通知",★真正停下来是你的责任★
class P(asyncio.Protocol):
def pause_writing(self):
self._can_write.clear() # ★停止产生数据★
def resume_writing(self):
self._can_write.set()
★ 不实现这两个方法会怎样:
→ 默认实现是空的 → ★你会一直写下去 → 缓冲无限增长★
→ 现象:内存缓慢增长,慢客户端拖垮服务
★ 读方向的流控:
transport.pause_reading() ★暂停读(不再触发 data_received)★
transport.resume_reading()
→ 用于"我处理不过来了,让内核的接收缓冲堆着"
→ ★TCP 窗口会自动收缩 → 压力传递到发送方★(这才是真正的端到端背压)
★ Streams 是怎么用它的(★理解封装★):
writer.write(data) # 同样不阻塞
await writer.drain() # ★如果处于 paused 状态就挂起,直到 resume★
→ StreamWriter.drain() 内部:
if 不是 paused: 立即返回
else: 创建一个 Future,等 resume_writing 时 set_result
→ ★所以"忘记 await drain()"和"低层不实现 pause_writing"是同一个错误★
★ 完整的流控示例(低层):
class Producer(asyncio.Protocol):
def connection_made(self, transport):
self.transport = transport
self.transport.set_write_buffer_limits(high=256*1024)
self._paused = asyncio.Event()
self._paused.set() # 初始可写
asyncio.create_task(self._produce())
def pause_writing(self):
self._paused.clear() # ★阻止 _produce 继续★
def resume_writing(self):
self._paused.set()
async def _produce(self):
async for chunk in data_source():
await self._paused.wait() # ★背压点★
self.transport.write(chunk)
★ 水位线怎么设:
high 太小 → 频繁 pause/resume,吞吐下降
high 太大 → 单连接占用内存多(★N 个连接就是 N 倍★)
★ 默认 64KB high / 16KB low 对大多数场景合适
★ 万级连接时要调小(64KB × 10000 = 640MB)
背压在低层是「一对回调」而不是 await。问题的根源是 transport.write() 不阻塞——数据先进写缓冲,对端读得慢时缓冲会无限增长直到内存耗尽。asyncio 的机制是水位线 + 回调:写缓冲超过 high 时框架调用你的 pause_writing()、降到 low 以下时调 resume_writing()——但这只是「通知」,真正停下来是你的责任;不实现这两个方法的话默认是空实现,你会一直写下去、缓冲无限增长。读方向则用 transport.pause_reading()(不再触发 data_received,让内核接收缓冲堆着,TCP 窗口会自动收缩把压力传回发送方,这才是真正的端到端背压)。理解封装关系很有帮助:await writer.drain() 内部就是「如果处于 paused 状态就挂起,等 resume_writing 时唤醒」——所以**「忘记 await drain()」和「低层不实现 pause_writing」本质上是同一个错误**。水位线的取舍是:太小会频繁 pause/resume 降低吞吐,太大则单连接占内存多(万级连接时 64KB × 10000 = 640MB,要调小)。
四、生命周期与错误处理
★ 回调的调用顺序(★保证★):
connection_made(transport) ★一定第一个★
↓
data_received(data) × N (可能 0 次)
pause_writing / resume_writing (交错发生)
↓
eof_received() (可选,对端 FIN 时)
↓
connection_lost(exc) ★一定最后一个,且★一定会被调用★★
★ connection_lost 是清理的唯一可靠位置:
- exc is None → 正常关闭(本端 close 或对端正常断开)
- exc 不为 None → ★异常断开(连接重置、超时等)★
★ eof_received 的语义(★容易搞错★):
对端调用了 shutdown(SHUT_WR) → 只关闭了它的写方向(半关闭)
返回值决定后续:
★False / None★ → 框架关闭 transport → 接着调 connection_lost
★True★ → ★保持连接打开★(你还能继续 write)
→ 返回 True 用于"请求-响应"型协议:对端说完了,我还要回复
★ 关闭的三种方式:
transport.close() ★优雅★:不再读,把写缓冲刷完,然后关闭
→ ★之后仍会调 connection_lost★
transport.abort() ★立即★:丢弃缓冲,直接关闭
→ 用于"检测到协议错误/攻击"
transport.write_eof() 半关闭:发 FIN,但还能继续读
★ 异常处理(★低层最容易踩的坑★):
★回调里抛出的异常不会传播到"调用方"★(因为没有调用方!)
→ 它会被送到 ★loop.set_exception_handler★ 注册的处理器
→ ★没注册的话只会打印到 stderr,很容易被忽略★
✓ 每个回调内部自己 try/except:
def data_received(self, data):
try:
self._parse(data)
except Exception:
logging.exception("协议解析失败")
self.transport.abort() # ★出错就断开,别让状态机错乱★
★ 在回调里做异步操作(★关键限制★):
★回调是同步函数,不能 await★
def data_received(self, data):
await self.db.save(data) # ✗ SyntaxError
✓ 创建任务:
task = asyncio.create_task(self.db.save(data))
self._tasks.add(task) # ★保存引用防 GC★
task.add_done_callback(self._tasks.discard)
★ 但要注意:这样就★失去了背压★(数据来多快就创建多少任务)
→ 需要自己加 Semaphore 限流,或 pause_reading()
★ 一个健壮的 Protocol 骨架:
class MyProtocol(asyncio.Protocol):
def connection_made(self, transport):
self.transport = transport
self._buf = bytearray()
self._tasks: set[asyncio.Task] = set()
self._closing = False
def data_received(self, data):
if self._closing: return
try:
self._buf.extend(data)
self._process_buffer()
except ProtocolError:
self.transport.abort()
except Exception:
logging.exception("未预期的错误")
self.transport.abort()
def connection_lost(self, exc):
self._closing = True
for t in self._tasks: t.cancel() # ★清理未完成的任务★
if exc: logging.warning("连接异常断开: %r", exc)
生命周期有明确的保证:connection_made 一定第一个、connection_lost 一定最后一个且一定会被调用(所以它是清理的唯一可靠位置,exc 为 None 表示正常关闭、非 None 表示异常断开)。eof_received 的返回值容易搞错:返回 False/None 会让框架关闭连接,返回 True 则保持连接打开(用于「对端说完了但我还要回复」的请求-响应型协议)。关闭有三种:close()(优雅,刷完写缓冲)、abort()(立即,丢弃缓冲,用于检测到协议错误时)、write_eof()(半关闭)。最容易踩的坑是异常处理:回调里抛出的异常不会传播到「调用方」(因为根本没有调用方),它会被送到 loop.set_exception_handler,没注册的话只打印到 stderr、极易被忽略——所以每个回调内部都要自己 try/except,出错就 abort() 避免状态机错乱。还有个关键限制:回调是同步函数、不能 await,需要异步操作只能 create_task(记得保存引用防 GC),但这样会失去背压,需要自己加限流或 pause_reading()。
五、什么时候该用低层 API
★ 该用的四种情况:
① ★UDP★(唯一选择)
asyncio 没有 UDP 的 Streams API
→ 必须用 DatagramProtocol + create_datagram_endpoint
② ★极致吞吐★
省掉 StreamReader 的中间缓冲和协程调度
→ 实测在小消息高频场景能快 ★20%~50%★
→ 但★只在框架开销占主导时才明显★(和 uvloop 的道理一样)
③ ★自定义二进制协议 / 框架开发★
数据库驱动(asyncpg)、消息队列客户端、RPC 框架
→ 需要精细控制解析、复用缓冲、零拷贝
④ ★需要 Transport 的低层能力★
get_extra_info("socket") 设置 TCP_NODELAY、SO_KEEPALIVE
精细的流控时机
★ 不该用的情况(★绝大多数业务★):
✗ 普通的 TCP 客户端/服务端 → ★用 Streams★
✗ HTTP 服务 → ★用 aiohttp/FastAPI★(它们已经优化过了)
✗ 只是想"看起来更底层更专业"
→ ★低层 API 的复杂度是实打实的★:状态机、拆包、背压、错误处理全要自己写
★ 真实项目里的选择:
asyncpg(PostgreSQL 驱动) → ★Protocol★(自定义二进制协议 + 极致性能)
aiohttp → ★Protocol★(自己实现 HTTP 解析)
redis-py 的异步部分 → 混合
你的业务代码 → ★Streams 或更高层的库★
★ 性能差距的量化(★别过度期待★):
echo 服务器,1KB 消息:
Streams: ~80k msg/s
Protocol: ~110k msg/s (★快约 35%★)
加上业务逻辑(解析 JSON + 查缓存,每条 50μs)后:
Streams: ~18k msg/s
Protocol: ~19k msg/s (★差距缩到 5%★)
→ ★业务逻辑越重,低层 API 的优势越小★
★ 折中方案:
✓ 用 Streams 但避免它的低效用法:
- readexactly / readuntil 而不是循环 read(1)
- 批量写 + 一次 drain,而不是每条都 drain
- 调大 limit 减少内部拷贝
✓ 或者用现成的高性能库(asyncpg、aiohttp)而不是自己写 Protocol
★ 学习价值(★即使不直接用★):
理解 Transport/Protocol 能让你搞懂:
- Streams 的背压是怎么实现的(drain 等的是 resume_writing)
- 为什么"忘了 await drain()"会内存爆
- ★TCP 粘包/分包为什么必须自己处理★
- 连接的完整生命周期和清理时机
→ ★这些认知在用高层 API 时同样重要★
低层 API 该用的只有四种情况:UDP(唯一选择,asyncio 没有 UDP 的 Streams API)、极致吞吐(省掉中间缓冲和协程调度,小消息高频场景能快 20%~50%)、自定义二进制协议或框架开发(asyncpg、aiohttp 都是这么做的)、需要 Transport 的低层能力(拿底层 socket 设置 TCP_NODELAY)。绝大多数业务不该用——低层 API 的复杂度是实打实的(状态机、拆包、背压、错误处理全要自己写)。性能差距也要有客观认识:echo 场景 Protocol 比 Streams 快约 35%,但一旦加上业务逻辑(每条 50μs),差距就缩到 5%——业务逻辑越重,低层 API 的优势越小(和 uvloop 是一个道理)。折中方案是用 Streams 但避免低效用法(用 readexactly 而不是循环 read(1)、批量写后一次 drain)。最后强调学习价值:理解 Transport/Protocol 能让你搞懂 Streams 的背压是怎么实现的、为什么忘了 await drain() 会内存爆、TCP 粘包为什么必须自己处理——这些认知在用高层 API 时同样重要。
六、完整示例与实践清单
★ 一个生产可用的长度前缀协议:
import asyncio, logging, struct
HEADER = struct.Struct(">I") # 4 字节大端长度
MAX_MSG = 8 * 1024 * 1024 # ★8MB 上限★
class MessageProtocol(asyncio.Protocol):
def connection_made(self, transport):
self.transport = transport
self._buf = bytearray()
self._tasks: set[asyncio.Task] = set()
self._can_write = asyncio.Event()
self._can_write.set()
transport.set_write_buffer_limits(high=256*1024, low=64*1024)
sock = transport.get_extra_info("socket")
if sock:
sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) # ★低延迟★
# ---- 读 ----
def data_received(self, data):
self._buf.extend(data)
try:
while True:
if len(self._buf) < HEADER.size:
return
(length,) = HEADER.unpack_from(self._buf)
if length > MAX_MSG: # ★防攻击★
logging.warning("消息过大 %d,断开", length)
self.transport.abort()
return
if len(self._buf) < HEADER.size + length:
return
msg = bytes(self._buf[HEADER.size:HEADER.size + length])
del self._buf[:HEADER.size + length]
self._dispatch(msg)
except Exception:
logging.exception("解析失败")
self.transport.abort()
def _dispatch(self, msg):
t = asyncio.create_task(self._handle(msg))
self._tasks.add(t) # ★防 GC★
t.add_done_callback(self._tasks.discard)
async def _handle(self, msg):
try:
reply = await process(msg)
await self.send(reply)
except Exception:
logging.exception("处理失败")
# ---- 写(带背压)----
async def send(self, payload: bytes):
await self._can_write.wait() # ★背压点★
self.transport.write(HEADER.pack(len(payload)) + payload)
def pause_writing(self):
self._can_write.clear()
def resume_writing(self):
self._can_write.set()
# ---- 生命周期 ----
def eof_received(self):
return False # 对端关了 → 我们也关
def connection_lost(self, exc):
for t in self._tasks:
t.cancel() # ★清理未完成任务★
if exc:
logging.warning("连接异常: %r", exc)
★ 实践清单:
□ ★默认用 Streams★,只有明确理由才下沉到 Protocol
□ data_received 里★必须循环拆包★(一次可能多条)
□ ★必须限制最大消息长度★(防内存攻击)
□ 缓冲用 ★bytearray★(bytes 拼接是 O(n²))
□ ★实现 pause_writing / resume_writing★(否则没有背压)
□ 回调里 ★try/except★(异常不会传播,只会进 exception_handler)
□ 回调里创建的任务要★保存引用 + connection_lost 时 cancel★
□ ★connection_lost 是唯一可靠的清理点★
□ 检测到协议错误用 ★abort()★ 而不是 close()
□ 需要低延迟时设 ★TCP_NODELAY★(禁用 Nagle)
□ 万级连接时★调小 write buffer 水位线★
★ 一句话总结:
★"Transport 管连接和写、Protocol 管协议和回调;
Streams 只是基于它们的封装。
低层更快但要自己处理拆包、背压、状态机和错误——
除非做 UDP、自定义协议或框架,否则用 Streams。"★
那个完整示例覆盖了生产协议的全部要素:长度前缀拆包 + 最大长度校验 + bytearray 缓冲 + 背压回调 + 任务引用管理 + connection_lost 清理 + TCP_NODELAY。实践清单里最关键的六条:默认用 Streams、data_received 里必须循环拆包、必须限制最大消息长度、缓冲用 bytearray、实现 pause_writing/resume_writing、回调里必须 try/except(异常不会传播,只会进 exception handler)。还有两个细节:检测到协议错误要用 abort() 而不是 close()(不要把可能已损坏的缓冲数据发出去),以及万级连接时要调小写缓冲水位线(64KB × 10000 = 640MB)。
记忆钩子:「asyncio 的网络 API 是★分层★不是并列:★Streams 就建立在 Protocol 之上★(内部是 StreamReaderProtocol)。职责划分:★Transport 由框架实现、代表『连接』★(write 是★同步调用不阻塞★、close 优雅关闭会刷完缓冲、★abort 立即关闭丢弃缓冲★、get_extra_info 拿底层 socket);★Protocol 由你实现、代表『协议逻辑』★,框架回调 connection_made → data_received×N → eof_received → ★connection_lost(一定最后且一定会被调用,是清理的唯一可靠位置)★。本质区别是控制流方向:★Streams 是拉模型(await 读、状态在协程栈里),Protocol 是推模型(数据来了框架调你、状态必须自己存在 self 上)★。★三个必须理解的点★:①★data_received 收到的绝不是『一条完整消息』★——TCP 是字节流,会粘包也会分包,协议必须自己定边界(长度前缀/分隔符/定长),拆包模板要点是★循环拆(一次可能多条)+ 限制最大长度(防 OOM 攻击)+ 用 bytearray(bytes 拼接是 O(n²))★;②★背压是一对回调而不是 await★——transport.write 不阻塞,写缓冲超过 high 水位时框架调 ★pause_writing()★、降到 low 调 resume_writing(),★但真正停下来是你的责任★,不实现就是缓冲无限增长;读方向用 ★pause_reading()★ 让 TCP 窗口收缩把压力传回发送方;★await writer.drain() 内部就是等 resume_writing★,所以『忘了 drain』和『不实现 pause_writing』是同一个错误;③★回调是同步函数不能 await★,要异步就 create_task(★记得保存引用防 GC、connection_lost 时 cancel★),但那样★会失去背压★,要自己限流。★异常处理是低层最容易踩的坑:回调里抛的异常不会传播(没有调用方!),只会进 loop.set_exception_handler,没注册就只打印到 stderr★——所以每个回调都要自己 try/except,出错用 abort()。★什么时候用低层:UDP(唯一选择)、极致吞吐(echo 场景快 35%,但★加上业务逻辑后只快 5%★)、自定义二进制协议/框架(asyncpg、aiohttp 都是)★——绝大多数业务用 Streams 就对了。」
七、常见误区与追问
- 误区:
data_received(data)收到的是对端发送的一条完整消息。 TCP 是字节流,没有任何消息边界的概念。发送方两次send(b"HELLO")和send(b"WORLD"),接收方可能收到一次b"HELLOWORLD"(粘包),也可能收到b"HEL"和b"LOWORLD"两次(分包),碰巧一一对应也纯属偶然、绝不能依赖。粘包和分包的成因很多:Nagle 算法合并小包、MTU 分片、网络拥塞、接收缓冲区状态。所以任何基于 TCP 的协议都必须自己定义消息边界(长度前缀、分隔符或固定长度),并在data_received里维护缓冲区 + 循环拆包。用 Streams 时readexactly(n)/readuntil(sep)已经帮你封装了这件事——这正是高层 API 的价值之一。 - 误区:
transport.write(data)会等数据真的发出去,所以不会有内存问题。write()是同步的非阻塞调用——它只是把数据放进 Transport 的写缓冲区就立刻返回,真正的发送由事件循环在 socket 可写时完成。如果对端读取很慢(慢客户端、网络拥塞)而你一直write,写缓冲会无限增长直到内存耗尽。asyncio 的解法是水位线 + 一对回调:缓冲超过high时框架调用protocol.pause_writing()、降到low以下时调resume_writing()——但框架只负责通知,真正停止产生数据是你的责任;asyncio.Protocol的默认实现是空方法,不重写就等于完全没有背压。这也是为什么用 Streams 时必须await writer.drain()——它内部等的就是resume_writing的信号。 - 误区:Protocol 的回调里抛异常,会被上层的
try/except捕获。 不会——因为回调根本没有「调用方」:它是事件循环在 IO 就绪时直接调用的,调用栈上没有你的业务代码。回调里抛出的异常会被送到loop.set_exception_handler()注册的处理器,没有注册的话只会用默认处理器打印到 stderr——在容器和生产日志流里极易被淹没,表现为「连接莫名其妙断了/协议状态错乱但没有任何错误信息」。所以每个回调内部都必须自己try/except:解析类错误直接transport.abort()断开(避免带着损坏的状态继续处理),未预期的异常记录完整堆栈。同时建议全局注册set_exception_handler做兜底上报。 - 误区:在
data_received里可以直接await处理数据。 回调是普通的同步函数,await在里面是语法错误。需要做异步处理只能asyncio.create_task(self._handle(msg)),但这会引入两个新问题:① 任务引用——create_task的返回值必须保存(事件循环只持有弱引用),否则任务可能被 GC 掉;而且connection_lost时应该 cancel 掉这些未完成的任务。② 失去背压——数据来得多快就创建多少个任务,对端一直发你就一直堆任务,内存和下游都会被打爆。解法是配合transport.pause_reading()(在途任务超过阈值时暂停读取,让 TCP 窗口收缩把压力传回发送方)或用Semaphore限制并发。这也是 Streams 更好写的原因之一:在协程里await处理是天然顺序的、自带背压。 - 误区:低层 API 比 Streams 快很多,所以性能敏感的服务都该用它。 差距只在框架开销占主导时才显著。实测数据:纯 echo(1KB 消息)场景 Protocol 比 Streams 快约 35%(110k vs 80k msg/s),但一旦加上真实业务逻辑(解析 JSON、查缓存,每条约 50μs),两者变成 19k vs 18k,差距缩到 5%。这和 uvloop 的道理完全一样:业务逻辑越重,框架层优化的相对收益越小。而低层 API 的复杂度代价是实打实的——拆包状态机、背压回调、错误处理、任务生命周期全要自己写,每一处都是潜在的 bug。所以正确的顺序是:先用 Streams 写对,profile 确认瓶颈真的在框架层,再考虑下沉;更常见的情况是换用已经优化好的高性能库(asyncpg、aiohttp)而不是自己写 Protocol。
- 追问:
eof_received()返回True和False有什么区别? 它在对端关闭了写方向(发送 FIN,即半关闭) 时被调用。返回False或None(默认):框架认为连接可以结束了,会关闭 transport,接着调用connection_lost。返回True:保持连接打开——你还可以继续transport.write()往对端发数据,直到自己调用close()。这个设计对应一类真实场景:请求-响应型协议中,客户端发完请求后立刻shutdown(SHUT_WR)表示「我说完了」,但服务端还需要把响应写回去(经典的例子是某些 HTTP/1.0 客户端和nc -N这类工具)。如果这时返回False,连接会被立即关闭、响应发不出去。需要注意的是:并非所有传输都支持半关闭(SSL/TLS 上的语义更复杂),而且返回True后你有责任在适当的时候关闭连接,否则会造成连接泄漏。 - 追问:
transport.close()和transport.abort()该怎么选?close()是优雅关闭:停止接收新数据,把写缓冲里剩余的数据全部发送出去,然后关闭连接、最后调用connection_lost(None)。abort()是立即关闭:丢弃写缓冲中未发送的数据,直接关闭 socket,同样会调用connection_lost(exc通常为None)。选择标准很清晰:正常的业务结束用close()(确保响应完整发出);检测到协议错误、非法数据、疑似攻击时用abort()——因为此时缓冲区里的数据可能是基于错误状态生成的,把它发出去只会让对端更困惑,而且攻击场景下你希望立刻释放资源而不是继续为对方服务。另外要记住:close()之后connection_lost不是立刻被调用的(要等缓冲刷完),期间连接仍处于关闭中状态,所以要用一个_closing标志避免重复处理。 - 追问:既然大多数场景该用 Streams,学习 Transport/Protocol 还有意义吗? 有,而且很实际。第一,它解释了 Streams 的行为:为什么
await writer.drain()有时立刻返回、有时挂起(因为它等的是resume_writing回调);为什么忘记drain()会导致内存暴涨(写缓冲无人节流);为什么readuntil有 64KB 的 limit(防止恶意的超长行耗尽内存)。第二,它是理解 TCP 编程的必经之路:粘包分包、半关闭、优雅关闭 vs 强制关闭、流控与 TCP 窗口的关系——这些认知在用任何高层库时都需要。第三,排查问题时用得上:transport.get_extra_info("socket")能拿到底层 socket 设置TCP_NODELAY/SO_KEEPALIVE,get_extra_info("peername")/"sslcontext"在调试连接问题时很有用(这些在 Streams 里通过writer.get_extra_info()同样可用)。第四,读第三方库源码(asyncpg、aiohttp 都基于 Protocol)时不至于卡住。
八、加强记忆
asyncio 的网络 API 是分层而不是并列的:Streams 就建立在 Protocol 之上(内部是 StreamReaderProtocol)。职责划分:Transport 由框架实现、代表「连接」——write() 是同步调用且不阻塞、close() 优雅关闭会刷完缓冲、abort() 立即关闭并丢弃缓冲、get_extra_info() 能拿到底层 socket;Protocol 由你实现、代表「协议逻辑」,框架按 connection_made → data_received×N → eof_received → connection_lost(一定最后、且一定会被调用,是清理的唯一可靠位置) 的顺序回调。本质区别是控制流方向:Streams 是拉模型(await 读,状态自然保存在协程栈里),Protocol 是推模型(数据来了框架调你,状态必须自己存在 self 上)。三个必须理解的点:① data_received 收到的绝不是「一条完整消息」——TCP 是字节流,会粘包也会分包,协议必须自己定义边界(长度前缀/分隔符/定长),拆包的要点是循环拆(一次可能含多条)+ 限制最大长度(防 OOM 攻击)+ 用 bytearray(bytes 拼接是 O(n²));② 背压是一对回调而不是 await——transport.write 不阻塞,写缓冲超过 high 水位时框架调 pause_writing()、降到 low 调 resume_writing(),但真正停下来是你的责任(默认实现是空的),读方向用 pause_reading() 让 TCP 窗口收缩把压力传回发送方;await writer.drain() 内部等的就是 resume_writing,所以「忘了 drain」和「不实现 pause_writing」是同一个错误;③ 回调是同步函数、不能 await,要异步处理只能 create_task(记得保存引用防 GC、connection_lost 时 cancel),但那样会失去背压,需要自己限流或 pause_reading()。异常处理是低层最容易踩的坑:回调里抛出的异常不会传播(根本没有调用方),只会进 loop.set_exception_handler,没注册就只打印到 stderr——所以每个回调都要自己 try/except,检测到协议错误用 abort() 而不是 close()。什么时候用低层:UDP(唯一选择)、极致吞吐(echo 场景快约 35%,但加上业务逻辑后只快 5%)、自定义二进制协议或框架开发(asyncpg、aiohttp 都是这么做的)、需要设置 TCP_NODELAY 等底层选项——绝大多数业务用 Streams 就对了,但理解这一层能解释 Streams 的所有行为。