云计算百科
云计算领域专业知识百科平台

Python Asyncio 网络编程实战指南:从 TCP/UDP 服务器到 SSL/TLS 与连接池(pysheeet)

  • 文档
  • 教程
  • 开发工具

【免费下载链接】pysheeet

Python Cheat Sheet

项目地址:
https://gitcode.com/gh_mirrors/py/pysheeet

点击查看 免费下载

导读

本文基于 pysheeet 仓库的 docs/notes/asyncio/python-asyncio-server.rst 文档,系统讲解 Python asyncio 在网络编程中的完整实战路径:从基于 Streams 高层 API 的 TCP 服务器/客户端、底层 Socket 与 Transport/Protocol 控制、UDP 数据报服务,到 HTTPS 客户端与 SSL/TLS 加密服务端,再到 DNS 解析、sendfile 零拷贝传输与连接池等高阶模式。读完本文,你将能够用单线程事件循环构建并发网络服务,理解高层 Streams API 与底层 Transport/Protocol API 的取舍,并掌握生产级客户端常用的连接复用与并发控制手段。

为什么 asyncio 天生适合网络编程

网络 I/O 本质上是异步的:你发出一个请求,然后等待响应。在传统阻塞式编程中,等待期间线程被白白占用;而 asyncio 采用协作式多任务(cooperative multitasking),任务在 await 挂起时主动让出控制权,让事件循环去调度其他就绪任务。正如 docs/notes/asyncio/index.rst 所述,asyncio 通过协程在单线程内多路复用 socket 及其它资源的 I/O,非常适合 web 服务器、数据库客户端这类以等待外部资源为主要瓶颈的应用。

仓库中的 src/basic/asyncio_.py 提供了大量配套测试(如 test_gather、test_wait_for_timeout、test_semaphore),可作为本文所有代码示例的可运行验证;而 docs/notes/network/python-socket-async.rst 则从 select/epoll/kqueue 的角度解释了 asyncio 底层 I/O 多路复用的原理。两种模型的对比如下图所示:

事件循环模型与多线程模型的 TCP 服务器实现对比

左侧事件循环模型用 await loop.sock_accept / loop.sock_recv 加 loop.create_task 在单线程内并发处理连接;右侧多线程模型则为每个连接开一个线程。理解这张图,就理解了 asyncio 网络编程的核心:用协作式调度替代线程开销。

基于 Streams 的 TCP Echo 服务器

asyncio.start_server() 与 asyncio.open_connection() 组成的 Streams API 是 TCP 网络编程的推荐高层接口,它自动处理缓冲、编码与连接管理。下面是一个完整的 TCP Echo 服务器:

import asyncio

async def handle_client(reader, writer):
addr = writer.get_extra_info('peername')
print(f"Connected: {addr}")

while True:
data = await reader.read(1024)
if not data:
break
message = data.decode()
print(f"Received: {message!r} from {addr}")
writer.write(data)
await writer.drain()

print(f"Disconnected: {addr}")
writer.close()
await writer.wait_closed()

async def main():
server = await asyncio.start_server(
handle_client, 'localhost', 8888
)
addr = server.sockets[0].getsockname()
print(f"Serving on {addr}")

async with server:
await server.serve_forever()

asyncio.run(main())

要点说明:

  • writer.get_extra_info('peername') 获取对端地址((host, port) 元组),这是 Streams API 中获取连接元信息的标准方式;
  • await reader.read(1024) 返回空字节串(b'')表示对端已关闭连接,据此跳出循环;
  • await writer.drain() 在写入大量数据时等待底层缓冲区排空,避免背压问题;
  • async with server 保证服务关闭时正确清理;
  • server.sockets[0].getsockname() 可打印实际监听地址,便于确认端口绑定。

TCP 客户端:open_connection

客户端通过 asyncio.open_connection() 建立连接,返回的 reader / writer 提供异步收发方法:

import asyncio

async def tcp_client(message):
reader, writer = await asyncio.open_connection(
'localhost', 8888
)

print(f"Sending: {message!r}")
writer.write(message.encode())
await writer.drain()

data = await reader.read(1024)
print(f"Received: {data.decode()!r}")

writer.close()
await writer.wait_closed()

asyncio.run(tcp_client("Hello, Server!"))

客户端与服务器的读写模式完全对称:write + drain 发送,read 接收。值得注意的是,writer.close() 后建议 await writer.wait_closed(),确保底层连接真正关闭完成后再退出协程——这在并发场景下可避免资源泄漏。

底层 Socket 控制:事件循环的 sock_* 方法

当你需要精细控制 socket 选项(如 SO_REUSEADDR),或需要与既有 socket 代码集成时,可以直接使用事件循环的 loop.sock_* 系列方法操作非阻塞 socket:

