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

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()

关键修改说明

  1. 移除encoder_T的线程继承与run循环:pigpio回调已在独立线程运行,无需额外线程,减少资源消耗。
  2. 优化UdpSender的队列读取逻辑:用get()阻塞等待数据,替代先判断空再读取的逻辑,避免竞争条件。
  3. pigpio实例移至run方法初始化:确保子进程独立建立GPIO连接,避免跨进程资源冲突。
  4. 添加队列满异常处理:防止高频率脉冲下队列溢出导致回调线程终止。
  5. 设置队列最大长度:避免无限制堆积数据导致内存占用过高。

内容的提问来源于stack exchange,提问作者USMAN SIDDIQUI

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 11:52:44