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

opcua-asyncio看门狗循环报错解决方案求助

OPC UA客户端连接中断后程序崩溃的修复方案

问题背景

用Python + opcua-asyncio在Raspberry Pi 4上搭建OPC UA客户端,订阅同网络内KEPServerEX6服务器的数据(服务器每2秒发布新值,客户端接收后存入JSON)。但每隔几天连接会中断,报错后程序直接停止,必须手动重启。曾尝试在SubHandler中检测错误后重启程序,既不优雅还经常失效。

报错信息

ERROR:asyncua.client.client:Error in watchdog loop
Traceback (most recent call last):
  File "/home/pi/.local/lib/python3.7/site-packages/asyncua/client/client.py", line 457, in _monitor_server_loop
    _ = await self.nodes.server_state.read_value()
  File "/home/pi/.local/lib/python3.7/site-packages/asyncua/common/node.py", line 179, in read_value
    result = await self.read_data_value()
  File "/home/pi/.local/lib/python3.7/site-packages/asyncua/common/node.py", line 190, in read_data_value
    return await self.read_attribute(ua.AttributeIds.Value, None, raise_on_bad_status)
  File "/home/pi/.local/lib/python3.7/site-packages/asyncua/common/node.py", line 304, in read_attribute
    result = await self.session.read(params)
  File "/home/pi/.local/lib/python3.7/site-packages/asyncua/client/ua_client.py", line 397, in read
    data = await self.protocol.send_request(request)
  File "/home/pi/.local/lib/python3.7/site-packages/asyncua/client/ua_client.py", line 160, in send_request
    data = await asyncio.wait_for(self._send_request(request, timeout, message_type), timeout if timeout else None)
  File "/usr/lib/python3.7/asyncio/tasks.py", line 423, in wait_for
    raise futures.TimeoutError()
concurrent.futures._base.TimeoutError
INFO:asyncua.common.subscription:Publish callback called with result: PublishResult(SubscriptionId=9, AvailableSequenceNumbers=[], MoreNotifications=True, NotificationMessage_=NotificationMessage(SequenceNumber=0, PublishTime=datetime.datetime(2023, 5, 30, 17, 39, 59, 457322), NotificationData=[StatusChangeNotification(Status=2148270080, DiagnosticInfo_=DiagnosticInfo(SymbolicId=None, NamespaceURI=None, Locale=None, LocalizedText=None, AdditionalInfo=None, InnerStatusCode=None, InnerDiagnosticInfo=None))]), Results=[], DiagnosticInfos=[])
WARNING:__main__:status_change_notification 2148270080
WARNING:__main__:Some bad status change, restarting program
ERROR:asyncua.common.subscription:Exception calling status change handler

修复方案及代码示例

核心修改点

  1. 捕获未处理的TimeoutError,避免程序崩溃
  2. 用重连标志替代强制重启,让主循环统一处理重连逻辑
  3. 优化订阅生命周期,避免资源泄漏
  4. 明确状态码判断,提高代码可读性
import sys
import os
import asyncio
import concurrent.futures
from asyncua import Client, ua
from asyncua.common.node import Node
import logging

_logger = logging.getLogger(__name__)

# 假设write_file是已实现的文件写入函数
def write_file(filename, value):
    with open(f"{filename}.json", "w") as f:
        f.write(str(value))

class SubHandler:
    def __init__(self, crude, acid):
        self.crude = crude
        self.acid = acid
        self.need_reconnect = False  # 新增重连标志

    def datachange_notification(self, node: Node, val, data):
        if node == self.crude:
            write_file("crude_flow", val)
        elif node == self.acid:
            write_file("acid_flow", val)
            
    def event_notification(self, event: ua.EventNotificationList):
        _logger.warning("event_notification %s", event)
        pass
    
    def status_change_notification(self, status: ua.StatusChangeNotification):
        _logger.warning("status_change_notification %s", status.Status)
        # 匹配连接/会话异常关闭的状态码
        bad_status_codes = (
            ua.StatusCodes.BadSessionClosed,
            ua.StatusCodes.BadConnectionClosed,
            2148270080  # 兼容硬编码状态码
        )
        if status.Status in bad_status_codes:
            _logger.warning("Connection lost, triggering reconnect")
            self.need_reconnect = True

async def main():
    url_scada = "opc.tcp://your-server-ip:port"  # 替换为实际服务器地址
    nodeID_crude = "ns=2;s=CrudeFlow"  # 替换为实际节点ID
    nodeID_acid = "ns=2;s=AcidFlow"

    while True:
        client = Client(url=url_scada)
        handler = None
        sub = None
        try:
            async with client:
                _logger.warning("Connected to OPC UA server")

                crude = client.get_node(nodeID_crude)
                acid = client.get_node(nodeID_acid)
                nodeIDs = [crude, acid]

                handler = SubHandler(crude, acid)
                sub = await client.create_subscription(500, handler)

                # 创建数据变更过滤器
                filter = ua.DataChangeFilter(Trigger=ua.DataChangeTrigger.StatusValueTimestamp)

                # 订阅节点数据变更(使用官方推荐的subscribe_data_change替代内部方法_subscribe)
                for node in nodeIDs:
                    await sub.subscribe_data_change(
                        node,
                        attr=ua.AttributeIds.Value,
                        mfilter=filter
                    )

                while True:
                    await asyncio.sleep(0.5)
                    await client.check_connection()
                    # 检测重连标志,触发重连
                    if handler and handler.need_reconnect:
                        _logger.warning("Initiating reconnect from main loop")
                        break
        except (ConnectionError, ua.Ua.Error, concurrent.futures.TimeoutError):
            _logger.warning("Connection error occurred, reconnecting in 1 second")
            await asyncio.sleep(1)
        finally:
            # 确保订阅被正确清理
            if sub:
                try:
                    await sub.delete()
                    _logger.warning("Subscription deleted successfully")
                except Exception as e:
                    _logger.warning("Failed to delete subscription: %s", str(e))

if __name__ == "__main__":
    logging.basicConfig(level=logging.WARNING)
    asyncio.run(main())

方案说明

  1. 捕获TimeoutError:原代码未捕获concurrent.futures.TimeoutError,导致watchdog超时后程序直接崩溃,新增该异常类型后,错误会被纳入重连逻辑。
  2. 重连标志替代强制重启:SubHandler中设置need_reconnect标志,主循环定期检测该标志后主动断开连接,由async with自动清理客户端资源,避免os.execl在异步环境下的资源泄漏问题。
  3. 优化订阅管理:使用官方推荐的subscribe_data_change替代内部方法_subscribe,并在finally块中主动删除订阅,确保服务器端无残留无效订阅。
  4. 明确状态码判断:结合官方枚举值和硬编码状态码,覆盖不同场景下的连接异常。

额外建议

  • 调整client.check_connection()的间隔,或在客户端初始化时设置更长的watchdog超时时间,减少误判。
  • 给write_file添加异常处理,避免文件写入失败导致程序崩溃。
  • 增加重连次数统计,达到阈值时发送告警(如邮件、短信),避免无限重连。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 02:37:05