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

asyncua线程中读取节点BrowseName延迟2秒的解决方法咨询

问题:asyncua订阅变量后读取节点BrowseName存在2秒延迟

使用asyncua订阅部分变量,因需定期动态更新订阅列表采用线程实现,但收到数据变更通知后读取节点BrowseName时会出现约2秒延迟,将该读取操作移至datachange_notification函数中也存在同样问题。

原代码:

import asyncio

from asyncua import Client, Node
from asyncua.common.subscription import DataChangeNotif, SubHandler
import pymssql
from threading import Thread

ENDPOINT = 'opc.tcp://localhost:4840'
NAMESPACE = 'http://examples.freeopcua.github.io'


class MyHandler(SubHandler):
    n = 1
    def __init__(self, variables_len=0):
        self._queue = asyncio.Queue()
        self.variables_len = variables_len
        conn = pymssql.connect(server='(local)', user='sa',
                               password='sa', database='AdventureWorks2005')
        self.cursor = conn.cursor()

    def datachange_notification(self, node: Node, value, data: DataChangeNotif) -> None:
        if self.n > self.variables_len:
            self._queue.put_nowait([node, value, data])
            print(f'Data change notification was received and queued.')
            self.n = self.variables_len+1
        self.n += 1

    async def process(self) -> None:
        try:
            while True:
                # Get data in a queue.
                [node, value, data] = self._queue.get_nowait()

                # *** Write your processing code ***
                print(await node.read_browse_name())   #delay here
                self.cursor.execute('select TOP(1) * from tool_list')  #test code
                row = self.cursor.fetchall()
                print(row)

        except asyncio.QueueEmpty:
            pass


async def main() -> None:
    async with Client(url=ENDPOINT) as client:
        # Get a variable node.
        idx = await client.get_namespace_index(NAMESPACE)
        obj = await client.get_objects_node().get_child([f'{idx}:MyObject'])
        variables = await obj.get_variables()
        # Subscribe data change.
        handler = MyHandler(variables_len=len(variables))
        subscription = await client.create_subscription(period=0, handler=handler)
        handle = await subscription.subscribe_data_change(variables)


        # Process data change every 100ms
        async def insert_to_sql():
            while True:
                await handler.process()
                #service_level = await client.get_node('i=2267').get_value()
                await asyncio.sleep(0.1)

        def between_callback():
            loop = asyncio.new_event_loop()
            asyncio.set_event_loop(loop)
            loop.run_until_complete(some_callback())
            loop.close()

        _thread = Thread(target=between_callback)
        _thread.daemon = True
        _thread.start()

        while True:
            # Get a variable node.
            idx = await client.get_namespace_index(NAMESPACE)
            obj = await client.get_objects_node().get_child([f'{idx}:MyObject'])
            variables_new = await obj.get_variables()
            if variables != variables_new:
                handler.variables_len = len(variables_new)
                handler.n = 1
                await subscription.unsubscribe(handle)
                handle = await subscription.subscribe_data_change(variables_new)
                variables = variables_new
                print('update done')
                await asyncio.sleep(100)



if __name__ == '__main__':
    asyncio.run(main())

解决方案

延迟核心原因是线程与asyncio事件循环冲突、同步数据库操作阻塞事件循环、重复发起OPC UA读取请求,可通过以下步骤解决:

1. 移除无用线程,统一在主asyncio事件循环处理逻辑

原代码中创建的线程未执行有效逻辑(some_callback()未定义),反而干扰事件调度,直接删除线程相关代码,所有异步任务放在主循环中。

2. 提前缓存节点BrowseName,避免实时读取

在订阅或更新订阅列表时,批量读取并缓存所有变量节点的BrowseName,收到通知时直接使用缓存值,无需再发起OPC UA读取请求。

3. 替换同步数据库驱动为异步版本

pymssql是同步库,会阻塞事件循环,改用asyncpymssql异步驱动,或用线程池处理同步数据库操作。

4. 优化订阅更新逻辑,减少重复操作

缓存命名空间索引和父节点,减少重复的OPC UA服务调用,提升效率。

修改后的示例代码

import asyncio
from concurrent.futures import ThreadPoolExecutor

from asyncua import Client, Node
from asyncua.common.subscription import DataChangeNotif, SubHandler
import asyncpymssql  # 使用异步数据库驱动

ENDPOINT = 'opc.tcp://localhost:4840'
NAMESPACE = 'http://examples.freeopcua.github.io'


class MyHandler(SubHandler):
    def __init__(self):
        self._queue = asyncio.Queue()
        self.node_browse_names = {}  # 缓存节点与BrowseName的映射
        self.db_pool = ThreadPoolExecutor(max_workers=2)  # 备用:线程池处理同步操作

    def datachange_notification(self, node: Node, value, data: DataChangeNotif) -> None:
        # 直接使用缓存的BrowseName
        browse_name = self.node_browse_names.get(node, "Unknown Node")
        print(f"Data change from {browse_name}: {value}")
        self._queue.put_nowait([node, value, browse_name])

    async def process(self) -> None:
        # 异步数据库连接
        async with asyncpymssql.connect(
            server='(local)',
            user='sa',
            password='sa',
            database='AdventureWorks2005'
        ) as conn:
            async with conn.cursor() as cursor:
                while True:
                    try:
                        node, value, browse_name = await self._queue.get()
                        # 异步执行数据库查询
                        await cursor.execute('select TOP(1) * from tool_list')
                        row = await cursor.fetchall()
                        print(f"DB result for {browse_name}: {row}")
                    except asyncio.QueueEmpty:
                        await asyncio.sleep(0.01)  # 空队列时短暂休眠,避免CPU空转


async def main() -> None:
    async with Client(url=ENDPOINT) as client:
        # 缓存命名空间索引和父节点
        idx = await client.get_namespace_index(NAMESPACE)
        obj = await client.get_objects_node().get_child([f'{idx}:MyObject'])
        
        handler = MyHandler()
        subscription = await client.create_subscription(period=0, handler=handler)
        
        # 初始订阅并缓存BrowseName
        variables = await obj.get_variables()
        handle = await subscription.subscribe_data_change(variables)
        # 批量读取并缓存BrowseName
        browse_names = await asyncio.gather(*[var.read_browse_name() for var in variables])
        for var, bn in zip(variables, browse_names):
            handler.node_browse_names[var] = bn.Name
        
        # 启动处理任务
        asyncio.create_task(handler.process())
        
        # 定期检查并更新订阅列表
        while True:
            variables_new = await obj.get_variables()
            if variables != variables_new:
                # 取消旧订阅
                await subscription.unsubscribe(handle)
                # 新订阅并更新缓存
                handle = await subscription.subscribe_data_change(variables_new)
                browse_names_new = await asyncio.gather(*[var.read_browse_name() for var in variables_new])
                handler.node_browse_names.clear()
                for var, bn in zip(variables_new, browse_names_new):
                    handler.node_browse_names[var] = bn.Name
                variables = variables_new
                print('订阅列表更新完成')
            await asyncio.sleep(10)  # 可根据需求调整检查间隔


if __name__ == '__main__':
    asyncio.run(main())

关键优化点说明

  • 缓存BrowseName:在订阅时批量读取并存储节点的BrowseName,彻底消除OPC UA读取延迟。
  • 异步数据库操作:用asyncpymssql替代同步库,避免阻塞事件循环。
  • 移除无用线程:删除无效线程逻辑,所有任务在主事件循环调度,避免冲突。
  • 优化订阅更新:缓存重复使用的OPC UA节点和索引,减少服务调用次数。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 03:42:53