一、为什么多线程需要 queue?

在多线程编程中,多个线程可能会同时访问和修改同一个共享数据。如果不对这种访问进行控制,就会导致数据错乱、程序崩溃等严重问题,这就是所谓的“线程不安全”。

queue模块提供的是一种线程安全的数据结构——队列。你可以把它想象成一个管道,一端的线程(生产者)往里放东西,另一端的线程(消费者)从里面取东西。这个管道内部自带了锁机制,确保了在任何时刻,只有一个线程能够对它进行操作,从而避免了数据竞争问题。

因此,在 Python 多线程编程中,只要涉及线程间的数据传递、任务分配或状态同步,queue 就是首选的、最标准、最安全、最高效的解决方案


二、queue模块的核心组件

queue模块主要提供了三种类型的队列,它们的核心功能和方法是共通的,只是存取数据的顺序不同。

  1. queue.Queue(maxsize=0)先进先出 (FIFO) 队列

    1. 这是最常用的一种队列。

    2. 行为:最先放入队列的元素,最先被取出。就像排队买票,先来的人先买。

    3. maxsize:一个整数,用于设置队列的容量上限。如果为0或负数,表示队列大小无限制。

  2. queue.LifoQueue(maxsize=0)后进先出 (LIFO) 队列

    1. 行为:最后放入队列的元素,最先被取出。这就像一叠盘子,你总是先用最上面(最后放上去)的那个。这种结构在数据结构中也称为“栈”。

  3. queue.PriorityQueue(maxsize=0)优先级队列

    1. 行为:存入的元素会按照其优先级进行排序,优先级最低的元素最先被取出。

    2. 存储内容:通常存入的是一个元组 (priority_number, data),队列会根据 priority_number 的大小来决定谁先出队。数字越小,优先级越高。


三、queue对象的常用方法和属性

假设我们创建了一个队列实例 q = queue.Queue(),以下是你可以调用的核心方法:

1、 状态检查
  • q.qsize()

    • 作用:返回队列中当前元素的数量。

    • 注意:这个数字不总是完全精确的,因为在你获取到这个值和使用它之间,其他线程可能已经修改了队列。它主要用于调试,不建议在多线程逻辑中依赖它。

  • q.empty()

    • 作用:判断队列是否为空。如果为空,返回 True;否则返回 False

    • 注意:同样存在qsize()的非精确性问题。

  • q.full()

    • 作用:判断队列是否已满。如果已满,返回 True;否则返回 False。只有在创建队列时指定了 maxsize,这个方法才有意义。

2、 放入元素 (生产者)
  • q.put(item, block=True, timeout=None)

    • 作用:将一个元素 item 放入队列中。

    • block=True (默认):这是阻塞模式。如果队列已满(达到了 maxsize),当前线程会被阻塞(暂停执行),直到队列中有空位可以放入新元素为止。

    • block=False:非阻塞模式。如果队列已满,不会等待,而是立即抛出 queue.Full 异常。

    • timeout:当 block=True 时,可以设置一个超时时间(秒)。如果线程被阻塞超过了这个时间,队列依然是满的,那么程序会抛出 queue.Full 异常。

3、 取出元素 (消费者)
  • q.get(block=True, timeout=None)

    • 作用:从队列中取出一个元素并返回。

    • block=True (默认):阻塞模式。如果队列为空,当前线程会被阻塞,直到队列中有新元素可取为止。

    • block=False:非阻塞模式。如果队列为空,不会等待,而是立即抛出 queue.Empty 异常。

    • timeout:当 block=True 时,可以设置一个超时时间(秒)。如果线程被阻塞超过了这个时间,队列依然为空,那么程序会抛出 queue.Empty 异常。

4、 生产者-消费者模型

以下两个方法通常配对使用,用于优雅地管理任务流程:

  • q.task_done()

    • 作用:由消费者线程调用。每当消费者线程完成了一个从队列中获取的任务(get() 出来的任务)时,就调用此方法。它像是在通知队列:“嘿,我处理完一个任务了。”

    • 关键:每次 get() 一个任务,处理完后,必须对应调用一次 task_done()

  • q.join()

    • 作用:由主线程(或生产者线程)调用。它会阻塞当前线程,直到队列中所有的任务都已经被 get() 并且都已经调用了 task_done()

    • 内部机制join() 内部有一个计数器。每次 put() 一个元素,计数器加一。每次调用 task_done(),计数器减一。join() 方法就是持续检查这个计数器,直到它归零,才解除阻塞。

import threading
import queue
import time
import random

# 创建一个容量为 5 的 FIFO 队列
q = queue.Queue(maxsize=5)


# 生产者线程
def producer(name):
    """
    生产者负责生成任务并放入队列。
    """
    for i in range(10):  # 每个生产者生产10个任务
        item = f"任务_{name}_{i}"
        # 模拟生产过程耗时
        time.sleep(random.uniform(0.1, 0.5))

        # 将任务放入队列,如果队列满了会阻塞等待
        q.put(item)
        print(f"✅ 生产者 [{name}] 生产了: {item} (队列大小: {q.qsize()})")


