Python 基础体系 · 第 53/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python asyncio 网络流:TCP、Reader、Writer、TLS 与背压

asyncio 是 Python 标准库中基于 async/await 的异步编程库,常用于 I/O 密集型和网络程序。它的核心执行模型不是“每个连接一个线程”,而是由事件循环在一个线程中轮流运行多个任务:任务运行到 await 时主动挂起,事件循环再调度其他就绪任务。一个任务如果长时间执行不让出控制权,同一事件循环中的其他任务和 I/O 都会被延迟。(docs.python.org)

本文聚焦 asyncio 的 TCP 流接口:

  • TCP 字节流究竟提供了什么保证;
  • StreamReader 如何读取数据;
  • StreamWriter 如何写出数据;
  • 为什么必须设计应用层协议帧;
  • TLS 如何建立、升级和关闭;
  • write()drain() 与背压之间的关系;
  • 超时、取消、半关闭、异常和资源回收如何组合;
  • 什么时候使用高层 Streams,什么时候下沉到 Transport/Protocol。

一、先建立整体模型:事件循环、任务与 TCP 流

1. 事件循环并不等于并行执行

事件循环可以抽象为:

while not stopped:
    ready = poll_os_io()
    run_callbacks_and_tasks(ready)

在单线程事件循环中,同一时刻只有一个 Python 回调或任务正在执行。假设有两个任务:

async def task_a():
    compute_for_a_long_time()
    await reader.read(1024)

async def task_b():
    await asyncio.sleep(0)
    print("task b")

如果 compute_for_a_long_time() 执行 5 秒,task_a() 在这 5 秒内没有遇到 await,那么 task_b() 不会因为“异步”而自动获得运行机会。

因此,异步 I/O 的基本条件不是“代码里出现了 async”,而是:

T阻塞0T_{\text{阻塞}} \approx 0

或者更准确地说,事件循环线程中的单次同步执行片段必须足够短。网络等待、定时器等待、队列等待可以通过 await 让出控制权;CPU 密集型计算、阻塞式文件操作、同步日志和同步 DNS 等操作则可能阻塞整个事件循环。官方文档明确建议把阻塞代码放入线程、解释器或进程执行器中。(docs.python.org)

2. Task 是并发单位,Stream 是连接视图

一个 TCP 连接通常对应一组对象:

Task
 ├── StreamReader  ← 从连接读取字节
 └── StreamWriter  → 向连接写入字节

asyncio.start_server() 接受客户端连接后,会向回调传入 (reader, writer);如果回调本身是协程函数,asyncio 会自动把它调度为 Task。客户端则通常通过 asyncio.open_connection() 获得同样的一对对象。(docs.python.org)

这几个概念需要区分:

概念 作用
Event loop 监控 I/O 和调度任务
Task 执行一个协程
TCP socket 操作系统提供的双向字节通道
StreamReader 对接收方向的异步读取封装
StreamWriter 对发送方向的异步写入封装
应用层协议 规定字节如何表示消息、请求、响应和错误

StreamReaderStreamWriter 不是消息队列,也不会自动知道你的业务消息边界。


二、TCP 是字节流,不是消息队列

1. TCP 提供有序、可靠的字节序列

从应用程序角度,可以把 TCP 连接抽象为:

B=b1b2b3bnB = b_1 b_2 b_3 \ldots b_n

发送方写入字节序列,接收方最终按相同顺序读取这些字节。TCP 负责:

  • 按顺序传输;
  • 检测丢失并重传;
  • 检测校验错误;
  • 在连接关闭或异常时报告结果。

但 TCP 不保留发送方的 write() 调用边界。

假设发送端执行:

writer.write(b"hello")
writer.write(b"world")
await writer.drain()

接收端可能看到:

b"helloworld"

也可能看到:

b"he"
b"lloworld"

或者:

b"hello"
b"world"

这些结果在 TCP 语义上完全等价,因为最终字节序列都是:

b"helloworld"

2. read(n) 不保证返回 n 字节

StreamReader.read(n) 的语义是“最多读取 n 字节”。当 n > 0 时,只要内部缓冲区中至少有一个字节可用,它就可能返回少于 n 的数据;如果连接已经 EOF 且缓冲区为空,则返回 b""。(docs.python.org)

因此,下面的代码不能表达“读取一个完整请求”:

data = await reader.read(1024)
request = decode_request(data)

它只表达:

当前最多取 1024 字节,取到多少算多少。

如果协议要求固定长度,应使用:

data = await reader.readexactly(8)

如果协议要求以换行结尾,应使用:

line = await reader.readline()

如果协议要求某个分隔符,应使用:

payload = await reader.readuntil(b"\r\n\r\n")

3. 正确性依赖“消息边界条件”

从 TCP 字节流恢复消息,必须满足下面至少一个条件:

  1. 消息长度固定;
  2. 消息携带长度字段;
  3. 消息由明确分隔符结束;
  4. 消息生命周期由连接 EOF 结束;
  5. 使用外层协议定义的解析规则。

可以形式化为:

字节流+边界规则消息序列\text{字节流} + \text{边界规则} \rightarrow \text{消息序列}

只有字节流而没有边界规则时,消息无法唯一确定。

例如字节流:

PINGPONG

可能表示:

["PING", "PONG"]

也可能表示:

["PINGPONG"]

甚至是:

["P", "INGP", "ONG"]

TCP 本身无法替你做出判断。


三、StreamReader:从字节流恢复协议数据

1. read():读取任意长度

data = await reader.read(4096)
if not data:
    print("对端已关闭连接")

这里的 not data 通常表示读到了 EOF,但只有在本次读取返回 b"" 时才成立。空消息和 EOF 不是同一个概念:

  • b"":没有任何字节,通常表示 EOF;
  • 非空 bytes:成功读取了一段数据;
  • read(0):立即返回 b"",不表示连接关闭。

官方文档规定,read(-1) 会一直读取到 EOF 后返回全部数据;如果 EOF 已到达且缓冲区为空,则返回空字节串。(docs.python.org)

read(-1) 适合“连接关闭才算数据结束”的场景,例如:

async def receive_whole_file(reader):
    return await reader.read(-1)

但它不适合长期连接。若对端保持连接打开,read(-1) 会一直等待,直到连接关闭。

2. readexactly(n):固定长度协议

假设协议头固定为 8 字节:

0       4       8
+-------+-------+
| magic | length|
+-------+-------+

其中:

  • magic:4 字节魔数;
  • length:4 字节无符号大端整数;
  • 后面跟随 length 字节正文。

读取过程必须严格按协议长度推进:

import struct

HEADER = struct.Struct("!4sI")
MAX_PAYLOAD = 1024 * 1024


async def read_frame(reader):
    header = await reader.readexactly(HEADER.size)
    magic, length = HEADER.unpack(header)

    if magic != b"WR01":
        raise ValueError(f"invalid magic: {magic!r}")

    if length > MAX_PAYLOAD:
        raise ValueError(f"payload too large: {length}")

    payload = await reader.readexactly(length)
    return payload

!4sI 中:

  • ! 表示网络字节序,即大端;
  • 4s 表示 4 字节字符串;
  • I 表示 4 字节无符号整数。

如果对端在正文未收完时关闭连接,readexactly() 会抛出 asyncio.IncompleteReadError,部分已读取数据位于异常的 partial 属性中。(docs.python.org)

try:
    payload = await read_frame(reader)
except asyncio.IncompleteReadError as exc:
    print(f"帧不完整,只收到 {len(exc.partial)} 字节")

这是和“把异常简单转换成 EOF”不同的故障语义:

  • 正常 EOF:对端在消息边界关闭;
  • 不完整帧:对端在消息中间关闭,协议数据损坏或连接异常。

3. readline():行协议

行协议以 \n 作为消息结束标记:

async def read_command(reader):
    line = await reader.readline()

    if not line:
        return None

    if len(line) > 8192:
        raise ValueError("command line too long")

    return line.rstrip(b"\r\n").decode("utf-8")

需要注意,readline() 返回的内容包括结尾的换行符。如果 EOF 到达前没有遇到 \n,它会返回已经读取到的部分数据;如果缓冲区为空,则返回 b""。(docs.python.org)

行协议简单,但必须处理两个边界:

  1. 客户端永远不发送换行符;
  2. 客户端发送无限长的一行。

第二种情况会造成内存和连接占用,因此不能只依赖 readline()。可以在协议层设置最大行长,或者使用带限制的读取方案。

4. readuntil():分隔符协议与限制

data = await reader.readuntil(b"\r\n\r\n")

它会读取直到找到分隔符,并将分隔符一起返回。若读取量超过 StreamReader 的配置限制,会抛出 LimitOverrunError,且已读取数据仍保留在内部缓冲区中,可以重新读取。若 EOF 在分隔符出现前到达,则抛出 IncompleteReadError。(docs.python.org)

Python 3.13 起,separator 还可以是分隔符元组:

data = await reader.readuntil((b"\n", b"\r\n"))

但生产协议通常不应仅凭多个分隔符“猜测”格式。协议解析器应明确:

  • 哪些分隔符合法;
  • 哪个分隔符优先;
  • 最大消息长度是多少;
  • 分隔符是否属于消息正文;
  • EOF 出现在何处时代表正常结束。

5. limit 的含义不是“连接总内存上限”

创建连接时可以设置:

reader, writer = await asyncio.open_connection(
    host,
    port,
    limit=64 * 1024,
)

limit 是返回的 StreamReader 使用的缓冲区限制,默认值为 64 KiB。服务器端 start_server() 也有同名参数。(docs.python.org)

它不是:

  • 单个 TCP 连接所有对象的总内存上限;
  • StreamWriter 写缓冲上限;
  • 业务消息大小限制;
  • 所有连接的全局内存上限。

因此仍然需要在协议层检查长度字段:

if length > MAX_PAYLOAD:
    raise ValueError("frame too large")

四、StreamWriter:写入、排空与关闭

