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
修复方案及代码示例
核心修改点
- 捕获未处理的
TimeoutError,避免程序崩溃 - 用重连标志替代强制重启,让主循环统一处理重连逻辑
- 优化订阅生命周期,避免资源泄漏
- 明确状态码判断,提高代码可读性
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())
方案说明
- 捕获TimeoutError:原代码未捕获
concurrent.futures.TimeoutError,导致watchdog超时后程序直接崩溃,新增该异常类型后,错误会被纳入重连逻辑。 - 重连标志替代强制重启:
SubHandler中设置need_reconnect标志,主循环定期检测该标志后主动断开连接,由async with自动清理客户端资源,避免os.execl在异步环境下的资源泄漏问题。 - 优化订阅管理:使用官方推荐的
subscribe_data_change替代内部方法_subscribe,并在finally块中主动删除订阅,确保服务器端无残留无效订阅。 - 明确状态码判断:结合官方枚举值和硬编码状态码,覆盖不同场景下的连接异常。
额外建议
- 调整
client.check_connection()的间隔,或在客户端初始化时设置更长的watchdog超时时间,减少误判。 - 给
write_file添加异常处理,避免文件写入失败导致程序崩溃。 - 增加重连次数统计,达到阈值时发送告警(如邮件、短信),避免无限重连。
内容的提问来源于stack exchange,提问作者heMe
相关产品推荐
相关产品推荐