# 消费者线程
def consumer(name):
    """
    消费者负责从队列中取出任务并处理。
    """
    while True:
        try:
            # 从队列中获取任务,如果队列空了会阻塞等待
            # 设置 timeout=2,如果2秒内拿不到任务,就会抛出 queue.Empty 异常
            item = q.get(timeout=2)

            # 模拟处理任务耗时
            time.sleep(random.uniform(0.5, 1.0))
            print(f"🔵 消费者 [{name}] 处理了: {item} (队列大小: {q.qsize()})")

            # 通知队列,这个任务已经处理完毕
            q.task_done()
        except queue.Empty:
            # 如果在超时时间内没有取到任务,说明任务可能已经处理完了
            print(f"🟡 消费者 [{name}] 等待超时,认为任务已完成,退出。")
            break  # 退出循环


if __name__ == "__main__":
    # 记录开始时间
    start_time = time.time()

    # 创建并启动线程
    producer_thread1 = threading.Thread(target=producer, args=("P1",))
    consumer_thread1 = threading.Thread(target=consumer, args=("C1",))
    consumer_thread2 = threading.Thread(target=consumer, args=("C2",))

    producer_thread1.start()
    consumer_thread1.start()
    consumer_thread2.start()

    # 等待生产者完成所有生产任务
    producer_thread1.join()
    print("\n--- 生产者 P1 已经结束生产 ---\n")

    # 主线程在此处阻塞,等待队列中的所有任务都被处理完毕
    # (即,等待所有 task_done() 调用完成)
    q.join()

    # 记录结束时间
    end_time = time.time()

    print(f"\n🎉 所有任务处理完毕!总耗时: {end_time - start_time:.2f} 秒")


# # 输出结果 (可能会略有不同):
# ✅ 生产者 [P1] 生产了: 任务_P1_0 (队列大小: 1)
# ✅ 生产者 [P1] 生产了: 任务_P1_1 (队列大小: 1)
# ✅ 生产者 [P1] 生产了: 任务_P1_2 (队列大小: 1)
# 🔵 消费者 [C2] 处理了: 任务_P1_1 (队列大小: 1)
# 🔵 消费者 [C1] 处理了: 任务_P1_0 (队列大小: 0)
# ✅ 生产者 [P1] 生产了: 任务_P1_3 (队列大小: 1)
# ✅ 生产者 [P1] 生产了: 任务_P1_4 (队列大小: 1)
# ✅ 生产者 [P1] 生产了: 任务_P1_5 (队列大小: 2)
# 🔵 消费者 [C2] 处理了: 任务_P1_2 (队列大小: 2)
# ✅ 生产者 [P1] 生产了: 任务_P1_6 (队列大小: 2)
# 🔵 消费者 [C1] 处理了: 任务_P1_3 (队列大小: 2)
# ✅ 生产者 [P1] 生产了: 任务_P1_7 (队列大小: 2)
# ✅ 生产者 [P1] 生产了: 任务_P1_8 (队列大小: 3)
# 🔵 消费者 [C1] 处理了: 任务_P1_5 (队列大小: 3)
# 🔵 消费者 [C2] 处理了: 任务_P1_4 (队列大小: 2)
# ✅ 生产者 [P1] 生产了: 任务_P1_9 (队列大小: 2)
#
# --- 生产者 P1 已经结束生产 ---
#
# 🔵 消费者 [C1] 处理了: 任务_P1_6 (队列大小: 2)
# 🔵 消费者 [C2] 处理了: 任务_P1_7 (队列大小: 1)
# 🔵 消费者 [C1] 处理了: 任务_P1_8 (队列大小: 0)
# 🔵 消费者 [C2] 处理了: 任务_P1_9 (队列大小: 0)
#
# 🎉 所有任务处理完毕!总耗时: 4.56 秒

代码解读:

  • 生产者 (producer):循环生成任务,并通过 q.put() 放入队列。

  • 消费者 (consumer):在一个无限循环中,通过 q.get() 尝试获取任务。获取到就处理,处理完后必须调用 q.task_done()。如果长时间(2秒)拿不到任务,就认为工作结束并退出。

  • 主线程

    • 启动所有生产者和消费者线程。

    • producer_thread1.join(): 等待生产者线程结束。但这只表示生产者不再放新东西进去了,不代表队列里的东西被处理完了。

    • q.join():这是关键!主线程会在这里“暂停”,直到两个消费者把队列里所有的东西都取出来处理完(即所有put进去的item都有一个对应的task_done调用),才会继续往下执行。

这个模型完美地解耦了生产者和消费者,它们之间通过队列进行通信,互不干扰,实现了高效的并发处理。


4、queue和Lock、Rlock对比

