如何检测asyncio TCP连接已断开?Python代码问题排查
问题解答
问题1:设备单连接限制下,asyncio.open_connection()仍返回成功
这是TCP协议的特性:只要硬件的TCP服务完成三次握手,open_connection()就会返回连接成功——硬件的单连接限制是应用层规则,不属于TCP层的管控范围。
你目前等待硬件状态消息的方案是合理的,因为只有硬件能告知当前连接是否有效。可以优化该逻辑:
- 连接成功后设置超时时间,等待硬件的状态确认消息,超时则主动断开连接,避免无效连接占用资源。
- 收到硬件的"连接被拒绝"类消息时,立即关闭连接并触发重连流程。
问题2:物理断开后无法检测连接状态
你的代码依赖writer.is_closing()判断连接状态,但这个方法仅在**主动调用writer.close()**时才会返回True,被动断开(比如拔USB、网络中断)不会触发状态变化。同时,write()和drain()在数据写入发送缓冲区后就会返回,不会立即检测底层连接是否失效——只有当TCP栈尝试发送数据失败(比如收到RST包、超时)时才会抛出异常,这个过程可能存在延迟。
解决方案:添加主动读取任务+维护自定义连接状态
TCP连接的断开可以通过读取端检测:当连接断开时,reader.read()会返回空字节串,或者抛出ConnectionResetError等异常。我们需要在后台运行一个读取任务,一旦检测到异常或空数据,就标记连接为断开并关闭writer。
修改后的代码示例:
import asyncio import logging from typing import TypeVar logger = logging.getLogger(__name__) host = '169.254.13.95' port = 51717 timeout_sec = 10 lock = asyncio.Lock() Self = TypeVar("Self", bound="TcpConnection") class TcpConnection: """TCP连接类,支持连接状态检测和心跳""" def __init__(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: self.reader: asyncio.StreamReader = reader self.writer: asyncio.StreamWriter = writer self._is_connected: bool = True # 自定义连接状态标志 @classmethod async def connect(cls, host: str, port: int) -> Self | None: connection = None logger.info(f'Connecting to {host}:{port}') try: reader, writer = await asyncio.open_connection(host=host, port=port) logger.info('TCP连接建立成功') connection = TcpConnection(reader, writer) # 启动后台读取任务,检测连接状态 asyncio.create_task(connection._monitor_connection()) except ConnectionRefusedError: logger.info(f'连接被拒绝 ({host}:{port})') except OSError as e: logger.info(f'连接失败 ({host}:{port}): {e}') except Exception as e: logger.warning(f'未知异常:\n{e}') finally: return connection def is_connected(self) -> bool: return self._is_connected and not self.writer.is_closing() async def _monitor_connection(self) -> None: """后台监听连接状态,断开时更新标志""" try: # 持续读取数据,连接断开时会返回空或抛出异常 while self._is_connected: # 可根据硬件协议调整读取逻辑(比如读取固定长度/直到分隔符) data = await self.reader.read(1024) if not data: logger.info('连接已断开(读取到空数据)') break # 若硬件有状态消息,可在此处处理并更新连接状态 logger.debug(f'收到硬件数据: {data.hex()}') except (ConnectionResetError, OSError) as e: logger.info(f'连接异常断开: {e}') finally: self._is_connected = False # 主动关闭writer if not self.writer.is_closing(): self.writer.close() await self.writer.wait_closed() logger.info('连接已关闭') async def keep_alive(self) -> None: logger.info('启动心跳任务') keep_alive_msg = b'\x00' while self.is_connected(): try: async with lock: self.writer.write(keep_alive_msg) await self.writer.drain() logger.debug('发送心跳消息') except (ConnectionResetError, OSError) as e: logger.info(f'心跳发送失败,连接断开: {e}') self._is_connected = False break await asyncio.sleep(4.5) logger.info('终止心跳任务') async def main() -> None: while 1: tcp = await TcpConnection.connect(host, port) if tcp and tcp.is_connected(): try: keep_alive_task = asyncio.create_task(tcp.keep_alive()) await keep_alive_task except Exception as e: logger.info(f'连接任务异常: {e}') logger.info(f'{timeout_sec}秒后尝试重连') await asyncio.sleep(timeout_sec) if __name__ == '__main__': logging.basicConfig(format='%(asctime)s,%(msecs)d %(name)s %(levelname)s %(message)s', datefmt='%Y-%m-%d %H:%M:%S', level=logging.DEBUG) try: logger.info('启动应用') asyncio.run(main()) except KeyboardInterrupt: logger.info('退出应用')
修改要点说明
- 自定义连接状态标志:用
_is_connected变量跟踪连接状态,替代依赖writer.is_closing()的不可靠判断。 - 后台连接监听任务:
_monitor_connection持续读取数据,一旦读取到空或异常,立即标记连接为断开并关闭writer。 - 心跳任务异常捕获:在心跳发送流程中捕获连接异常,及时更新状态并终止任务。
- 硬件消息处理:如果硬件有连接状态通知,可以在
_monitor_connection中处理这些消息,比如收到"已被其他用户连接"的消息时主动断开。
内容的提问来源于stack exchange,提问作者Steven Krick
相关产品推荐
相关产品推荐

