Python queue.Queue同步异常:pigpio编码器脉冲队列无法被UdpSender清空
解决pigpio编码器脉冲读取与UDP发送的队列阻塞问题
使用pigpio库读取编码器脉冲时,单编码器场景下脉冲读取正常,但将编码器线程(通过pigpio.callback()写入队列)和UDP发送线程(从队列读取数据)封装到multiprocessing.Process后,出现队列无法被清空的问题。即使使用threading.Lock保证队列访问互斥,高频率脉冲下队列仍只会被填充一次后停止更新,扩展到双编码器场景问题依旧。
核心问题与解决方案
修正锁与队列操作逻辑
queue.Queue.empty()并非线程安全的可靠判断,即使加锁,判断空和执行get_nowait()之间仍可能出现状态变化。建议直接使用get()阻塞等待队列数据,或捕获Empty异常处理空队列情况,避免无效循环判断。优化pigpio回调线程模型
pigpio.callback()会在pigpio库内部的独立线程中执行回调函数,自定义的encoder_T线程的run循环无实际作用,反而浪费资源,需移除该循环,让回调线程独立工作。调整多进程中pigpio实例初始化时机
在multiprocessing.Process的__init__方法中初始化pigpio.pi()可能因父进程与子进程的资源继承问题导致异常,应将pigpio的初始化移到run方法中,确保子进程独立建立GPIO连接。处理高频率脉冲下的队列异常
高频率脉冲场景下,put_nowait()可能因队列满抛出Full异常,导致数据丢失。需捕获该异常,或设置队列的最大长度并根据业务场景处理溢出,避免回调线程因异常终止。
修改后的代码示例
import pigpio import threading import queue import socket import config from multiprocessing import Process class encoder_T: def __init__(self, pinNumber, scale_value, gpio, queue, lock): self.gpio = gpio self.pinNumber = pinNumber self.cb1 = self.gpio.callback(self.pinNumber, pigpio.RISING_EDGE, self.callback_func) self.scale_value = scale_value self.position_pulses = 0 self.queue = queue self.position_degrees = 0 self.direction = 0 self.lock = lock def callback_func(self, pinNumber, level, tick): # 读取方向引脚状态 self.direction = 1 if self.gpio.read(config.DIRECTION_TRAINING) == 0 else -1 self.position_pulses += self.direction # 计算角度(处理循环校准) mod_value = config.TR_CALIBRATION_VALUE self.position_degrees = ((self.position_pulses % mod_value + mod_value) % mod_value * self.scale_value) / mod_value # 写入队列,处理满队列异常 try: with self.lock: self.queue.put_nowait(format(self.position_degrees,".2f")) except queue.Full: # 可根据需求添加日志或丢弃策略 pass class UdpSender(threading.Thread): def __init__(self, queue1, lock): super().__init__() self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.queue1 = queue1 self.num1 = 0 self.lock = lock def run(self): while True: try: # 直接阻塞获取队列数据,无需提前判断空 with self.lock: self.num1 = self.queue1.get() # 发送UDP数据 message = "{}, {}".format(self.num1, self.num1) print(message) self.sock.sendto(message.encode('utf-8'), (config.MULTICAST_ADDRESS, config.MULTICAST_PORT)) except queue.Empty: # 空队列时短暂休眠,避免CPU占用过高 threading.Event().wait(0.0001) class Encoder_Process(Process): def __init__(self): super().__init__() self.lock = threading.Lock() # 设置队列最大长度,避免内存溢出 self.queue1 = queue.Queue(maxsize=100) def run(self): # 在子进程内初始化pigpio连接 gpio = pigpio.pi() if not gpio.connected: print("GPIO连接失败") return # 初始化编码器和UDP发送线程 self.num1 = encoder_T(8, 360, gpio, self.queue1, self.lock) self.udp1 = UdpSender(self.queue1, self.lock) self.udp1.start() # 编码器无需启动线程,由pigpio回调自动处理 self.udp1.join() if __name__ == '__main__': p2 = Encoder_Process() p2.start() p2.join()
关键修改说明
- 移除encoder_T的线程继承与run循环:pigpio回调已在独立线程运行,无需额外线程,减少资源消耗。
- 优化UdpSender的队列读取逻辑:用
get()阻塞等待数据,替代先判断空再读取的逻辑,避免竞争条件。 - pigpio实例移至run方法初始化:确保子进程独立建立GPIO连接,避免跨进程资源冲突。
- 添加队列满异常处理:防止高频率脉冲下队列溢出导致回调线程终止。
- 设置队列最大长度:避免无限制堆积数据导致内存占用过高。
内容的提问来源于stack exchange,提问作者USMAN SIDDIQUI
相关产品推荐
相关产品推荐