queue 是一个更高层次的抽象,而 Lock RLock 是更低层次的。用一个比喻来解释:

  • queue.Queue:就像一辆整车。它已经为你设计好了引擎、变速箱、刹车和方向盘,并且它们完美地协同工作。你只需要踩油门(put)和刹车(get)就能安全地把货物(数据)从 A 点运到 B 点。

  • threading.Lock/RLock:就像是螺丝、齿轮和刹车片这些零件。它们是制造汽车的基础,非常强大和灵活。你可以用它们来组装任何东西,但如果你不是一个熟练的机械师,你很容易把东西装错,导致汽车在路上解体。

所以,您的原则可以这样理解:如果只是想解决“把货物从A点运到B点”的问题,那就直接买一辆整车 (queue),而不是自己从零件 (Lock) 开始造车。


4.1 为什么“能用 queue 就不用 Lock”?

因为 queue 已经为你完美地处理了所有底层的 Lock 管理。当你调用 q.put(item) 时,queue 内部会自动执行以下操作:

  1. 获取一个锁 (lock.acquire())。

  2. 将数据放入内部的列表。

  3. 释放这个锁 (lock.release())。

  4. (如果之前有线程因队列为空而等待)通知等待的线程“有新货了”。

q.get() 的过程类似。这一切都是原子性的、线程安全的,并且经过了专家级的优化。

使用 queue 的好处:

  1. 更简单:你的业务代码只关心放数据和取数据,完全不用理会复杂的加锁、释放锁、线程通知等细节。

  2. 更安全:避免了忘记释放锁而导致的死锁问题。queue 的设计保证了同步的正确性。

  3. 目标明确queue 的设计目标就是用于线程间的数据/任务传递。当你的需求是这个时,它就是最对口的工具。

4.2 什么时候必须使用 LockRLock

既然 queue 这么好,为什么还需要 Lock 呢?

因为有些问题不是“运送货物”,而是“保护一块共享区域,确保一段时间内的多个操作不被打断”queue 解决的是数据传递问题,而 Lock 解决的是状态共享原子性**问题。以下是必须使用 LockRLock 的场景:

场景一:修改共享状态 (Modifying Shared State)

当多个线程需要修改同一个全局变量共享对象时。

例子:网站的在线用户计数器 假设有一个全局变量 online_users_count = 0

  • 一个线程(用户登录)需要执行 online_users_count += 1

  • 另一个线程(用户登出)需要执行 online_users_count -= 1

  • 还有一个线程(后台统计)可能需要读取这个值。

这个问题不是把“用户+1”这个任务传递给谁,而是要保证在执行 += 1 这个非原子操作时,不被其他线程干扰。这时就必须用 Lock

import threading

online_users_count = 0
lock = threading.Lock()

def user_login():
    global online_users_count
    with lock:  # with 语句会自动获取和释放锁
        online_users_count += 1
        print(f"用户登录,当前在线人数: {online_users_count}")

def user_logout():
    global online_users_count
    with lock:
        online_users_count -= 1
        print(f"用户登出,当前在线人数: {online_users_count}")
场景二:保证一系列操作的原子性 (Ensuring Atomicity)

当一个操作包含多个步骤,而你必须保证这些步骤作为一个整体完成,中间不能有其他线程插入。

经典例子:银行转账 从账户 A 转 100 元到账户 B,包含两个步骤:

  1. 账户 A 的余额减少 100。

  2. 账户 B 的余额增加 100。

你必须保证这两步之间,不能有其他线程来读取账户 A 或 B 的余额,否则可能会读到“钱已经从A扣了,但还没到B”的中间错误状态。

import threading

accounts = {'A': 1000, 'B': 1000}
lock = threading.Lock()

def transfer(from_acc, to_acc, amount):
    with lock:
        # --- 临界区开始 ---
        if accounts[from_acc] >= amount:
            print(f"开始转账... {from_acc} -> {to_acc} ${amount}")
            accounts[from_acc] -= amount
            # 在这里模拟一个耗时操作,增加线程切换的可能性
            time.sleep(0.001)
            accounts[to_acc] += amount
            print("转账成功!")
        else:
            print("余额不足,转账失败!")
        # --- 临界区结束 ---

在这个例子里,queue 完全无用武之地。需要的是一个 Lock 来保护“转账”这一系列操作的完整性。

4.3 总结
工具 抽象层次 核心用途 何时使用
queue.Queue 线程间的数据/任务传递 (Producer-Consumer) 当你的关注点是“把东西从一个线程安全地交给另一个线程”时。
threading.Lock 保护共享状态,保证操作原子性 (Critical Section) 当你的关注点是“保护一段代码或一个变量,确保同一时间只有一个线程能访问它”时。
threading.RLock 与Lock相同,但允许同一线程重复获取 (Re-entrant) 在 Lock 的使用场景中,如果代码涉及递归或复杂的嵌套调用,可能会导致同一线程重复加锁。

更多推荐