import asyncio
import socket

async def handle_client(loop, conn):
while True:
data = await loop.sock_recv(conn, 1024)
if not data:
break
await loop.sock_sendall(conn, data)
conn.close()

async def server():
loop = asyncio.get_event_loop()

sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
sock.setblocking(False)
sock.bind(('localhost', 8888))
sock.listen(100)

print("Server listening on localhost:8888")
while True:
conn, addr = await loop.sock_accept(sock)
print(f"Connected: {addr}")
asyncio.create_task(handle_client(loop, conn))

asyncio.run(server())

关键点:

  • sock.setblocking(False) 是必须的——事件循环的 sock_* 方法要求非阻塞 socket;
  • SO_REUSEADDR 允许端口在服务重启后立即复用(避免 TIME_WAIT 状态导致的绑定失败),仓库 src/basic/socket_.py 的 test_reuseaddr 测试验证了该选项的读写行为;
  • listen(100) 设置等待队列长度;
  • 每个新连接用 asyncio.create_task 创建独立任务,这正是前文图中"事件循环模型"的典型写法。

UDP Echo 服务器:DatagramProtocol

UDP 无连接,API 与 TCP 截然不同:使用 loop.create_datagram_endpoint() 配合协议类处理数据报。每个数据报相互独立,可能乱序到达,也可能丢失:

import asyncio

class EchoUDPProtocol(asyncio.DatagramProtocol):
def connection_made(self, transport):
self.transport = transport

def datagram_received(self, data, addr):
message = data.decode()
print(f"Received {message!r} from {addr}")
self.transport.sendto(data, addr)

async def main():
loop = asyncio.get_event_loop()
transport, protocol = await loop.create_datagram_endpoint(
EchoUDPProtocol,
local_addr=('localhost', 9999)
)
print("UDP server listening on localhost:9999")

try:
await asyncio.sleep(3600) # Run for 1 hour
finally:
transport.close()

asyncio.run(main())

注意这里的协议是 asyncio.DatagramProtocol 而非 asyncio.Protocol,回调是 datagram_received(data, addr)(比 TCP 的 data_received(data) 多出来源地址),回包使用 self.transport.sendto(data, addr) 明确指定目标地址。由于 UDP 无连接语义,transport 不会关闭,示例用定时睡眠保持服务运行。

HTTPS 客户端:SSL 上下文与并发抓取

发起 HTTPS 请求需要配置 SSL 上下文。以下示例用底层 Streams 手写 HTTP/1.1 请求,并通过 asyncio.gather 并发抓取多个站点:

import asyncio
import ssl

async def fetch_https(host, path="/"):
# Create SSL context with certificate verification
ctx = ssl.create_default_context()

reader, writer = await asyncio.open_connection(
host, 443, ssl=ctx
)

# Send HTTP request
request = f"GET {path} HTTP/1.1\\r\\nHost: {host}\\r\\nConnection: close\\r\\n\\r\\n"
writer.write(request.encode())
await writer.drain()

# Read response
response = await reader.read()
writer.close()
await writer.wait_closed()

return response.decode()

async def main():
urls = [
("www.python.org", "/"),
("github.com", "/"),
]
tasks = [fetch_https(host, path) for host, path in urls]
responses = await asyncio.gather(*tasks)

for (host, _), resp in zip(urls, responses):
status = resp.split('\\r\\n')[0]
print(f"{host}: {status}")

asyncio.run(main())

要点:

  • ssl.create_default_context() 使用系统 CA 证书库做服务端证书校验,这是安全默认值;
  • ssl=ctx 参数传给 open_connection 即启用 TLS 握手(443 端口);
  • Connection: close 告知服务器响应结束后关闭连接,配合 await reader.read() 无参读取直到 EOF,可一次性拿全响应;
  • asyncio.gather(*tasks) 并发执行多个请求,这正是 asyncio 提升 I/O 密集任务吞吐的核心手段——仓库 src/basic/asyncio_.py 的 test_gather 验证了 gather 的并发语义。

HTTPS 服务器:TLS 加密静态内容服务

创建 HTTPS 服务器需要 SSL 证书。下面用 ssl.SSLContext 加载证书链,将 TLS 层直接嵌入 start_server:

import asyncio
import ssl

async def handle_request(reader, writer):
request = await reader.read(1024)

response = b"HTTP/1.1 200 OK\\r\\n"
response += b"Content-Type: text/html\\r\\n\\r\\n"
response += b"<html><body><h1>Hello HTTPS!</h1></body></html>"

writer.write(response)
await writer.drain()
writer.close()
await writer.wait_closed()

async def main():
# Create SSL context
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
ctx.load_cert_chain('cert.pem', 'key.pem')

