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

Python中含Lock的While循环无法终止问题求助

排查并解决Python多线程Lock循环无法终止的问题

嘿,我来帮你搞定这个循环卡死的问题!你的代码里的循环无法终止,核心原因出在非线程安全的队列空检查和阻塞式的get()调用上,咱们一步步拆解:

问题根源

  1. empty()方法的线程不安全特性:你在锁外面用while not my_queue.empty()判断队列是否有元素,这在多线程场景下完全不可靠。比如队列只剩最后一个元素时,可能有多个线程同时通过这个判断(因为此时队列还没被取空),然后它们会排队等待获取锁。
  2. 阻塞式get()导致线程挂起:当第一个线程取走最后一个元素后,队列已经空了,但后续拿到锁的线程会执行my_queue.get()——而queue.Queue.get()默认是阻塞的,会一直等待新元素进入,这些线程就会一直卡在这里,循环自然无法终止。

解决方案

这里给你两种实用的修复方案,根据你的场景选择:

方案1:捕获queue.Empty异常(推荐用于一次性任务)

把队列获取改成非阻塞模式,队列为空时抛出Empty异常,我们捕获这个异常来终止循环,同时用finally确保锁一定会被释放:

修改后的完整代码:

import threading
import queue
import time

my_queue = queue.Queue()
lock = threading.Lock()

for i in range(5):
    my_queue.put(i)

def something_useful(CPU_number):
    while True:
        item = None
        try:
            lock.acquire()
            # 非阻塞获取元素,队列为空时直接抛出Empty异常
            item = my_queue.get(block=False)
            print(f"\n CPU_C {CPU_number}: {item}")
        except queue.Empty:
            # 队列为空,退出循环
            break
        finally:
            # 无论是否成功取到元素,都必须释放锁
            lock.release()
    print(f"\n CPU_C {CPU_number}: the next line is the return")
    return

number_of_threads = 8
# 启动所有线程
threads = []
for i in range(number_of_threads):
    t = threading.Thread(target=something_useful, args=(i,))
    threads.append(t)
    t.start()

# 等待所有线程执行完毕
for t in threads:
    t.join()

方案2:使用结束标记(适合持续任务场景)

如果以后你需要线程一直等待新任务,直到收到明确的结束信号,可以给队列添加特定的结束标记(比如None),每个线程取到标记就退出:

修改后的代码:

import threading
import queue
import time

my_queue = queue.Queue()
lock = threading.Lock()

for i in range(5):
    my_queue.put(i)
# 给每个线程添加一个结束标记
for _ in range(number_of_threads):
    my_queue.put(None)

def something_useful(CPU_number):
    while True:
        lock.acquire()
        item = my_queue.get()
        lock.release()
        if item is None:
            # 收到结束信号,退出循环
            break
        print(f"\n CPU_C {CPU_number}: {item}")
    print(f"\n CPU_C {CPU_number}: the next line is the return")
    return

number_of_threads = 8
# 启动并等待线程
threads = []
for i in range(number_of_threads):
    t = threading.Thread(target=something_useful, args=(i,))
    threads.append(t)
    t.start()

for t in threads:
    t.join()

关键提醒

queue.Queue本身是线程安全的,但empty()方法的结果不能作为可靠判断——因为在多线程环境下,从你调用empty()到执行后续操作的间隙,队列的状态可能已经被其他线程修改了。永远要通过get()的返回值或异常来判断是否还有任务需要处理。

内容的提问来源于stack exchange,提问作者Mohamad Zeina

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:32:58