1. write() 通常不等待网络发送完成

StreamWriter.write(data) 会尝试立即写入底层 socket;如果不能立即完成,数据会排入内部写缓冲区。它接受 bytesbytearray 或符合要求的 memoryview。官方文档建议将 write()drain() 配合使用。(docs.python.org)

writer.write(b"hello\n")
await writer.drain()

write() 返回值为 None,不能用返回值判断“发送了多少字节”:

# 错误理解
sent = writer.write(data)
print(sent)  # None

即使 drain() 成功,也只表示当前写入过程没有报告错误,并不意味着远端应用已经读取或处理了数据。数据可能仍处于:

应用 write()
    ↓
asyncio 写缓冲
    ↓
操作系统 socket 缓冲
    ↓
TCP 发送队列
    ↓
网络
    ↓
远端内核
    ↓
远端应用

drain() 不能提供应用层确认。若业务需要确认,必须设计 ACK 或响应消息。

2. drain() 是写方向的流量控制点

设写缓冲区当前大小为 BB,高水位为 HH,低水位为 LL,其中:

L<HL < H

当:

B<HB < H

调用 await writer.drain() 通常立即返回;当缓冲区达到高水位后,drain() 会等待缓冲区下降到低水位附近,之后写入任务才能继续。官方文档将其定义为与底层 I/O 写缓冲交互的流控制方法。(docs.python.org)

数据流可以表示为:

flowchart LR
    A[业务生产者] -->|write| B[asyncio 写缓冲]
    B --> C[操作系统发送缓冲]
    C --> D[TCP]
    D --> E[网络]
    E --> F[对端接收缓冲]
    B -->|达到高水位| G[drain 挂起]
    F -->|持续读取| H[缓冲逐渐下降]
    H -->|低于低水位| G

关键因果关系是:

  1. 生产者调用 write()
  2. 消费者读取速度较慢;
  3. 本端和内核发送缓冲逐渐积累;
  4. 写缓冲达到高水位;
  5. drain() 挂起;
  6. 生产者停止继续产生数据;
  7. 对端读取后,缓冲下降;
  8. drain() 恢复。

这就是背压。

3. 背压解决的是“发送过快”,不是所有内存问题

错误写法:

async def bad_send(writer, messages):
    for message in messages:
        writer.write(message)

如果 messages 很多,代码会持续把数据塞入写缓冲,而没有等待流控点。更合理的写法是:

async def send_messages(writer, messages):
    for message in messages:
        writer.write(message)
        await writer.drain()

但这仍然不是完整的内存控制,因为 messages 本身可能已经是一个巨大列表。更完整的生产者—消费者模型应限制队列长度:

import asyncio


async def producer(queue):
    for i in range(100_000):
        await queue.put(f"event-{i}\n".encode())

    await queue.put(None)


async def writer_loop(queue, writer):
    while True:
        item = await queue.get()
        try:
            if item is None:
                return

            writer.write(item)
            await writer.drain()
        finally:
            queue.task_done()

这里有两层背压:

  • asyncio.Queue(maxsize=N) 限制业务生产者积压;
  • writer.drain() 限制网络写缓冲积压。

如果省略第一层,业务对象仍可能在内存中堆积;如果省略第二层,队列消费速度可能超过网络实际发送速度。

4. drain() 不一定主动让出事件循环

当写缓冲没有达到高水位时,drain() 可以直接返回而不挂起当前任务。官方文档特别指出,反复执行 write()await drain() 的代码,在没有背压时可能持续占用事件循环;需要时可以显式使用 await asyncio.sleep(0) 让出执行权。(docs.python.org)

例如广播场景:

async def broadcast(writer, payloads):
    for payload in payloads:
        writer.write(payload)
        await writer.drain()

        # 当单个任务可能长时间连续发送时,主动让出控制权
        await asyncio.sleep(0)

这不是每次写入都必须加入的固定模板。它的适用条件是:

  • 单个任务可能连续处理大量数据;
  • 每次 drain() 都立即返回;
  • 其他任务需要获得及时调度;
  • 业务允许发送任务在消息之间短暂让出。

五、一个完整的长度前缀 TCP 协议示例

下面实现一个可运行的 TCP echo 服务。协议格式为:

+--------+--------+----------------+
| 4 字节 | 4 字节 | length 字节    |
| WR01   | length | payload        |
+--------+--------+----------------+

1. 服务端

# server.py
import asyncio
import struct

HEADER = struct.Struct("!4sI")
MAGIC = b"WR01"
MAX_PAYLOAD = 1024 * 1024


async def read_frame(reader: asyncio.StreamReader) -> bytes:
    header = await reader.readexactly(HEADER.size)
    magic, length = HEADER.unpack(header)

    if magic != MAGIC:
        raise ValueError(f"invalid magic: {magic!r}")

    if length > MAX_PAYLOAD:
        raise ValueError(f"payload too large: {length}")

    return await reader.readexactly(length)


