如何控制两个依赖线程的执行时机?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
相关产品推荐
相关产品推荐

