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

如何检测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('退出应用')

修改要点说明

  1. 自定义连接状态标志:用_is_connected变量跟踪连接状态,替代依赖writer.is_closing()的不可靠判断。
  2. 后台连接监听任务:_monitor_connection持续读取数据,一旦读取到空或异常,立即标记连接为断开并关闭writer。
  3. 心跳任务异常捕获:在心跳发送流程中捕获连接异常,及时更新状态并终止任务。
  4. 硬件消息处理:如果硬件有连接状态通知,可以在_monitor_connection中处理这些消息,比如收到"已被其他用户连接"的消息时主动断开。

内容的提问来源于stack exchange,提问作者Steven Krick

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 14:55:02