async def write_frame(
    writer: asyncio.StreamWriter,
    payload: bytes,
) -> None:
    if len(payload) > MAX_PAYLOAD:
        raise ValueError(f"payload too large: {len(payload)}")

    writer.write(HEADER.pack(MAGIC, len(payload)))
    writer.write(payload)
    await writer.drain()


async def handle_client(
    reader: asyncio.StreamReader,
    writer: asyncio.StreamWriter,
) -> None:
    peer = writer.get_extra_info("peername")
    print(f"connected: {peer}")

    try:
        while True:
            try:
                payload = await asyncio.wait_for(
                    read_frame(reader),
                    timeout=300,
                )
            except asyncio.TimeoutError:
                print(f"idle timeout: {peer}")
                break
            except asyncio.IncompleteReadError:
                print(f"peer closed during frame: {peer}")
                break

            print(f"received {payload!r} from {peer}")
            await write_frame(writer, payload)

    except (ConnectionError, ValueError) as exc:
        print(f"connection error from {peer}: {exc}")
    finally:
        writer.close()
        await writer.wait_closed()
        print(f"closed: {peer}")


async def main() -> None:
    server = await asyncio.start_server(
        handle_client,
        host="127.0.0.1",
        port=8888,
        limit=64 * 1024,
    )

    addresses = ", ".join(
        str(sock.getsockname())
        for sock in server.sockets or []
    )
    print(f"serving on {addresses}")

    async with server:
        await server.serve_forever()


if __name__ == "__main__":
    asyncio.run(main())

启动:

python server.py

预期输出类似:

serving on ('127.0.0.1', 8888)

asyncio.start_server() 启动监听 socket,并在每个连接建立时向 handle_client() 传入 StreamReaderStreamWriter。使用 async with server 可以确保离开上下文时关闭服务器;serve_forever() 则持续处理连接。(docs.python.org)

2. 客户端

# client.py
import asyncio
import struct

HEADER = struct.Struct("!4sI")
MAGIC = b"WR01"


async def write_frame(writer, payload: bytes) -> None:
    writer.write(HEADER.pack(MAGIC, len(payload)))
    writer.write(payload)
    await writer.drain()


async def read_frame(reader) -> bytes:
    header = await reader.readexactly(HEADER.size)
    magic, length = HEADER.unpack(header)

    if magic != MAGIC:
        raise ValueError("bad magic")

    return await reader.readexactly(length)


async def main() -> None:
    reader, writer = await asyncio.open_connection(
        "127.0.0.1",
        8888,
    )

    try:
        for message in ("hello", "asyncio", "tcp"):
            payload = message.encode("utf-8")
            await write_frame(writer, payload)

            response = await asyncio.wait_for(
                read_frame(reader),
                timeout=5,
            )
            print("response:", response.decode("utf-8"))

    finally:
        writer.close()
        await writer.wait_closed()


if __name__ == "__main__":
    asyncio.run(main())

启动:

python client.py

预期输出:

response: hello
response: asyncio
response: tcp

这个例子中,readexactly() 的必要性来自协议头和正文长度之间的因果关系:

  1. 先读取固定 8 字节头;
  2. 从头部解析出 length
  3. 校验 length 是否超过协议允许值;
  4. 再读取恰好 length 字节;
  5. 这些字节才构成一个完整帧。

如果把第 1 步换成 read(8),即使通常能读到 8 字节,也不能把它当作 TCP 或 Streams 的保证。


六、连接生命周期:建立、使用、关闭与半关闭

1. 建立连接

客户端:

reader, writer = await asyncio.open_connection(
    "example.com",
    443,
    ssl=True,
)

服务端:

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

open_connection() 返回 (reader, writer)start_server() 的连接回调接收同样的一对对象。两者默认的 StreamReader 限制都是 64 KiB。(docs.python.org)

2. 正常关闭

推荐关闭流程:

writer.close()
await writer.wait_closed()

close() 关闭流和底层 socket;wait_closed() 等待底层连接真正关闭,并确保关闭前的缓冲数据得到处理。官方文档建议在 close() 后调用 wait_closed()。(docs.python.org)

因此不要只写:

writer.close()

尤其在短生命周期客户端中,主协程可能马上结束,程序退出时尚未完成的关闭过程和缓冲处理可能无法按预期完成。

3. 半关闭:只关闭写方向

TCP 是双向连接。关闭写方向而继续读取,称为半关闭:

if writer.can_write_eof():
    writer.write_eof()

can_write_eof() 用于判断底层传输是否支持 write_eof()write_eof() 会在已缓冲的写数据刷新后关闭写端。(docs.python.org)

典型用途是:

客户端:发送完请求
客户端:半关闭写端
服务端:读到 EOF,知道请求结束
服务端:继续发送响应
客户端:继续读取响应

但半关闭并非所有传输都支持。特别是 TLS 连接通常不能简单等同于普通 TCP 半关闭,因此必须先判断 can_write_eof(),并根据协议定义决定是否允许这种状态。

4. 关闭状态不是一瞬间完成的

writer.is_closing() 在连接已关闭或正在关闭时返回 True。(docs.python.org)

在并发发送场景中,可能出现:

发送任务 A:write()
连接管理任务:close()
发送任务 B:drain()

因此共享一个 StreamWriter 时,应明确所有权:

  • 哪个任务负责关闭;
  • 关闭后是否禁止新消息;
  • 发送任务如何感知连接已关闭;
  • 关闭时是否取消排队中的发送任务。

仅依赖 is_closing() 不能消除竞态,因为检查和写入之间仍可能发生状态变化。


七、超时与取消:必须区分等待超时和连接关闭

1. 对一次读取设置超时

Python 3.14 中可以使用 asyncio.wait_for()

try:
    data = await asyncio.wait_for(
        reader.readexactly(8),
        timeout=5,
    )
except asyncio.TimeoutError:
    print("5 秒内没有收到完整头部")

该超时针对的是这次等待。如果读取已经收到部分数据后超时,协议状态可能处于“半帧”状态。此时不能总是简单继续读取,必须由协议决定:

  • 保留部分数据并继续;
  • 发送错误响应;
  • 直接关闭连接。

对于长度帧,最安全的策略通常是:头部或正文读取超时后关闭连接,因为连接中已经存在一个未完成的协议帧。

2. 空闲超时和操作超时不是一回事

下面的代码设置的是“单次 read 操作最多等待 300 秒”:

while True:
    frame = await asyncio.wait_for(
        read_frame(reader),
        timeout=300,
    )

它不是整个连接最多存活 300 秒。每次成功收到一帧后,计时器都会重新开始。

两种策略语义不同:

连接生命周期超时:
连接建立后最多存在 N 秒

空闲超时:
连续 N 秒没有收到完整业务消息就关闭

单次操作超时:
某个连接操作最多等待 N 秒

如果把这三者混在一起,容易造成连接提前断开,或者慢客户端无限占用资源。

3. 取消不是异常清理的替代品

任务可能因为超时、上层关闭或服务停止而收到 CancelledError。无论是普通异常还是取消,都应在 finally 中释放 writer:

async def handle_client(reader, writer):
    try:
        await serve_protocol(reader, writer)
    except asyncio.CancelledError:
        raise
    except (ConnectionError, ValueError) as exc:
        print("protocol failure:", exc)
    finally:
        writer.close()
        await writer.wait_closed()

这里重新抛出 CancelledError 很重要:清理完成后仍需保留取消语义,让上层知道任务确实被取消,而不是被误判为正常完成。


八、TLS:在 TCP 流上增加身份与机密性

1. TLS 不改变应用层“仍然是字节流”

TLS 建立在 TCP 之上:

应用协议帧
    ↓
TLS 记录层
    ↓
TCP 字节流
    ↓
网络

TLS 提供:

  • 加密;
  • 完整性保护;
  • 服务器身份认证;
  • 在适当配置下的客户端身份认证。

但 TLS 不会给应用消息自动加边界。对应用层而言,依然必须使用 readexactly()readline()readuntil() 或自定义协议解析器。

2. 客户端使用 HTTPS/TLS 连接

最简单的客户端方式:

import asyncio
import ssl


async def main():
    context = ssl.create_default_context()

    reader, writer = await asyncio.open_connection(
        "example.com",
        443,
        ssl=context,
        server_hostname="example.com",
        ssl_handshake_timeout=10,
        ssl_shutdown_timeout=10,
    )

    try:
        request = (
            b"GET / HTTP/1.1\r\n"
            b"Host: example.com\r\n"
            b"Connection: close\r\n"
            b"\r\n"
        )
        writer.write(request)
        await writer.drain()

        while True:
            chunk = await reader.read(4096)
            if not chunk:
                break
            print(chunk.decode("latin-1"), end="")
    finally:
        writer.close()
        await writer.wait_closed()


asyncio.run(main())

ssl=True 可以使用默认 TLS 配置;更严谨的客户端应显式使用 ssl.create_default_context(),并传入正确的 server_hostname,这样证书主机名校验才对应目标主机。

open_connection() 支持 ssl_handshake_timeoutssl_shutdown_timeout;Python 3.11 增加了关闭超时参数。(docs.python.org)

3. 服务端启用 TLS

import asyncio
import ssl


async def handle_client(reader, writer):
    try:
        writer.write(b"hello over tls\n")
        await writer.drain()
    finally:
        writer.close()
        await writer.wait_closed()


async def main():
    context = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
    context.load_cert_chain(
        certfile="server.crt",
        keyfile="server.key",
    )

    server = await asyncio.start_server(
        handle_client,
        "127.0.0.1",
        8443,
        ssl=context,
        ssl_handshake_timeout=10,
        ssl_shutdown_timeout=10,
    )

    async with server:
        await server.serve_forever()


asyncio.run(main())

前置条件:

  • server.crt 是服务端证书;
  • server.key 是对应私钥;
  • 私钥文件权限应限制为服务进程可读;
  • 客户端必须信任该证书或其签发链;
  • 生产环境不能为了绕过证书问题而关闭验证。