server = await asyncio.start_server(
handle_request, 'localhost', 8443, ssl=ctx
)
print("HTTPS server on https://localhost:8443")

async with server:
await server.serve_forever()

# Generate self-signed cert:
# openssl req -x509 -newkey rsa:4096 -keyout key.pem -out cert.pem -days 365 -nodes
asyncio.run(main())

服务端 SSL 上下文必须使用 ssl.PROTOCOL_TLS_SERVER(而非客户端的 default context),并用 load_cert_chain(certfile, keyfile) 加载证书与私钥。文档注释中给出的 openssl 命令可生成本地自签名证书用于测试:

openssl req -x509 -newkey rsa:4096 -keyout key.pem -out cert.pem -days 365 -nodes

Transport/Protocol 底层 API

Streams API 之下是 Transport/Protocol 分层:Transport 负责实际 I/O 传输,Protocol 负责数据处理。这种职责分离让网络代码更灵活、可复用。用 loop.create_server() 直接注册协议类:

import asyncio

class EchoProtocol(asyncio.Protocol):
def connection_made(self, transport):
self.transport = transport
peername = transport.get_extra_info('peername')
print(f"Connection from {peername}")

def data_received(self, data):
print(f"Received: {data.decode()!r}")
self.transport.write(data)

def connection_lost(self, exc):
print("Connection closed")

async def main():
loop = asyncio.get_event_loop()
server = await loop.create_server(
EchoProtocol, 'localhost', 8888
)

async with server:
await server.serve_forever()

asyncio.run(main())

与 Streams 回调风格对比:connection_made(连接建立)、data_received(数据到达)、connection_lost(连接断开)是协议生命周期回调;transport.write 直接写回数据。前文的 UDP 示例 EchoUDPProtocol 正是同一 API 家族在数据报协议上的变体。当需要精细控制(如自定义协议解析、流量控制)时,此 API 比 Streams 更直接。

异步 DNS 解析

DNS 解析(getaddrinfo)同样是阻塞操作,asyncio 提供异步版本避免阻塞事件循环:

import asyncio
import socket

async def resolve_host(host, port=80):
loop = asyncio.get_event_loop()
infos = await loop.getaddrinfo(
host, port,
family=socket.AF_UNSPEC,
type=socket.SOCK_STREAM
)

for family, type_, proto, canonname, sockaddr in infos:
ip, port = sockaddr[:2]
family_name = "IPv4" if family == socket.AF_INET else "IPv6"
print(f"{host} -> {ip} ({family_name})")

async def main():
hosts = ["python.org", "github.com", "google.com"]
await asyncio.gather(*[resolve_host(h) for h in hosts])

asyncio.run(main())

参数说明:family=socket.AF_UNSPEC 表示同时查询 IPv4 与 IPv6;type=socket.SOCK_STREAM 限定流式(TCP)套接字。返回的 sockaddr 元组中前两项即 IP 与端口。仓库 src/basic/socket_.py 的 test_getaddrinfo 测试验证了同步版 socket.getaddrinfo 的返回结构,可作为对照。注意 loop.getaddrinfo 在 Python 3.7+ 中实际在默认执行器中运行,因此不会阻塞事件循环。

手写极简 HTTP 服务器

一个最小的 HTTP 服务器可以展示请求解析与响应的完整流程。生产环境请使用 aiohttp、FastAPI 等框架,但理解手写实现有助于掌握协议细节:

import asyncio

async def handle_http(reader, writer):
request = await reader.read(1024)
request_line = request.decode().split('\\r\\n')[0]
method, path, _ = request_line.split(' ')

print(f"{method} {path}")

# Simple routing
if path == '/':
body = b"<h1>Home</h1>"
status = "200 OK"
elif path == '/about':
body = b"<h1>About</h1>"
status = "200 OK"
else:
body = b"<h1>404 Not Found</h1>"
status = "404 Not Found"

response = f"HTTP/1.1 {status}\\r\\n"
response += f"Content-Length: {len(body)}\\r\\n"
response += "Content-Type: text/html\\r\\n\\r\\n"

writer.write(response.encode() + body)
await writer.drain()
writer.close()
await writer.wait_closed()

async def main():
server = await asyncio.start_server(
handle_http, 'localhost', 8080
)
print("HTTP server on http://localhost:8080")

async with server:
await server.serve_forever()

asyncio.run(main())

这里手工构造了 HTTP 响应的三要素:状态行(HTTP/1.1 200 OK)、头部(Content-Length 与 Content-Type)、空行分隔符(\\r\\n\\r\\n)加消息体。正确设置 Content-Length 是浏览器正常渲染的关键。对照前文 HTTPS 服务器示例可见:加密版与此处结构完全一致,区别仅在 start_server 是否传入 ssl 上下文。

