树莓派4B上Python多线程异常:可编程逻辑线程计算错误
问题
使用Raspberry Pi 4B搭配BLE传感器开发套件,树莓派作为边缘设备负责:
- 读取BLE传感器数据
- 数据滤波与处理,同时用两个机器学习模型和可编程逻辑处理同一数据,获取电梯轿厢垂直位置及门状态
单独运行可编程逻辑或ML模型代码均正常,但整合后:
- 可编程逻辑线程计算结果错误
- 禁用ML线程后,可编程逻辑线程恢复正常,排除函数本身问题
htop状态截图:
代码
from bluepy import btle import threading import numpy as np import tensorflow as tf import time import variables as var import functions as func class MyDelegate(btle.DefaultDelegate): def __init__(self): btle.DefaultDelegate.__init__(self) def handleNotification(self, cHandle, data): var.acc_z[0] = int.from_bytes(data[6:8],"little",signed=True) var.acc_y[0] = int.from_bytes(data[4:6],"little",signed=True) var.acc_x[0] = int.from_bytes(data[2:4],"little",signed=True) var.mag_x[0] = int.from_bytes(data[14:16],"little",signed=True)/1000#Convert mGAuss to Gauss var.mag_y[0] = int.from_bytes(data[16:18],"little",signed=True)/1000#Convert mGAuss to Gauss var.mag_z[0] = int.from_bytes(data[18:20],"little",signed=True)/1000#Convert mGAuss to Gauss func.Filter_Data() if var.calculate_average == 1: func.Calculate_Average() else: var.ready_read = True if var.count_buffer_DOORS < var.WINDOW_STEP_DOORS - 1: var.count_buffer_DOORS += 1 var.window_buffer_DOORS = var.window_buffer_DOORS[1:]+[var.f_acc_x[0]-var.average_acc_x] else: var.count_buffer_DOORS = 0 var.window_buffer_DOORS = var.window_buffer_DOORS[1:]+[var.f_acc_x[0]-var.average_acc_x] var.ready_model_DOORS = 1 if var.count_buffer_CAB < var.WINDOW_STEP_CAB - 1: var.count_buffer_CAB += 1 var.window_buffer_CAB = var.window_buffer_CAB[1:]+[var.f_acc_y[0]-var.average_acc_y] else: var.count_buffer_CAB = 0 var.window_buffer_CAB = var.window_buffer_CAB[1:]+[var.f_acc_y[0]-var.average_acc_y] var.ready_model_CAB = 1 #Thread where programmed logic calculations run def programmedLogic_THREAD(): while var.stop_thread == False: if (var.ready_read == True): func.Calculate_Desv_Est_Mag() func.Calculate_Position_Cab() if var.flag_ready == 1 and var.learn_mag == 1: func.Learn_Mag() func.Calculate_Position_Doors() if var.update_average == 1: func.Update_Average() else: var.count_average = 0 var.sum_acc_z = 0 var.sum_acc_y = 0 var.sum_acc_x = 0 if var.state == 1: var.DB_LP_cab = "Up" elif var.state == -1: var.DB_LP_cab = "Down" elif var.state == 0: var.DB_LP_cab = "Still" if var.move_doors == 1: var.DB_LP_doors = "Opening" elif var.move_doors == -1: var.DB_LP_doors = "Closing" elif var.move_doors == 0: var.DB_LP_doors = "Still" var.ready_read = False #Thread where ML Doors Model runs and predicts def doorsML_THREAD(): while var.stop_thread == False: if (var.ready_model_DOORS == 1) and (var.ML_flag == True): var.ready_model_DOORS = 0 x = np.expand_dims(var.window_buffer_DOORS, axis = 1) var.result_DOORS = custom_model_puertas.predict(x.T.tolist(), verbose = 0) var.doors_State = np.argmax(var.result_DOORS) var.array_doors_State = var.array_doors_State[1:]+[var.doors_State] if all(element == var.array_doors_State[0] for element in var.array_doors_State): var.print_doors_State = var.array_doors_State[0] if var.print_doors_State == 0: var.DB_ML_doors = "Opening" elif var.print_doors_State == 1: var.DB_ML_doors = "Closing" elif var.print_doors_State == 2: var.DB_ML_doors = "Still" #Thread where ML Cab Model runs and predicts def cabML_THREAD(): while var.stop_thread == False: if (var.ready_model_CAB == 1) and (var.ML_flag == True): var.ready_model_CAB = 0 x = np.expand_dims(var.window_buffer_CAB, axis = 1) var.result_CAB = custom_model_cabina.predict(x.T.tolist(), verbose = 0) var.cab_State = np.argmax(var.result_CAB) var.array_cab_State = var.array_cab_State[1:]+[var.cab_State] if all(element == var.array_cab_State[0] for element in var.array_cab_State): var.print_cab_State = var.array_cab_State[0] if var.print_cab_State == 0: var.DB_ML_cab = "Down" elif var.print_cab_State == 1: var.DB_ML_cab = "Still" elif var.print_cab_State == 2: var.DB_ML_cab="Up" #Thread that prints results in terminal def print_THREAD(): count = 0 while (var.stop_thread == False): count+=1 if count>=20: count=0 var.programmedLogic_flag = not var.programmedLogic_flag var.ML_flag = not var.ML_flag if var.programmedLogic_flag == True: print(f"\n\tLP Data:\nCab state:\tDoors state:\n{var.DB_LP_cab}\t\t{var.DB_LP_doors}") elif var.ML_flag == True: print(f"\n\tML Data:\nCab state:\tDoors state:\n{var.DB_ML_cab}\t\t{var.DB_ML_doors}") time.sleep(0.5) #Create threads thread1 = threading.Thread(target=programmedLogic_THREAD) thread2 = threading.Thread(target=doorsML_THREAD) thread3 = threading.Thread(target=cabML_THREAD) thread4 = threading.Thread(target=print_THREAD) #Upload ML Models print("Uploading ML Models...") custom_model_puertas = tf.keras.models.load_model('/home/aag/Desktop/Edge_Impulse_Model/doors_MODEL.h5') custom_model_cabina = tf.keras.models.load_model('/home/aag/Desktop/Edge_Impulse_Model/cab_MODEL.h5') print("Successfully uploaded") #Connect to BLE device print("Connecting ...") p = btle.Peripheral(var.MAC_ADDR,addrType=btle.ADDR_TYPE_RANDOM) p.setDelegate( MyDelegate() ) print("Connected to DA:BB:C5:28:12:A7") # Setup to turn notifications on svc = p.getServiceByUUID("00000000-0001-11e1-9ab4-0002a5d5c51b") ch = svc.getCharacteristics("00e00000-0001-11e1-ac36-0002a5d5c51b")[0] p.writeCharacteristic(ch.valHandle+1, b"\x01\x00") #Initialize threads thread1.start() thread2.start() thread3.start() thread4.start() while True: try: if p.waitForNotifications(1.0): continue except KeyboardInterrupt: p.disconnect() print("\nDevice disconnected") var.stop_thread = True thread1.join() thread2.join() thread3.join() thread4.join() break
分析与解决建议
1. 多线程共享变量无同步机制(核心问题)
代码中所有线程直接读写variables.py中的全局变量(如var.ready_read、var.window_buffer_DOORS、var.state等),但未使用线程锁(threading.Lock)保护临界区。当ML线程和可编程逻辑线程同时读写同一变量时,会出现竞态条件,导致数据被篡改、计算逻辑混乱。
修复步骤:
- 在
variables.py中定义全局锁:import threading data_lock = threading.Lock() - 所有读写共享变量的代码块,用
with var.data_lock:包裹:
比如在handleNotification中:
再比如with var.data_lock: var.acc_z[0] = int.from_bytes(data[6:8],"little",signed=True) var.acc_y[0] = int.from_bytes(data[4:6],"little",signed=True) var.acc_x[0] = int.from_bytes(data[2:4],"little",signed=True) # ... 后续修改共享变量的代码都放在锁内programmedLogic_THREAD中:while var.stop_thread == False: with var.data_lock: if (var.ready_read == True): func.Calculate_Desv_Est_Mag() func.Calculate_Position_Cab() # ... 后续处理代码都放在锁内
2. TensorFlow模型预测的线程冲突
TensorFlow/Keras模型在多线程环境下可能存在内部资源竞争,尤其是predict调用。可以尝试:
- 为每个ML线程创建独立的模型实例,避免共享同一个模型对象
- 或者在ML线程中调用
predict时,也使用全局锁限制同一时间只有一个模型在预测
3. 树莓派CPU/内存资源过载
从htop截图看,ML模型运行会占用大量CPU资源,导致可编程逻辑线程被抢占,无法及时处理数据,进而引发计算错误。可以优化:
- 降低ML模型的推理频率,比如调整
WINDOW_STEP_DOORS/WINDOW_STEP_CAB参数,减少预测次数 - 将TensorFlow模型转换为TensorFlow Lite格式,轻量化模型,降低资源占用
- 为线程设置优先级,给可编程逻辑线程更高的优先级:
thread1 = threading.Thread(target=programmedLogic_THREAD) thread1.setPriority(threading.Thread.MAX_PRIORITY) # 仅支持部分系统,Linux下可能需要修改权限
4. 数据缓冲区的非原子操作
代码中修改缓冲区的操作(如var.window_buffer_DOORS = var.window_buffer_DOORS[1:]+[...])不是原子操作,多线程同时修改会导致缓冲区数据混乱。必须用锁包裹这类操作。
内容的提问来源于stack exchange,提问作者aag
相关产品推荐
相关产品推荐