启用 ssl 后,连接回调收到的仍然是 StreamReaderStreamWriter,应用层协议代码通常无需改变。这是 Streams 抽象的价值:将 TCP 或 TLS 连接统一成异步字节流接口。

4. 连接建立时 TLS 与连接建立后升级

有两种不同场景:

场景一:连接建立时直接 TLS

reader, writer = await asyncio.open_connection(
    host,
    port,
    ssl=context,
    server_hostname=host,
)

服务端:

server = await asyncio.start_server(
    handle_client,
    host,
    port,
    ssl=context,
)

这是最常见的 TLS TCP 服务方式。

场景二:已有连接升级为 TLS

StreamWriter.start_tls() 可以将已有的基于 Streams 的连接升级为 TLS。该方法要求传入配置好的 ssl.SSLContext,并支持服务器主机名、握手超时和关闭超时参数;该接口在 Python 3.11 增加。(docs.python.org)

await writer.start_tls(
    context,
    server_hostname="example.com",
    ssl_handshake_timeout=10,
    ssl_shutdown_timeout=10,
)

这种模式适合协议明确规定“先以明文协商,再切换 TLS”的场景,例如某些 STARTTLS 风格协议。它不是“随时把任意 TCP 连接改成安全连接”的通用魔法:

  1. 客户端和服务端必须在完全相同的协议状态切换;
  2. 升级前必须完成明文协议阶段;
  3. 升级时不能还有未处理的应用层数据;
  4. 两端必须对 TLS 握手的开始位置达成一致。

否则一端会把 TLS ClientHello 当普通业务数据解析,另一端则把普通文本当作 TLS 握手,最终表现为握手失败、协议错误或连接关闭。


九、背压的完整推导:从速度不匹配到资源耗尽

设:

  • PP:业务生产者产生数据的平均速率;
  • CC:网络和对端能够消费数据的平均速率;
  • Q(t)Q(t):时刻 tt 的待发送数据量。

当:

P>CP > C

积压量会增长:

dQdt=PC>0\frac{dQ}{dt} = P - C > 0

如果没有上限,经过时间 TT 后:

Q(T)=Q(0)+(PC)TQ(T) = Q(0) + (P-C)T

这意味着即使 TCP 连接本身没有立即报错,进程内存也可能不断增长。

例如:

  • 生产速率:10 MB/s;
  • 实际发送速率:2 MB/s;
  • 速度差:8 MB/s;
  • 持续 60 秒。

则理论积压量约为:

8×60=480 MB8 \times 60 = 480\text{ MB}

这不是 asyncio 特有问题,而是所有异步发送系统都必须面对的排队问题。

1. 没有背压的广播

async def bad_broadcast(clients, payload):
    for writer in clients:
        writer.write(payload)

问题有三层:

  1. 慢客户端的写缓冲会增长;
  2. 一个连接写失败可能在后续才暴露;
  3. 多个客户端的待发送数据总量没有统一上限。

2. 对每个客户端独立发送

async def send_to_client(writer, payload):
    try:
        writer.write(payload)
        await writer.drain()
    except ConnectionError:
        writer.close()
        await writer.wait_closed()

广播时:

await asyncio.gather(
    *(send_to_client(writer, payload) for writer in clients),
    return_exceptions=True,
)

这样一个慢客户端只会让自己的发送任务等待,不会直接阻塞其他客户端。但是,如果广播频率非常高,仍可能为每次消息创建大量任务,或在每个客户端内部积累未完成的发送操作。

3. 更清晰的客户端发送队列

每个客户端拥有一个有限队列和一个专属发送任务:

import asyncio
from dataclasses import dataclass


@dataclass
class Client:
    writer: asyncio.StreamWriter
    queue: asyncio.Queue[bytes | None]


async def client_sender(client: Client) -> None:
    try:
        while True:
            payload = await client.queue.get()
            try:
                if payload is None:
                    return

                client.writer.write(payload)
                await client.writer.drain()
            finally:
                client.queue.task_done()
    except (ConnectionError, asyncio.CancelledError):
        raise


async def publish(client: Client, payload: bytes) -> bool:
    try:
        client.queue.put_nowait(payload)
        return True
    except asyncio.QueueFull:
        return False

QueueFull 的处理策略必须是协议和业务的一部分:

  • 丢弃最旧消息;
  • 丢弃最新消息;
  • 合并多条状态消息;
  • 暂停生产者;
  • 降级客户端;
  • 关闭慢客户端。

实时状态广播常常允许“丢旧状态但保留最新状态”;审计日志和支付事件则通常不能丢弃,应该使用持久化队列或明确的失败重试机制。不能把所有消息都套用同一种策略。


十、背压不是 sleep():两个反例

反例一:用固定 sleep 代替 drain()

async def wrong_send(writer, messages):
    for message in messages:
        writer.write(message)
        await asyncio.sleep(0.01)

sleep(0.01) 只是让出事件循环,并没有观察写缓冲状态:

  • 网络快时,可能无谓降低吞吐;
  • 网络慢时,仍然可能持续积压;
  • 不同连接的实际消费能力不同,固定延迟无法适配。

