asyncio.create_datagram_endpoint 不能直接 await 是因为它返回协程对象,必须用 await 驱动执行,否则静默卡住;它底层调用 socket 并注册到事件循环,DNS 解析等耗时操作在此 await 阶段完成。

asyncio.create_datagram_endpoint 为什么不能直接 await?
因为 create_datagram_endpoint 返回的是一个协程对象,不是可等待的“完成态”,必须用 await 驱动它执行,否则会卡在未调度状态——常见错误是漏掉 await,结果程序静默卡住,连异常都不抛。
它底层调用系统 socket 接口并注册到事件循环,耗时操作(如地址解析、绑定端口)都在 await 期间完成。如果传入的 remote_addr 是域名,DNS 解析也在此阶段同步阻塞(除非你提前用 loop.getaddrinfo 异步解析好)。
- 必须搭配
async with或显式await loop.create_datagram_endpoint(...)使用 - 不支持
ssl=True:UDP 本身无连接、无握手,create_datagram_endpoint不提供 TLS 封装能力 - 若绑定
0.0.0.0:0,系统分配临时端口,需从返回的transport.get_extra_info('sockname')中读取实际端口
如何正确实现 UDP 客户端发送 + 接收逻辑?
UDP 是无连接的,但 asyncio 的 create_datagram_endpoint 仍需指定协议类(继承 asyncio.DatagramProtocol),接收逻辑全靠重写 datagram_received 方法——它不是回调函数,而是由 transport 在数据到达时自动调用,且**不在用户协程栈中执行**。
这意味着:你在 datagram_received 里不能直接 await 协程(会报 "not inside async function"),必须用 loop.create_task() 或 loop.call_soon_threadsafe() 转交。
立即学习“Python免费学习笔记(深入)”;
class EchoClientProtocol(asyncio.DatagramProtocol):
def __init__(self, message, on_con_lost):
self.message = message
self.on_con_lost = on_con_lost
def connection_made(self, transport):
self.transport = transport
self.transport.sendto(self.message.encode())
def datagram_received(self, data, addr):
print(f"Received: {data.decode()}")
self.transport.close()
def error_received(self, exc):
print('Error received:', exc)
def connection_lost(self, exc):
self.on_con_lost.set_result(True)
- 发送必须在
connection_made中触发,不能在__init__里 —— 此时 transport 还没就绪 -
datagram_received和error_received可能被并发调用,注意共享状态的线程安全(比如用asyncio.Lock或原子操作) - 客户端通常不长期监听,收到响应后调用
transport.close()即可;服务端则应保持 transport 开放
UDP 服务端如何避免丢包和资源泄漏?
asyncio 的 UDP transport 默认使用系统 UDP 缓冲区,一旦接收速度超过处理速度,内核缓冲区满就会静默丢包——这跟 TCP 的流控完全不同,没有任何重传或通知机制。
根本解法不是加大 recv buffer(SO_RCVBUF),而是确保 datagram_received 尽可能轻量,并把实际业务逻辑异步派发出去。同时,transport 没有内置超时,必须自己用 loop.call_later 管理空闲连接(虽然 UDP 本无连接概念,但你可以按客户端 IP:port 维护会话状态)。
- 不要在
datagram_received中做文件 I/O、HTTP 请求、数据库查询等阻塞操作 - 用
transport.get_extra_info('peername')获取来源地址,可用于限速或黑白名单 - 服务端不调用
transport.close(),除非明确要下线;否则应让 event loop 管理生命周期 - 若需高吞吐,考虑用
socket.SO_REUSEPORT(Linux/macOS)启动多个 worker 进程分担负载
为什么 recvfrom 不会阻塞,但 asyncio 里还是可能卡住?
因为 asyncio 的 transport 底层仍依赖 select/epoll/kqueue 等 I/O 多路复用机制,而 UDP socket 的就绪判断只看内核接收缓冲区是否有数据。如果缓冲区为空,事件循环会跳过该 socket;但如果缓冲区一直有数据(比如攻击者持续发包),而你的 datagram_received 处理太慢,就会积压、最终溢出丢包——此时现象是“看起来没收到新包”,实则是旧包还没来得及处理。
- 可通过
ss -u -i(Linux)查看 UDP socket 的rcvbuf使用量和丢包数(skmem_r,drop字段) - Python 层无法直接获取当前内核缓冲区长度,但可用
transport.get_extra_info('socket').getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF)查配置值 - 真正关键的不是缓冲区大小,而是
datagram_received的平均处理耗时是否稳定低于数据到达间隔
create_datagram_endpoint 提供的是单向传输通道,收发逻辑分离且执行上下文不同——这点最容易被当成普通同步 socket 去用,然后陷入难以复现的丢包或死锁。


















