You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python3.5 asyncio:如何等待transport.write完成或捕获错误?

没问题,我来给你梳理一下怎么用asyncio的底层Protocol API实现这个需求!

实现思路与代码示例

要实现这个「等待所有数据写入完成或抛出异常」的方法,核心是利用asyncio Transport的drain()方法跟踪缓冲区状态,同时监听连接生命周期来处理意外断开的情况。下面是具体的实现方案:

1. 自定义TrackingTCPProtocol类

这个类会帮我们维护连接状态、跟踪写入任务的完成情况,以及处理连接异常:

import asyncio

class TrackingTCPProtocol(asyncio.Protocol):
    def __init__(self):
        self.transport = None
        self._write_future = None
        self._is_closed = False

    def connection_made(self, transport):
        # 连接建立时保存transport对象,标记连接未关闭
        self.transport = transport
        self._is_closed = False

    def connection_lost(self, exc):
        # 连接断开时更新状态,如果有未完成的写入任务,触发异常
        self._is_closed = True
        if self._write_future is not None and not self._write_future.done():
            # 如果没有传入异常,默认用连接重置错误
            if exc is None:
                exc = ConnectionResetError("Connection closed unexpectedly")
            self._write_future.set_exception(exc)

    async def wait_all_data_have_been_written_or_throw(self):
        # 先检查连接是否已经关闭
        if self._is_closed:
            raise ConnectionResetError("Cannot wait for write: connection is already closed")
        
        # 创建一个Future来跟踪写入完成状态
        if self._write_future is None or self._write_future.done():
            self._write_future = asyncio.get_event_loop().create_future()
        
        try:
            # 等待transport排空缓冲区,这一步会阻塞直到所有数据发送完成
            await self.transport.drain()
            # 缓冲区排空后标记写入成功
            if not self._write_future.done():
                self._write_future.set_result(None)
        except Exception as e:
            # 写入过程中出错,设置Future异常并重新抛出
            if not self._write_future.done():
                self._write_future.set_exception(e)
            raise
        finally:
            # 清理Future
            self._write_future = None

    def write_data(self, data):
        # 写入前先检查连接状态
        if self._is_closed:
            raise ConnectionResetError("Cannot write to closed connection")
        # 将数据写入transport缓冲区
        self.transport.write(data)

2. 使用自定义Protocol创建TCP客户端

接下来用这个Protocol实现客户端逻辑,调用我们的等待方法:

async def tcp_client():
    loop = asyncio.get_event_loop()
    
    # 连接到目标服务器,替换成你的服务器地址和端口
    transport, protocol = await loop.create_connection(
        TrackingTCPProtocol, '127.0.0.1', 8888
    )
    
    try:
        # 写入测试数据
        protocol.write_data(b"Hello from asyncio TCP client!")
        print("Data sent to buffer, waiting for write completion...")
        
        # 等待所有数据写入完成或捕获异常
        await protocol.wait_all_data_have_been_written_or_throw()
        print("All data has been successfully written!")
    except Exception as e:
        print(f"Write failed with error: {str(e)}")
    finally:
        # 关闭连接
        transport.close()

# Python 3.5需要用loop.run_until_complete启动
if __name__ == "__main__":
    loop = asyncio.get_event_loop()
    loop.run_until_complete(tcp_client())
    loop.close()

关键细节说明

  • transport.drain():这是实现等待写入完成的核心,它返回的Future会在输出缓冲区所有数据发送完成后完成;如果连接在过程中断开,会直接抛出对应异常(比如ConnectionResetError)。
  • connection_lost处理:当连接意外断开(比如服务器主动关闭),connection_lost会被触发,我们在这里会终止未完成的写入任务,确保等待方法能及时捕获错误。
  • 流量控制:如果服务器接收速度慢导致缓冲区满,drain()会自动等待直到缓冲区有空闲空间,无需额外处理流量控制逻辑。

内容的提问来源于stack exchange,提问作者allan.simon

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 03:35:41