正确的流控点是:

writer.write(message)
await writer.drain()

必要时再额外增加队列长度和业务级速率控制。

反例二:只调用 drain(),不限制消息大小

async def wrong_frame(writer, payload):
    writer.write(payload)
    await writer.drain()

如果 payload 本身是 2 GB,drain() 不会把它变成安全消息。它只能控制写缓冲何时允许继续写入,不能替你检查业务对象的大小。

长度前缀协议必须在写入前校验:

if len(payload) > MAX_PAYLOAD:
    raise ValueError("payload too large")

十一、并发读取与并发写入的所有权问题

1. 一个连接通常应只有一个读取循环

错误结构:

asyncio.create_task(read_command(reader))
asyncio.create_task(read_response(reader))

两个任务同时从同一个 StreamReader 读取,会产生竞争:哪个任务先读到字节取决于调度时机,协议边界会被破坏。应用层应让一个读取循环拥有 reader,然后根据解析出的消息分发给其他任务:

async def read_loop(reader, inbox):
    while True:
        frame = await read_frame(reader)
        await inbox.put(frame)

2. 多个任务写入时需要串行化协议帧

多个任务同时执行:

writer.write(header)
writer.write(payload)

可能出现交错:

任务 A:写入 A 的 header
任务 B:写入 B 的 header
任务 A:写入 A 的 payload
任务 B:写入 B 的 payload

如果协议要求一个 header 紧跟一个 payload,就会损坏帧结构。

可用锁保护一次完整帧:

class FramedWriter:
    def __init__(self, writer):
        self.writer = writer
        self.lock = asyncio.Lock()

    async def send(self, payload: bytes):
        async with self.lock:
            self.writer.write(HEADER.pack(MAGIC, len(payload)))
            self.writer.write(payload)
            await self.writer.drain()

更常见的做法是只允许一个发送任务拥有 writer,其他任务通过有限队列提交消息。这样不仅避免交错,也更容易实现背压、关闭和失败通知。


十二、EOF、连接重置与协议错误

网络连接结束至少有几种不同路径:

1. 正常 EOF

data = await reader.read(4096)
if data == b"":
    # 对端完成关闭写端,且本地缓冲已读完
    return

如果协议以 EOF 表示消息结束,必须确认这确实是合法边界。

2. 不完整帧

try:
    payload = await read_frame(reader)
except asyncio.IncompleteReadError as exc:
    # 头或正文只收到一部分
    print("incomplete frame:", exc.partial)

这是协议层错误或异常关闭,不应当把部分数据当成完整消息处理。

3. 对端重置连接

写入或排空时可能出现:

except ConnectionResetError:
    ...
except BrokenPipeError:
    ...

也可以将它们归入 ConnectionError 处理,但日志中最好记录具体异常类型,因为“对端主动关闭”“网络重置”和“本端关闭时竞态”具有不同诊断价值。

4. 本地任务取消

服务停止时,处理任务可能被取消。清理代码必须仍然关闭 writer,但不能吞掉取消信号:

except asyncio.CancelledError:
    raise
finally:
    writer.close()
    await writer.wait_closed()

十三、低层 Transport/Protocol 与 Streams 的关系

Streams 是高层接口,官方文档将其描述为用于网络连接的高层异步原语,可以在不直接使用回调、低层协议和传输对象的情况下收发数据。(docs.python.org)

低层结构可以表示为:

事件循环
   ↓
Transport  ← 字节传输、写缓冲、流控
   ↕
Protocol   ← connection_made、data_received、connection_lost

Streams 则在其上封装出:

StreamReader + StreamWriter

适合使用 Streams 的场景

  • 普通 TCP 客户端和服务器;
  • 自定义请求—响应协议;
  • 行协议和长度前缀协议;
  • TLS TCP 服务;
  • 连接数中等、协议状态机不需要极致优化的程序。

适合使用 Transport/Protocol 的场景

  • 需要回调驱动的高吞吐协议实现;
  • 需要直接控制传输层流控;
  • 需要实现复杂协议状态机;
  • 需要与已有基于 Protocol 的 asyncio 组件集成;
  • 需要减少高层对象封装带来的管理复杂度。

下沉到低层并不会消除 TCP 的字节流性质,也不会自动提供消息边界。data_received(data) 每次收到的 data 仍然只是任意长度的一段字节,协议解析器仍必须维护自己的缓冲区和状态。


十四、诊断:确认问题发生在哪一层

遇到“异步 TCP 很慢”或“消息丢失”时,不应直接归因于 asyncio。可以沿数据路径逐层排查:

1. 应用层是否正确分帧

检查:

  • 是否错误地把一次 read() 当成一次消息;
  • 长度字段是否使用统一字节序;
  • 是否校验最大帧长度;
  • 是否正确处理半帧和 EOF;
  • 多个写任务是否可能交错。

2. 任务是否真的让出了事件循环

打开调试模式:

PYTHONASYNCIODEBUG=1 python server.py

或者:

asyncio.run(main(), debug=True)

