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

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()语句可缓解问题,但需要找到根本原因和彻底解决方案。

根本原因

  1. 多线程与Queue初始化的竞态:multiprocessing.Queue底层依赖操作系统管道和内部锁机制实现。当在子进程的构造函数(__init__)中同时启动多个线程时,子进程刚完成fork,Queue的内部feeder线程(负责将数据从用户空间传到管道)还未完成初始化,此时工作线程的put操作会和Queue的内部线程争抢锁资源,导致Queue的状态计数变量被错误更新,出现同时标记为“满”和“空”的矛盾状态。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 01:40:17