arm64平台multiprocessing.Queue多线程进程偶发阻塞问题问询
问题分析与解决方案
运行环境
- Jetson设备(arm64架构)
- Ubuntu 22.04系统
- Python 3.10.12版本
问题复现代码
import time from threading import Thread from multiprocessing import Process, Queue class ProcessClasses: def __init__(self, get_config_queue=None,): self.get_config_queue = get_config_queue self.get_config_queue_thread = Thread(target=self.get_config) self.log_thread = Thread(target=self.log) self.get_config_queue_thread.start() # start this thread will stuck #self.log_thread.start() def get_config(self): count = 0 while True: self.get_config_queue.put(count) print('get_config', count) count += 1 def log(self): #print('log') pass def debug(): queue = Queue(maxsize=1) process = Process(target=ProcessClasses, args=(queue,)) process.start() while True: res = queue.get() print('res', res) print('qsize, full, empty', queue.qsize(), queue.full(), queue.empty()) time.sleep(.5) debug_thread = Thread(target=debug) debug_thread.start() #debug_thread.join()
现象描述
当前代码运行正常,但启动log_thread后,multiprocessing.Queue有时会出现阻塞情况。阻塞时,queue.full()与queue.empty()会同时返回True,这种状态明显不符合预期。临时在循环前添加event.wait()或time.sleep()语句可缓解问题,但需要找到根本原因和彻底解决方案。
根本原因
- 多线程与Queue初始化的竞态:
multiprocessing.Queue底层依赖操作系统管道和内部锁机制实现。当在子进程的构造函数(__init__)中同时启动多个线程时,子进程刚完成fork,Queue的内部feeder线程(负责将数据从用户空间传到管道)还未完成初始化,此时工作线程的put操作会和Queue的内部线程争抢锁资源,导致Queue的状态计数变量被错误更新,出现同时标记为“满”和“空”的矛盾状态。 - arm64架构的内存可见性特性:Jetson设备的arm64架构内存模型与x86不同,内存屏障的语义更弱,锁状态的同步可能存在延迟,加剧了竞态条件引发的状态不一致问题。
彻底解决方案
方案1:延迟线程启动,避免在构造函数中启动线程
修改ProcessClasses类,新增start方法统一启动线程,确保Queue完成内部初始化后再启动工作线程:
import time from threading import Thread from multiprocessing import Process, Queue class ProcessClasses: def __init__(self, get_config_queue=None): self.get_config_queue = get_config_queue self.get_config_queue_thread = Thread(target=self.get_config) self.log_thread = Thread(target=self.log) def start(self): # 实例化完成后再启动所有线程 self.get_config_queue_thread.start() self.log_thread.start() def get_config(self): count = 0 while True: self.get_config_queue.put(count) print('get_config', count) count += 1 def log(self): # print('log') pass def debug(): queue = Queue(maxsize=1) # 用lambda包装,确保先实例化再启动线程 process = Process(target=lambda: ProcessClasses(queue).start()) process.start() while True: res = queue.get() print('res', res) print('qsize, full, empty', queue.qsize(), queue.full(), queue.empty()) time.sleep(.5) debug_thread = Thread(target=debug) debug_thread.start() # debug_thread.join()
方案2:使用单独的初始化函数替代构造函数启动线程
如果不想修改类结构,可以在子进程中先实例化ProcessClasses,再手动启动线程:
# 在debug函数中修改process的target def debug(): queue = Queue(maxsize=1) def init_and_start(): proc_cls = ProcessClasses(queue) proc_cls.log_thread.start() proc_cls.get_config_queue_thread.start() process = Process(target=init_and_start) process.start() # 后续代码不变
核心思路
避免在fork后的子进程构造阶段启动多线程,给Queue足够时间完成内部锁和feeder线程的初始化,确保所有线程操作Queue时,其内部状态已经稳定,从根源上消除竞态条件引发的状态不一致问题。
内容的提问来源于stack exchange,提问作者Tony
相关产品推荐
相关产品推荐

