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

树莓派4B上Python多线程异常:可编程逻辑线程计算错误

问题

使用Raspberry Pi 4B搭配BLE传感器开发套件,树莓派作为边缘设备负责:

  • 读取BLE传感器数据
  • 数据滤波与处理,同时用两个机器学习模型和可编程逻辑处理同一数据,获取电梯轿厢垂直位置及门状态

单独运行可编程逻辑或ML模型代码均正常,但整合后:

  • 可编程逻辑线程计算结果错误
  • 禁用ML线程后,可编程逻辑线程恢复正常,排除函数本身问题

htop状态截图:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 18:35:12