asyncio 调试模式可以报告错误线程调用、耗时过长的回调等问题;默认情况下,超过 100 毫秒的回调会被视为慢回调并记录,可通过 loop.slow_callback_duration 调整。(docs.python.org)

如果日志显示某个回调长时间运行,问题可能是:

  • 同步 CPU 计算;
  • 阻塞式文件或数据库调用;
  • 同步网络日志;
  • 不必要的大规模 JSON 编解码;
  • 某个循环没有 await 或显式让出。

3. 写缓冲是否持续触发背压

可以记录发送耗时:

import time


async def timed_send(writer, data):
    writer.write(data)

    started = time.monotonic()
    await writer.drain()
    elapsed = time.monotonic() - started

    if elapsed > 0.1:
        print(f"drain waited {elapsed:.3f}s")

这不能替代底层网络监控,但可以帮助区分:

  • 业务生产慢;
  • 编码和序列化慢;
  • drain() 经常等待,说明发送路径受到流控;
  • drain() 从不等待但内存增长,说明可能缺少队列上限、单次写入过大或生产任务没有正确串行化。

4. 是否存在未等待的协程和未消费的任务异常

以下代码不会执行 work()

async def main():
    work()

必须:

await work()

或者:

task = asyncio.create_task(work())
await task

如果创建的 Task 异常却没有被等待或读取,asyncio 可能记录“Task exception was never retrieved”。官方调试文档也将未等待协程和未消费异常列为常见错误。(docs.python.org)


十五、生产代码中的状态机

一个 TCP/TLS 流连接可以用如下状态描述:

stateDiagram-v2
    [*] --> Connecting
    Connecting --> TLSHandshake: ssl enabled
    Connecting --> Established: plain TCP
    TLSHandshake --> Established: handshake success
    TLSHandshake --> Closing: handshake failure
    Established --> Reading
    Reading --> Reading: complete frame
    Reading --> Writing: request requires response
    Writing --> Reading: drain completed
    Reading --> HalfClosed: peer EOF
    HalfClosed --> Closing: response finished
    Reading --> Closing: protocol error
    Writing --> Closing: connection error
    Established --> Closing: timeout/cancel/shutdown
    Closing --> Closed
    Closed --> [*]

每个状态都有明确含义:

  • Connecting:TCP 连接尚未完成;
  • TLSHandshake:TCP 已连接,但应用数据不能提前当作明文协议解析;
  • Established:可以进入应用协议;
  • Reading:等待或解析输入帧;
  • Writing:写入响应,并可能在 drain() 上等待;
  • HalfClosed:一侧已结束发送,但另一方向仍可能有数据;
  • Closing:禁止接收新的业务消息,清理已有资源;
  • Closed:socket 和相关任务均已终止。

实际代码不一定需要显式写出枚举状态,但错误处理必须体现这些状态差异。例如 TLS 握手失败不能作为普通业务帧解析错误;不完整帧不能当作正常 EOF;写缓冲排空超时也不能继续接受无限新的业务消息。


十六、核心取舍

使用 read() 的条件

适合:

  • 原始字节处理;
  • 上层已有完整缓冲区协议;
  • 读取到任意数据即可处理。

不适合:

  • 需要一次得到完整业务消息;
  • 将返回长度误认为消息长度。

使用 readexactly() 的条件

适合:

  • 固定长度字段;
  • 长度前缀协议;
  • 已知剩余数据长度。

代价是:对端不完整发送时会抛出 IncompleteReadError,调用者必须处理。

使用 readline()readuntil() 的条件

适合:

  • 文本行协议;
  • 明确的分隔符协议。

必须设置并执行长度限制,否则恶意或异常对端可以通过永不发送分隔符制造无限等待和内存压力。

使用 TLS 的条件

适合:

  • 传输中包含凭据、令牌或隐私数据;
  • 需要验证服务端身份;
  • 网络边界不可信。

必须配置证书验证、主机名和握手/关闭超时。TLS 只保护传输,不负责认证业务用户,也不负责应用消息格式。

使用 drain() 的条件

只要程序持续写入 StreamWriter,就应把 drain() 作为写路径的一部分。它能让写入路径感知底层缓冲区的背压,但不替代:

  • 有界业务队列;
  • 消息大小限制;
  • 慢客户端策略;
  • 应用层 ACK;
  • 连接级超时。

asyncio Streams 的抽象很小:StreamReader 读取字节,StreamWriter 写入字节,Task 在等待 I/O 时让出事件循环。真正决定程序是否正确的,是围绕这几个字节操作建立的协议边界、生命周期、并发所有权、TLS 状态和背压策略。

TCP 不提供消息边界,所以必须分帧;write() 不等于远端已收到,所以需要 drain() 和必要时的业务确认;drain() 只能控制传输缓冲,不能替代有界队列;TLS 改变的是安全属性,不改变流式读取模型;关闭、取消、EOF 和不完整帧则必须分别处理。掌握这些因果关系后,asyncio 的网络代码才不会停留在“能连上、能收发”的示例层面,而能形成可验证的协议实现。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。