sendfile:零拷贝高效文件传输

socket.sendfile(Python 3.5+,asyncio 的 loop.sendfile 为 3.7+)利用操作系统 sendfile 系统调用直接在文件与 socket 之间传输数据,避免数据拷贝进 Python 内存:

import asyncio

async def handle_request(reader, writer):
await reader.read(1024) # Read request

with open('index.html', 'rb') as f:
# Get file size
f.seek(0, 2)
size = f.tell()
f.seek(0)

# Send headers
headers = f"HTTP/1.1 200 OK\\r\\n"
headers += f"Content-Length: {size}\\r\\n"
headers += "Content-Type: text/html\\r\\n\\r\\n"
writer.write(headers.encode())

# Send file efficiently
loop = asyncio.get_event_loop()
await loop.sendfile(writer.transport, f)

writer.close()
await writer.wait_closed()

async def main():
server = await asyncio.start_server(
handle_request, 'localhost', 8080
)
async with server:
await server.serve_forever()

asyncio.run(main())

流程:先读取请求、计算文件大小(seek 到末尾再 tell)、发送头部,然后 await loop.sendfile(writer.transport, f) 让内核直接完成文件到 socket 的传输。注意 sendfile 要求文件以二进制模式打开,且传输后文件指针已移动到末尾,无需手动 reset。

连接池:复用连接提升客户端性能

连接池复用已建立的连接,避免为每个请求重新经历 TCP 三次握手与 TLS 握手,是高频请求客户端的必备优化。下面的实现用 deque 保存空闲连接、用 asyncio.Lock 保证池操作的原子性:

import asyncio
from collections import deque

class ConnectionPool:
def __init__(self, host, port, size=5):
self.host = host
self.port = port
self.size = size
self._pool = deque()
self._lock = asyncio.Lock()

async def get(self):
async with self._lock:
if self._pool:
return self._pool.popleft()

# Create new connection
reader, writer = await asyncio.open_connection(
self.host, self.port
)
return reader, writer

async def put(self, reader, writer):
async with self._lock:
if len(self._pool) < self.size:
self._pool.append((reader, writer))
else:
writer.close()
await writer.wait_closed()

async def close(self):
async with self._lock:
while self._pool:
reader, writer = self._pool.popleft()
writer.close()
await writer.wait_closed()

async def fetch(pool, message):
reader, writer = await pool.get()
try:
writer.write(message.encode())
await writer.drain()
data = await reader.read(1024)
return data.decode()
finally:
await pool.put(reader, writer)

async def main():
pool = ConnectionPool('localhost', 8888, size=3)
try:
tasks = [fetch(pool, f"msg{i}") for i in range(10)]
results = await asyncio.gather(*tasks)
for r in results:
print(r)
finally:
await pool.close()

asyncio.run(main())

设计要点:

  • get() 优先从池中取空闲连接,池空时新建连接——注意建连操作放在锁外,避免持有锁期间阻塞其他协程;
  • put() 在池未满时归还连接,池满则关闭多余连接(writer.close() + await writer.wait_closed());
  • fetch() 用 try/finally 保证无论请求成功与否连接都被归还;
  • close() 遍历池中所有连接逐一关闭,实现优雅清理。

asyncio.Lock 的使用与仓库 src/basic/asyncio_.py 中 test_lock 的语义一致:它保证同一时刻只有一个协程能修改池状态,但必须在 await 的上下文中使用,且只在同一事件循环内有效。

组合应用:一个可运行的完整链路

将上述组件串联即可得到完整验证:先启动 TCP Echo 服务器(端口 8888),再运行连接池示例,10 个并发请求经 3 条连接的池复用往返,输出 msg0~msg9 的回显。整个过程单线程完成,可配合 src/basic/asyncio_.py 中的 test_create_task、test_gather、test_semaphore 等测试理解底层调度语义。若需进一步学习,可参阅同目录的 python-asyncio-basic.rst(协程与事件循环基础)与 python-asyncio-advanced.rst(同步原语、队列、子进程与优雅关闭)。

赞

分享

  • 文档
  • 教程
  • 开发工具

【免费下载链接】pysheeet

Python Cheat Sheet

项目地址:
https://gitcode.com/gh_mirrors/py/pysheeet

点击查看 免费下载

上一篇:
英雄联盟智能助手League Akari:重新定义游戏自动化体验

下一篇:
DLSS Swapper完全指南:游戏画质优化大师深度解析

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

赞(0)
未经允许不得转载:网硕互联帮助中心 » Python Asyncio 网络编程实战指南:从 TCP/UDP 服务器到 SSL/TLS 与连接池(pysheeet)
分享到: 更多 (0)

评论 抢沙发

评论前必须登录!