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

如何控制两个依赖线程的执行时机?OPC UA客户端读写优化咨询

OPC UA 多线程读写日志问题分析与解决方案

需求背景

我有一个生成随机数值的dummy OPC UA服务器,需要用客户端读取数据,每10秒将读取的数据存入日志文件。因为服务器数据生成速度快,若读完10秒数据后立即执行存储操作,会导致客户端延迟读取新数据,所以计划用多线程方案:一个线程负责持续读取数据,每10秒攒一批数据后触发存储线程,且存储操作不能中断读取线程的运行。

现有代码的核心问题

  • 线程启动方式错误:使用t1.run()和t2.run()不会创建新线程,而是在当前主线程中同步执行函数,完全没有实现多线程并行的效果,正确做法是调用start()方法启动新线程。
  • 线程返回值获取错误:异步线程的run()方法无法直接返回值,logs,dt = t1.run()的写法不成立,线程间数据传递需要用线程安全的机制(比如队列)。
  • 读取逻辑阻塞:外层while True循环会等待读取线程执行完10秒的任务才继续,导致读取过程无法与存储操作并行,违背多线程设计初衷。
  • 冗余延迟:store_values函数中的time.sleep(10)完全多余,会白白占用存储线程的执行时间。

修正后的实现代码

from opcua import Client
import time
from datetime import datetime
import threading
import queue


class OPCUAClient:
    def __init__(self, url="opc.tcp://127.0.0.1:8080", store_cycle=10, read_cycle=1):
        self.url = url
        self.store_cycle = store_cycle
        self.read_cycle = read_cycle
        self.log_queue = queue.Queue()  # 线程安全队列,用于传递待存储日志
        self.stop_event = threading.Event()  # 用于优雅终止线程

    def connect_to_server(self):
        clt = Client(self.url)
        clt.connect()
        var1 = clt.get_node("ns=2;s=var1_val")
        var2 = clt.get_node("ns=2;s=var2_val")
        return var1, var2

    def read_loop(self, var1, var2):
        while not self.stop_event.is_set():
            iterator = 0
            logs = []
            dt = datetime.now().strftime("%Y-%m-%d_%Hh%Mm%Ss")
            # 攒够一个存储周期的数据
            while iterator < (self.store_cycle / self.read_cycle) and not self.stop_event.is_set():
                iterator += 1
                val1 = var1.get_value()
                val2 = var2.get_value()
                logs.append(f"datachange_notification ns=2;s=var1_val {val1}")
                logs.append(f"datachange_notification ns=2;s=var2_val {val2}")
                print(val1, val2)
                time.sleep(self.read_cycle)
            # 将日志提交给存储线程
            if logs:
                self.log_queue.put((logs, dt))
                print(f"已提交一批日志至队列,待存储")

    def store_loop(self):
        while not self.stop_event.is_set():
            try:
                # 阻塞等待队列数据,定期检查终止信号
                logs, dt = self.log_queue.get(timeout=1)
                with open(f"logs/message_{dt}.log", 'w+') as logfile:
                    logfile.write('\n'.join(logs))
                print(f"日志 {dt} 已存储完成")
                self.log_queue.task_done()
            except queue.Empty:
                continue


if __name__ == "__main__":
    client = OPCUAClient()
    var1, var2 = client.connect_to_server()

    # 启动守护线程,主线程退出时自动终止
    read_thread = threading.Thread(target=client.read_loop, args=(var1, var2), daemon=True)
    store_thread = threading.Thread(target=client.store_loop, daemon=True)
    read_thread.start()
    store_thread.start()

    # 等待用户中断程序
    try:
        while True:
            time.sleep(1)
    except KeyboardInterrupt:
        print("正在停止程序...")
        client.stop_event.set()
        read_thread.join()
        store_thread.join()
        print("程序已停止")

关键修正说明

  • 线程安全数据传递:用queue.Queue实现读取线程与存储线程的通信,避免多线程直接共享变量引发的线程安全问题。
  • 并行执行逻辑:读取线程持续运行,每完成一个存储周期的数据收集就提交到队列,立即开始下一轮读取,完全不等待存储操作完成,彻底消除读取延迟。
  • 优雅退出机制:通过threading.Event捕获用户中断信号,安全终止两个线程,避免程序异常退出。
  • 移除冗余操作:删除存储线程中不必要的sleep,让存储操作高效执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 23:27:02