← 返回题目列表

asyncio 的 Transport 和 Protocol 是什么?什么时候要用低层 API?

困难 第 27 / 27 题 更新于 2026/08/01
TransportProtocol低层API自定义协议

简化版

asyncio 提供了两层网络 API:高层的 StreamsStreamReader/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_madedata_receivedeof_receivedconnection_lost 以及背压相关的 pause_writing/resume_writing本质区别在控制流方向:Streams 是拉模型(你主动 await 读,状态自然保存在协程栈里),Protocol 是推模型(数据来了框架调你,状态必须自己存在 self)。还有一个实用信息:UDP 只有低层 APIcreate_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 一定最后一个且一定会被调用(所以它是清理的唯一可靠位置excNone 表示正常关闭、非 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。实践清单里最关键的六条:默认用 Streamsdata_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() 返回 TrueFalse 有什么区别? 它在对端关闭了写方向(发送 FIN,即半关闭) 时被调用。返回 FalseNone(默认):框架认为连接可以结束了,会关闭 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_lostexc 通常为 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_KEEPALIVEget_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_madedata_received×N → eof_receivedconnection_lost(一定最后、且一定会被调用,是清理的唯一可靠位置) 的顺序回调。本质区别是控制流方向Streams 是拉模型await 读,状态自然保存在协程栈里),Protocol 是推模型(数据来了框架调你,状态必须自己存在 self)。三个必须理解的点data_received 收到的绝不是「一条完整消息」——TCP 是字节流,会粘包也会分包,协议必须自己定义边界(长度前缀/分隔符/定长),拆包的要点是循环拆(一次可能含多条)+ 限制最大长度(防 OOM 攻击)+ 用 bytearraybytes 拼接是 O(n²))② 背压是一对回调而不是 await——transport.write 不阻塞,写缓冲超过 high 水位时框架调 pause_writing()、降到 lowresume_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 的所有行为。