Python 基础体系 · 第 42/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python collections:deque、Counter、defaultdict、ChainMap 与队列

collections 模块提供了一组针对特定数据组织方式优化的容器。它们不是“更复杂的字典或列表”,而是把常见的数据结构约束直接编码进 API:

  • deque:双端队列,适合从两端快速进出数据;
  • Counter:多重集合,适合统计可哈希对象出现的次数;
  • defaultdict:带缺省值工厂的字典,适合分组、聚合和索引构建;
  • ChainMap:多个映射的动态查找视图,适合配置覆盖和嵌套作用域;
  • queueasyncio.Queue:面向并发生产者—消费者模型的队列,提供同步、等待、任务跟踪、容量限制和关闭协议。

这些类型都建立在熟悉的 dictlistset 和迭代器之上。理解它们的关键,不是记住方法名,而是明确三件事:

  1. 数据的逻辑结构是什么;
  2. 哪些操作需要高效;
  3. 空、满、缺失、覆盖、关闭等边界状态应该如何表示。

一、先区分三类“容器问题”

假设有以下需求:

events = [
    ("login", "alice"),
    ("download", "alice"),
    ("login", "bob"),
    ("login", "alice"),
]

我们可能需要:

  1. 保留最近 100 条事件;
  2. 统计每种事件出现了多少次;
  3. 按用户分组事件;
  4. 让命令行配置覆盖环境变量,再覆盖默认配置;
  5. 让多个线程或协程异步处理事件。

这五个需求看起来都可以用 listdict 拼出来,但它们的核心约束不同:

需求 核心结构 合适工具
两端进出、先进先出、滑动窗口 双端序列 deque
元素到出现次数的映射 多重集合 Counter
键不存在时自动创建容器 缺省映射 defaultdict
多层配置按优先级查找 映射链 ChainMap
并发任务传递、等待和背压 生产者—消费者队列 queueasyncio.Queue

其中,“队列”有两个容易混淆的含义:

  • deque 是一种底层容器,可以被用来实现简单队列;
  • queue.Queueasyncio.Queue 是带并发协调语义的队列组件。

前者主要解决数据如何存放,后者还要解决线程或协程如何等待、何时唤醒、任务是否完成以及队列如何关闭。


二、deque:双端队列

2.1 逻辑模型

deque 的全称是 double-ended queue,即双端队列。它维护一个有序序列,并允许在左端和右端分别执行插入和删除:

左端                                      右端
  ↓                                        ↓
[ A ] <-> [ B ] <-> [ C ] <-> [ D ]

appendleft(X)                         append(Y)
popleft()                             pop()

典型操作是:

from collections import deque

d = deque(["B", "C"])

d.appendleft("A")  # [A, B, C]
d.append("D")      # [A, B, C, D]

left = d.popleft() # left == "A",剩余 [B, C, D]
right = d.pop()    # right == "D",剩余 [B, C]

deque 两端的追加和弹出具有近似 O(1) 的性能;相比之下,列表的 pop(0)insert(0, value) 通常需要移动大量元素,属于 O(n) 操作。这里的复杂度描述的是操作随元素数量增长的趋势,不代表每次操作具有固定的微秒级耗时。(docs.python.org)


2.2 为什么 list.pop(0) 不适合队列

列表的逻辑顺序如下:

index:  0    1    2    3
value: [ A,   B,   C,   D ]

执行:

items.pop(0)

删除 A 后,为了让新的第一个元素仍然位于索引 0,实现需要把后面的元素向左移动:

删除前:[ A, B, C, D ]
删除后:[ B, C, D ]
         ↑  ↑  ↑
         移动

当列表中有 n 个元素时,平均需要处理与 n 成正比的数据。因此,使用列表模拟大型 FIFO 队列会让出队成为瓶颈。

deque 不要求所有元素都连续位于一个需要整体搬移的数组区域中,因此两端操作不会因为队列长度增加而线性移动全部元素。随机访问则不同:访问两端较快,访问中间位置会逐渐变慢;需要大量随机访问时,列表更合适。(docs.python.org)


2.3 deque 的两种方向

可以把 deque 同时当作队列或栈:

FIFO 队列

from collections import deque

queue = deque()

queue.append("task-1")
queue.append("task-2")
queue.append("task-3")

while queue:
    print(queue.popleft())

输出:

task-1
task-2
task-3

数据从右端进入,从左端离开:

append(task-1) → append(task-2) → append(task-3)
popleft()       → task-1
popleft()       → task-2
popleft()       → task-3

LIFO 栈

from collections import deque

stack = deque()

stack.append("A")
stack.append("B")
stack.append("C")

while stack:
    print(stack.pop())

输出:

C
B
A

数据从右端进入,也从右端离开。deque 因而可以表达栈和队列,但它本身不提供优先级排序、任务完成计数或并发等待。


2.4 extendleft() 为什么会反转输入

下面的代码容易产生误解:

from collections import deque

d = deque([3])
d.extendleft([1, 2])

print(d)

输出是:

deque([2, 1, 3])

原因是 extendleft(iterable) 等价于依次执行左侧追加:

for item in iterable:
    d.appendleft(item)

执行过程:

初始:        [3]
appendleft(1) [1, 3]
appendleft(2) [2, 1, 3]

所以,如果希望最终左侧保持 [1, 2] 的顺序,应当传入反向迭代器:

d = deque([3])
d.extendleft(reversed([1, 2]))

print(d)
# deque([1, 2, 3])

这不是特殊的排序行为,而是“连续从左侧插入”自然产生的结果。(docs.python.org)


2.5 有界 deque:固定大小的滑动窗口

deque(maxlen=n) 表示一个最多容纳 n 个元素的双端队列:

from collections import deque

recent = deque(maxlen=3)

for value in [10, 20, 30, 40, 50]:
    recent.append(value)
    print(list(recent))

输出:

[10]
[10, 20]
[10, 20, 30]
[20, 30, 40]
[30, 40, 50]

当右端追加导致长度超过 maxlen 时,左端的旧元素会被自动丢弃:

追加 40:

[10, 20, 30] → [20, 30, 40]

这与 queue.Queue(maxsize=3) 的语义完全不同:

  • deque(maxlen=3):满了以后自动淘汰另一端的数据;
  • queue.Queue(maxsize=3):满了以后生产者等待,或者抛出 Full

前者适合“只关心最近数据”的日志窗口、最近访问记录和滑动统计;后者适合“每个任务都必须处理”的工作队列。deque 的有界模式与 Unix tail 保留最后若干行的行为类似。(docs.python.org)


2.6 用 deque 实现滑动窗口

例如,计算窗口大小为 3 的移动平均值:

from collections import deque

def moving_average(values, window_size):
    if window_size <= 0:
        raise ValueError("window_size must be positive")

    window = deque()
    total = 0

    for value in values:
        window.append(value)
        total += value

        if len(window) > window_size:
            total -= window.popleft()

        yield total / len(window)


values = [10, 20, 30, 40, 50]

print(list(moving_average(values, 3)))

输出:

[10.0, 15.0, 20.0, 30.0, 40.0]

以最后三个值为完整窗口时,计算过程是:

加入 10:窗口 [10]          总和 10      平均 10.0
加入 20:窗口 [10, 20]       总和 30      平均 15.0
加入 30:窗口 [10, 20, 30]   总和 60      平均 20.0
加入 40:移除 10             总和 90      平均 30.0
加入 50:移除 20             总和 120     平均 40.0

如果每一步都对窗口重新调用 sum(window),每次需要扫描窗口,单步成本为 O(k),其中 k 是窗口大小。维护一个增量变量 total 后,每次只需加上新元素、减去旧元素,单步成本降为近似 O(1)


2.7 rotate():循环调度与轮转

rotate(n) 将右端元素轮转到左端;负数表示向左轮转:

from collections import deque

d = deque(["A", "B", "C", "D"])

d.rotate(1)
print(d)  # deque(['D', 'A', 'B', 'C'])

d.rotate(-2)
print(d)  # deque(['B', 'C', 'D', 'A'])

右移一步可以理解为:

d.appendleft(d.pop())

左移一步可以理解为:

d.append(d.popleft())

因此可以用 deque 实现轮询调度:

from collections import deque

def round_robin(*iterables):
    iterators = deque(map(iter, iterables))

    while iterators:
        current = iterators[0]

        try:
            yield next(current)
            iterators.rotate(-1)
        except StopIteration:
            iterators.popleft()


print(list(round_robin("ABC", "D", "EF")))

输出:

['A', 'D', 'E', 'B', 'F', 'C']

每轮只从当前迭代器取一个元素,然后把它移到队尾;耗尽的迭代器则从左端删除。这里的公平性是“每个仍有数据的迭代器每轮获得一次机会”,不是对每个任务的执行时间提供严格保证。


2.8 deque 的边界和误用

空队列弹出会抛异常

from collections import deque

d = deque()

d.popleft()

结果:

IndexError: pop from an empty deque

如果不希望异常,可以先判断:

if d:
    item = d.popleft()

但在多线程场景中,“先判断再弹出”不是完整的同步协议:判断之后,另一个线程可能已经改变了队列状态。需要并发等待和唤醒时,应使用 queue.Queue,而不是自行用 deque 加锁拼装一个不完整的队列。

中间插入不是双端操作

d.insert(100, value)

虽然 deque 支持 insert(),但它的核心优势是两端操作,而不是任意位置操作。中间索引、搜索、删除和插入不应被误认为与 append()popleft() 具有相同成本。

maxlen 会静默丢数据

from collections import deque

d = deque(maxlen=2)
d.extend(["job-1", "job-2", "job-3"])

print(list(d))
# ['job-2', 'job-3']

如果 job-1 代表不可重试的业务任务,这种静默淘汰就是数据丢失,而不是“队列满时的自然行为”。


三、Counter:带计数的多重集合

3.1 Counter 的数学模型

Counterdict 的子类:

from collections import Counter

counts = Counter(["red", "blue", "red"])

print(counts)
# Counter({'red': 2, 'blue': 1})

它的逻辑结构可以写成:

C:KVC: K \rightarrow V

其中:

  • KK 是可哈希对象的集合;
  • VV 是计数值;
  • C[x] 表示元素 x 的计数。

与普通字典不同的是,读取不存在的键时,Counter 返回 0

counts = Counter(["red"])

print(counts["blue"])
# 0

这使得下面的增量统计自然成立:

counts["blue"] += 1

其概念过程是:

counts["blue"]      → 0
0 + 1               → 1
写回 counts["blue"] → 1

计数值允许为零或负数,但常见的多重集合运算主要面向正数计数。Counter 的键必须是可哈希对象;值通常是数值,但类本身并不强制值必须是整数。(docs.python.org)


3.2 初始化方式决定输入含义

从可迭代对象初始化

Counter("banana")

把字符串视为元素序列,结果是:

Counter({'a': 3, 'n': 2, 'b': 1})

从映射初始化

Counter({"apple": 3, "orange": 2})

这里的值已经是计数,不会把字典的键重复展开。

从关键字参数初始化

Counter(apples=3, oranges=2)

关键字参数的键必须是合法标识符,因此不能表达任意字符串键:

Counter({"red-blue": 2})  # 可以
# Counter(red-blue=2)    # 语法错误

3.3 update() 与普通 dict.update() 不同

普通字典的 update() 会覆盖旧值:

d = {"a": 2}
d.update({"a": 5})
print(d)
# {'a': 5}

Counter.update() 则是累加:

from collections import Counter

c = Counter({"a": 2})
c.update({"a": 5})

print(c)
# Counter({'a': 7})

对可迭代对象调用 update() 时,输入应当是元素序列,而不是 (key, value) 对:

c = Counter()
c.update(["a", "a", "b"])

print(c)
# Counter({'a': 2, 'b': 1})

若已有键值对,应使用映射或构造新的 Counter

c = Counter(dict([("a", 2), ("b", 3)]))

3.4 subtract() 可以产生负数

from collections import Counter

stock = Counter({"apple": 5, "orange": 2})
used = Counter({"apple": 3, "orange": 4, "banana": 1})

stock.subtract(used)

print(stock)
# Counter({'apple': 2, 'orange': -2, 'banana': -1})

这表示:

apple:5 - 3 =  2
orange:2 - 4 = -2
banana:0 - 1 = -1

负数可以表示“缺口”“撤销超过新增”“借贷差额”等状态。不要在业务语义不允许负数时,直接把 Counter 运算结果当作合法库存。

如果只想保留正数计数,可以使用一元加法:

print(+stock)
# Counter({'apple': 2})

+c 会移除零和负数;-c 则把负数转成对应的正数。(docs.python.org)


3.5 多重集合运算

设:

from collections import Counter

c = Counter(a=3, b=1)
d = Counter(a=1, b=2)

加法

c + d
# Counter({'a': 4, 'b': 3})

数学上:

(c+d)[x]=c[x]+d[x](c+d)[x] = c[x] + d[x]

正数截断的减法

c - d
# Counter({'a': 2})

虽然:

a = 3 - 1 = 2
b = 1 - 2 = -1

Counter 的多重集合减法只保留正数结果,因此 b 被丢弃。

交集

c & d
# Counter({'a': 1, 'b': 1})

交集取对应计数的最小值:

(cd)[x]=min(c[x],d[x])(c \cap d)[x] = \min(c[x], d[x])

并集

c | d
# Counter({'a': 3, 'b': 2})

并集取对应计数的最大值:

(cd)[x]=max(c[x],d[x])(c \cup d)[x] = \max(c[x], d[x])

这些运算的结果会排除零和负数计数。(docs.python.org)


3.6 most_common()elements()total()

统计前 N 项

from collections import Counter

c = Counter("abracadabra")

print(c.most_common(3))
# [('a', 5), ('b', 2), ('r', 2)]

计数相同的元素按首次遇到的顺序排列。该方法要求计数值可比较。

展开元素

c = Counter(a=3, b=1, c=0, d=-1)

print(list(c.elements()))
# ['a', 'a', 'a', 'b']

elements() 只处理正整数计数;零和负数会被忽略。它并不是把所有键值对机械展开,而是把每个键重复其正整数计数次。

计算总数

c = Counter(a=3, b=1, c=0, d=-1)

print(c.total())
# 3

total() 等价于对计数值求和,因此如果 Counter 中存在负数,它可能小于正数计数之和,甚至为负数。(docs.python.org)


3.7 零计数键仍然存在

from collections import Counter

c = Counter(a=2)
c["a"] = 0

print("a" in c)
# True

print(list(c))
# ['a']

将计数设置为零不会删除键:

del c["a"]

print("a" in c)
# False

这会影响:

  • len(c)
  • list(c)
  • dict(c)
  • 序列化结果;
  • 某些业务中的“是否出现过”判断。

Python 3.10 起,Counter 的相等比较把缺失键视为零,因此:

Counter(a=1) == Counter(a=1, b=0)
# True

但“数学上计数相同”与“内部是否保留了零计数键”仍然是两个不同问题。(docs.python.org)


四、defaultdict:缺失键的构造协议

4.1 defaultdict 到底重写了什么

defaultdictdict 的子类。它增加了一个可写属性 default_factory,并通过 __missing__() 处理使用 d[key] 访问时不存在的键:

from collections import defaultdict

d = defaultdict(list)

items = d["missing"]

print(items)
# []

print(d)
# defaultdict(<class 'list'>, {'missing': []})

缺失键访问的过程是:

1. 查找 key
2. 发现 key 不存在
3. 调用 default_factory()
4. 将返回值写入 d[key]
5. 返回该值

因此,defaultdict(list) 的行为不是“返回一个临时空列表”,而是“创建并保存一个新的空列表”。

如果 default_factoryNone,缺失访问仍会抛出 KeyError。如果工厂函数自身抛出异常,该异常会原样传播。(docs.python.org)


4.2 __missing__() 只对 d[key] 生效

这是最重要的边界之一:

from collections import defaultdict

d = defaultdict(list)

print(d["a"])
# []

print(d.get("b"))
# None

print("b" in d)
# False

d["a"] 会触发 __missing__(),而 d.get("b") 不会。inkeys()items() 等操作也不会因为查询缺失键而创建条目。

因此,下面的代码会改变字典:

if d["user-1"]:
    ...

即使列表为空,"user-1" 也已经被插入。只读探测应使用:

value = d.get("user-1")

或者:

if "user-1" in d:
    ...

4.3 分组:defaultdict(list)

使用普通字典分组时,通常需要显式初始化:

groups = {}

for key, value in [("red", 1), ("blue", 2), ("red", 3)]:
    if key not in groups:
        groups[key] = []
    groups[key].append(value)

defaultdict(list) 把“如果没有就创建空列表”编码进容器:

from collections import defaultdict

groups = defaultdict(list)

for key, value in [("red", 1), ("blue", 2), ("red", 3)]:
    groups[key].append(value)

print(dict(groups))
# {'red': [1, 3], 'blue': [2]}

处理第一条 ("red", 1) 时:

groups["red"] 不存在
→ 调用 list()
→ 写入 groups["red"] = []
→ append(1)
→ groups["red"] == [1]

处理第二条 ("red", 3) 时,键已存在,直接取得原来的列表并追加。

等价的 setdefault() 写法是:

groups = {}

for key, value in [("red", 1), ("blue", 2), ("red", 3)]:
    groups.setdefault(key, []).append(value)

两者结果相同,但 defaultdict 更明确地表达了“这个映射的值类型是列表”。(docs.python.org)


4.4 defaultdict(set):去重分组

如果每个分组中的值不能重复:

from collections import defaultdict

relations = [
    ("alice", "admin"),
    ("alice", "admin"),
    ("alice", "reader"),
    ("bob", "reader"),
]

roles = defaultdict(set)

for user, role in relations:
    roles[user].add(role)

print(dict(roles))
# {'alice': {'admin', 'reader'}, 'bob': {'reader'}}

这里要区分两个层次:

  • 外层键 user:哈希表键;
  • 内层值 set:用于去重。

由于集合不承诺业务排序,如果输出顺序重要,应在输出阶段排序,而不是依赖集合的遍历顺序:

for user, user_roles in sorted(roles.items()):
    print(user, sorted(user_roles))

4.5 defaultdict(int)Counter 的区别

二者都能计数:

from collections import Counter, defaultdict

text = "mississippi"

a = defaultdict(int)
for char in text:
    a[char] += 1

b = Counter(text)

print(dict(a))
print(b)

结果都表示:

{'m': 1, 'i': 4, 's': 4, 'p': 2}

但抽象不同:

  • defaultdict(int) 是“键到值的映射,缺失时生成 0”;
  • Counter 是“元素到计数的映射”,额外提供 most_common()elements()total() 以及多重集合运算。

如果只是聚合数值,defaultdict(int) 足够;如果需要计数领域的操作,Counter 更合适。


4.6 多层 defaultdict 与递归创建

可以构造嵌套映射:

from collections import defaultdict

tree = defaultdict(lambda: defaultdict(list))

tree["service-a"]["error"].append("timeout")
tree["service-a"]["info"].append("started")

print(tree["service-a"]["error"])
# ['timeout']

但嵌套缺省工厂有两个风险:

  1. 任何使用 [] 的查询都会创建中间层;
  2. 调试和序列化时,结构可能比预期更大。

如果数据结构层级复杂,显式定义工厂函数通常比长的匿名 lambda 更容易阅读:

def nested_list():
    return defaultdict(list)

tree = defaultdict(nested_list)

五、ChainMap:多个映射的动态查找视图

5.1 合并字典和链接字典不是一回事

假设有三层配置:

defaults = {
    "host": "localhost",
    "port": 8000,
}

environment = {
    "port": 9000,
}

command_line = {
    "host": "example.com",
}

常见的优先级是:

命令行参数 > 环境变量 > 默认值

使用 ChainMap

from collections import ChainMap

config = ChainMap(command_line, environment, defaults)

print(config["host"])
# example.com

print(config["port"])
# 9000

ChainMap 不会立即复制所有字典,而是保存一个映射列表。查找时依次搜索:

config["host"]
→ command_line 中找到 example.com
→ 停止

config["port"]
→ command_line 没有
→ environment 中找到 9000
→ 停止

底层映射会被按引用连接,因此原始字典发生变化时,ChainMap 的结果也会反映变化。查找从第一个映射到最后一个映射进行,而写入、更新和删除只作用于第一个映射。(docs.python.org)


5.2 ChainMap 的写入规则

from collections import ChainMap

local = {"debug": True}
env = {"port": 9000}
defaults = {"port": 8000, "host": "localhost"}

config = ChainMap(local, env, defaults)

config["port"] = 7000
config["host"] = "example.com"

print(local)
print(env)
print(defaults)

结果:

{'debug': True, 'port': 7000, 'host': 'example.com'}
{'port': 9000}
{'port': 8000, 'host': 'localhost'}

即使 port 原本来自 env,赋值也不会修改 env,而是在最前面的 local 中创建或覆盖 port

删除也只会尝试删除第一层:

del config["port"]

这会删除 local["port"]。如果第一层没有该键,即使后面的映射存在,也不会直接删除后层键。

因此,ChainMap 的语义是:

读取:从前往后查找
写入:只写第一层
删除:只删第一层

这与“一个合并后的普通字典”不同。普通字典复制完成后,后续原始映射的变化不会自动反映进去。


5.3 new_child()parents

new_child() 用于创建一个新的前置上下文:

from collections import ChainMap

root = ChainMap({"theme": "light"})
child = root.new_child({"user": "alice"})

print(child["user"])
# alice

print(child["theme"])
# light

child["theme"] = "dark"

print(child.maps)
# [{'user': 'alice', 'theme': 'dark'}, {'theme': 'light'}]

child 新增了一个最前面的映射,因此对 child 的修改不会直接写入 root 的第一层。

parents 则跳过当前第一层:

print(child.parents["theme"])
# light

它等价于:

ChainMap(*child.maps[1:])

这使 ChainMap 能够模拟嵌套作用域:

当前局部作用域
    ↓
外层作用域
    ↓
全局默认作用域

需要注意,new_child() 创建的是新的 ChainMap 视图和新的前置映射;后面的原始映射仍然是共享引用。(docs.python.org)


5.4 ChainMap 的迭代顺序

查找顺序和迭代顺序不应混为一谈。

from collections import ChainMap

baseline = {"music": "bach", "art": "rembrandt"}
adjustments = {"art": "van gogh", "opera": "carmen"}

config = ChainMap(adjustments, baseline)

print(config["art"])
# van gogh

print(list(config))
# ['music', 'art', 'opera']

art 的值来自第一层,但迭代顺序相当于从最后一个映射开始执行 dict.update(),再逐层向前更新。这样既保留首次出现的键顺序,又让前层值覆盖后层值。(docs.python.org)

如果需要一个独立的普通字典,可以显式展平:

flat = dict(config)

展平之后:

  • 后续原始映射变化不会影响 flat
  • flat 不再保留多层写入语义;
  • 复制可能产生额外内存成本。

5.5 配置覆盖的完整示例

下面的程序实现:

命令行参数 > 环境变量 > 默认配置
import argparse
import os
from collections import ChainMap


def load_config():
    parser = argparse.ArgumentParser()
    parser.add_argument("--host")
    parser.add_argument("--port", type=int)
    args = parser.parse_args()

    cli = {
        key: value
        for key, value in vars(args).items()
        if value is not None
    }

    env = {}
    if "APP_HOST" in os.environ:
        env["host"] = os.environ["APP_HOST"]
    if "APP_PORT" in os.environ:
        env["port"] = int(os.environ["APP_PORT"])

    defaults = {
        "host": "127.0.0.1",
        "port": 8000,
    }

    return ChainMap(cli, env, defaults)


config = load_config()
print(f"{config['host']}:{config['port']}")

例如:

APP_HOST=10.0.0.8 APP_PORT=9000 python app.py --port 7000

输出:

10.0.0.8:7000

host 没有命令行值,因此从环境变量读取;port 同时存在于命令行和环境变量中,前面的命令行映射优先。

生产代码中还应考虑环境变量类型转换失败:

try:
    env["port"] = int(os.environ["APP_PORT"])
except ValueError as exc:
    raise ValueError("APP_PORT must be an integer") from exc

ChainMap 只负责查找优先级,不负责配置校验、类型转换或必填字段检查。


六、从 deque 到并发队列

6.1 数据容器不等于并发协调器

下面的 deque 可以实现单线程 FIFO:

from collections import deque

tasks = deque()

tasks.append("task-1")
task = tasks.popleft()

但多线程生产者—消费者还需要处理:

  • 队列为空时,消费者等待;
  • 队列满时,生产者等待;
  • 一个线程放入数据后,唤醒等待的消费者;
  • 一个线程取走数据后,唤醒等待的生产者;
  • 所有任务处理完成后,通知 join()
  • 关闭时,阻止新的生产并唤醒阻塞线程。

queue 模块正是为多生产者、多消费者线程场景设计的同步队列。它内部使用锁来协调竞争线程,但不支持同一线程内的可重入操作。(docs.python.org)


6.2 queue.Queue 的核心状态

以有容量的 FIFO 队列为例:

from queue import Queue

q = Queue(maxsize=2)

可以把它抽象为:

容量:2
当前项目数:0
未完成任务数:0
状态:运行中

执行:

q.put("A")

状态变为:

当前项目数:1
未完成任务数:1

再执行:

q.put("B")

状态变为:

当前项目数:2
未完成任务数:2

此时继续执行阻塞式:

q.put("C")

生产者会等待,直到消费者通过 get() 取走一个项目。

消费者执行:

item = q.get()

状态变为:

当前项目数:1
未完成任务数:2

注意:get() 只表示“取到了任务”,不表示“任务处理完成”。处理结束后必须调用:

q.task_done()

状态才变为:

当前项目数:1
未完成任务数:1

join() 等待的是未完成任务数归零,而不是简单等待队列变空。每次 put() 增加一个未完成任务;每次匹配的 task_done() 减少一个。(docs.python.org)


6.3 线程生产者—消费者完整示例

from queue import Queue, ShutDown
from threading import Thread
import time


def producer(q):
    try:
        for i in range(5):
            item = f"task-{i}"
            q.put(item)
            print(f"produced: {item}")
    finally:
        # 通知队列不再接受新的任务。
        q.shutdown()


def consumer(q):
    while True:
        try:
            item = q.get()
        except ShutDown:
            print("consumer: queue is shut down")
            return

        try:
            print(f"consuming: {item}")
            time.sleep(0.05)
        finally:
            q.task_done()


q = Queue(maxsize=2)

producer_thread = Thread(target=producer, args=(q,))
consumer_thread = Thread(target=consumer, args=(q,))

consumer_thread.start()
producer_thread.start()

producer_thread.join()

# 等待已经放入队列的任务全部处理完成。
q.join()

consumer_thread.join()

运行逻辑如下:

生产者 put(task-0)
生产者 put(task-1)
队列达到容量上限
生产者 put(task-2) 等待
消费者 get(task-0)
生产者继续 put(task-2)
消费者处理完成并 task_done()
...
生产者完成生产并 shutdown()
消费者继续消费队列中剩余任务
队列耗尽后,get() 抛出 ShutDown
消费者退出

这里的 shutdown() 是 Python 3.13 引入的队列关闭机制,在 Python 3.14 中可用。关闭后,新的 put() 会抛出 ShutDown;非立即关闭时,队列中已经存在的项目仍可被取出处理,直到队列耗尽。(docs.python.org)

如果消费者在处理任务时抛出异常,task_done() 仍应放在 finally 中,否则 join() 可能永久等待:

try:
    item = q.get()
    process(item)
finally:
    q.task_done()

但要注意,task_done() 只表示“该任务的生命周期结束”,不表示业务处理一定成功。失败任务是否重试、转入死信队列或记录错误,需要由应用层决定。


6.4 qsize()empty()full() 不是同步判断

下面的写法存在竞态:

if not q.empty():
    item = q.get()

即使 empty() 返回 False,另一个消费者也可能在 get() 执行前取走最后一个项目。

同样:

if not q.full():
    q.put(item)

返回后,其他生产者可能已经填满队列。

queue.Queue.qsize()empty()full() 反映的是近似或瞬时状态,不能作为后续阻塞操作一定成功的保证。更可靠的方式是直接调用 get()put(),或者使用 block=False / 超时并处理 EmptyFull。(docs.python.org)

例如:

from queue import Full

try:
    q.put_nowait(item)
except Full:
    # 立即失败:可以丢弃、重试、降级或记录指标
    handle_overflow(item)

6.5 FIFO、LIFO 和优先级队列

queue 模块提供三种主要取出顺序:

Queue

先进先出:

from queue import Queue

q = Queue()
q.put("A")
q.put("B")

print(q.get())  # A
print(q.get())  # B

LifoQueue

后进先出:

from queue import LifoQueue

q = LifoQueue()
q.put("A")
q.put("B")

print(q.get())  # B
print(q.get())  # A

PriorityQueue

取出优先级最低的项目:

from queue import PriorityQueue

q = PriorityQueue()
q.put((20, "normal"))
q.put((1, "urgent"))
q.put((10, "important"))

print(q.get())
# (1, 'urgent')

默认情况下,队列比较元组。如果优先级相同而数据对象不可比较,可能出现 TypeError。可以把数据包裹在只比较优先级的对象中:

from dataclasses import dataclass, field
from queue import PriorityQueue
from typing import Any


@dataclass(order=True)
class PrioritizedItem:
    priority: int
    item: Any = field(compare=False)


q = PriorityQueue()
q.put(PrioritizedItem(1, {"request_id": "r-1"}))
q.put(PrioritizedItem(1, {"request_id": "r-2"}))

print(q.get().item)

field(compare=False) 使排序只看 priority,不会尝试比较字典对象。(docs.python.org)


6.6 SimpleQueue 的取舍

SimpleQueue 是无界 FIFO 队列,提供较少功能:

from queue import SimpleQueue

q = SimpleQueue()
q.put("A")
print(q.get())

它不提供 task_done()join() 这类任务跟踪功能。因此:

  • 只需要线程安全地传递项目:可以考虑 SimpleQueue
  • 需要容量限制、背压、任务完成等待:使用 Queue
  • 需要 LIFO 或优先级:使用对应队列类型。

“简单”并不表示“更适合所有场景”。没有任务跟踪时,调用方无法通过 join() 知道所有工作是否已处理完成。(docs.python.org)


七、背压:容量限制如何控制生产速度

7.1 无界队列的问题

设生产者平均每秒产生 λp\lambda_p 个任务,消费者平均每秒处理 λc\lambda_c 个任务。

当:

λp>λc\lambda_p > \lambda_c

队列长度会随着时间增长。忽略初始长度,经过时间 tt 后,积压量近似为:

Q(t)Q(0)+(λpλc)tQ(t) \approx Q(0) + (\lambda_p - \lambda_c)t

如果队列无界,短时间内系统看似“没有阻塞”,但代价是:

  • 内存持续增长;
  • 任务等待时间持续增加;
  • 进程可能因内存耗尽被终止;
  • 故障被延迟到更难诊断的地方。

如果队列容量为 MM,队列达到上限后,生产者必须采取一种策略:

阻塞等待
立即失败
丢弃
降级
写入外部持久化系统

queue.Queue(maxsize=M)asyncio.Queue(maxsize=M) 默认选择“等待空位”,这就是背压:下游处理不过来时,上游的生产速度被迫降低。


7.2 容量不是吞吐量

maxsize=100 只表示最多存放 100 个项目,不代表:

  • 每秒能处理 100 个任务;
  • 任务最多等待 100 秒;
  • 队列中项目占用的内存不超过固定值。

如果每个项目大小差异很大,按“项目数量”限制容量可能不够。一个包含大型字节串的项目和一个小整数都占用一个队列槽位。

更准确的容量控制可能需要应用层维护:

当前字节数 ≤ 最大字节预算

标准队列的 maxsize 通常按项目个数计数,而不是按项目内存大小计数。


八、asyncio.Queue:协程之间的异步队列

8.1 与 queue.Queue 的根本区别

asyncio.Queue 面向同一个事件循环中的协程,不是线程安全队列:

import asyncio

queue = asyncio.Queue(maxsize=2)

它的操作通常是可等待的:

item = await queue.get()
await queue.put(item)

与线程队列相比:

  • queue.Queue:线程阻塞;
  • asyncio.Queue:协程挂起,让事件循环运行其他任务;
  • queue.Queue:设计为线程安全;
  • asyncio.Queue:不提供线程安全保证;
  • asyncio.Queue:没有 timeout 参数,需要使用 asyncio.wait_for() 包裹队列操作。

asyncio.Queueqsize() 可以直接获得当前项目数;但它仍然不应被用作复杂并发判断的替代品。(docs.python.org)


8.2 异步生产者—消费者示例

import asyncio


async def producer(queue):
    for i in range(5):
        item = f"task-{i}"
        await queue.put(item)
        print(f"produced: {item}")

    # 不再接受新任务,但允许消费者处理队列中已有任务。
    queue.shutdown()


async def consumer(name, queue):
    while True:
        try:
            item = await queue.get()
        except asyncio.QueueShutDown:
            print(f"{name}: queue is shut down")
            return

        try:
            print(f"{name} consuming: {item}")
            await asyncio.sleep(0.05)
        finally:
            queue.task_done()


async def main():
    queue = asyncio.Queue(maxsize=2)

    consumers = [
        asyncio.create_task(consumer("worker-1", queue)),
        asyncio.create_task(consumer("worker-2", queue)),
    ]

    await producer(queue)

    # 等待所有已经入队的项目完成处理。
    await queue.join()

    await asyncio.gather(*consumers)


asyncio.run(main())

数据流是:

flowchart LR
    P[producer 协程] -->|await put| Q[asyncio.Queue]
    Q -->|await get| W1[worker-1]
    Q -->|await get| W2[worker-2]
    W1 -->|task_done| Q
    W2 -->|task_done| Q
    P -->|shutdown| Q
    Q -->|耗尽后 QueueShutDown| W1
    Q -->|耗尽后 QueueShutDown| W2

当队列达到 maxsize 时,await queue.put() 会挂起,而不是阻塞整个线程。事件循环可以在此期间运行消费者;消费者取走项目后,生产者才继续。


8.3 超时必须在外部实现

asyncio.Queueput()get() 没有 timeout 参数。需要超时时,可以使用:

import asyncio

async def get_with_timeout(queue):
    try:
        return await asyncio.wait_for(queue.get(), timeout=1.0)
    except asyncio.TimeoutError:
        return None

这表示:

最多等待 1 秒
→ 有项目:返回项目
→ 超时:抛出 TimeoutError

不要写成:

await queue.get(timeout=1)

这不是 asyncio.Queue 支持的 API。(docs.python.org)


8.4 正常关闭与立即关闭

Python 3.13 起,asyncio.Queue 提供 shutdown(immediate=False)

正常关闭

queue.shutdown()

语义是:

  1. 队列不再接受新的 put()
  2. 已经入队的项目仍可通过 get() 取出;
  3. 每个项目都调用 task_done() 后,join() 正常结束;
  4. 队列耗尽后,后续 get() 抛出 QueueShutDown

这适合“生产阶段结束,但要处理完已有任务”的场景。

立即关闭

queue.shutdown(immediate=True)

语义更强:

  1. 队列立即被清空;
  2. 尚未处理的任务被丢弃;
  3. 等待中的 get() 被唤醒并抛出 QueueShutDown
  4. join() 可能在任务实际没有完成的情况下解除阻塞。

因此,立即关闭会破坏通常的 join() 不变量:

join() 返回
≠
所有任务都已成功处理

它适合进程退出、不可恢复故障或明确允许丢弃剩余任务的场景,不适合作为普通的“生产者完成”信号。(docs.python.org)


8.5 取消协程时仍要维护任务计数

消费者从队列取出项目后,如果在处理过程中被取消,必须考虑 task_done() 是否仍然会执行:

async def consumer(queue):
    while True:
        item = await queue.get()
        try:
            await process(item)
        finally:
            queue.task_done()

如果取消发生在 get() 之后、task_done() 之前,且没有 finally,那么未完成任务计数可能永远不会归零,等待中的 join() 会一直挂起。

另一方面,task_done() 不应在没有成功取得项目时调用。否则会出现:

ValueError: task_done() called too many times

任务计数必须严格满足:

task_done 次数put 次数\text{task\_done 次数} \leq \text{put 次数}

并且对于每个已经 get() 的项目,最终应有且只有一次对应的 task_done()。(docs.python.org)


九、dequequeueasyncio.Queue 的选择

工具 并发安全 是否阻塞等待 是否有容量背压 是否任务跟踪 关闭协议
deque 仅适合容器级简单操作 maxlen 是淘汰,不是等待
queue.Queue 线程安全 shutdown()
queue.SimpleQueue 线程安全 否,无界 不提供任务跟踪
asyncio.Queue 非线程安全 协程挂起 shutdown()

可以按语义选择:

需要两端快速操作
    → deque

需要保留最近 N 条
    → deque(maxlen=N)

需要统计次数、排名或多重集合运算
    → Counter

需要缺失键自动创建列表、集合或计数
    → defaultdict

需要配置层级查找,且保留底层映射的动态变化
    → ChainMap

多个线程交换任务
    → queue.Queue / LifoQueue / PriorityQueue

多个协程交换任务
    → asyncio.Queue / LifoQueue / PriorityQueue

需要任务完成确认
    → Queue 或 asyncio.Queue,而不是 deque

允许满时丢弃旧数据
    → deque(maxlen=N)

满时必须等待或失败,并保留每个任务
    → 有界 queue.Queue 或 asyncio.Queue

最容易发生的错误,是用一种容器表达另一种语义:

  • deque(maxlen=N) 模拟不能丢任务的工作队列;
  • 用无界 Queue() 掩盖下游持续过载;
  • Counterc[key] == 0 判断键是否存在;
  • defaultdictd[key] 做只读查询,意外创建数据;
  • ChainMap 赋值,以为会修改命中键所在的底层映射;
  • qsize() 判断下一步 put()get() 必然不会阻塞;
  • asyncio.Queue 跨线程传递对象;
  • shutdown(immediate=True) 后仍然把 join() 返回当作任务成功完成。

十、把几种容器组合起来

下面的例子模拟一个事件处理管道:

  • asyncio.Queue:限制待处理事件数量;
  • Counter:统计事件类型;
  • defaultdict(set):维护每个用户访问过的资源;
  • deque(maxlen=3):保留最近三条事件;
  • ChainMap:管理配置覆盖。
import asyncio
from collections import ChainMap, Counter, defaultdict, deque


async def process_events(events, config):
    queue = asyncio.Queue(maxsize=config["queue_size"])
    counts = Counter()
    user_resources = defaultdict(set)
    recent_events = deque(maxlen=config["recent_size"])

    async def worker():
        while True:
            try:
                event = await queue.get()
            except asyncio.QueueShutDown:
                return

            try:
                event_type = event["type"]
                user = event["user"]
                resource = event["resource"]

                counts[event_type] += 1
                user_resources[user].add(resource)
                recent_events.append(event)

                await asyncio.sleep(0)
            finally:
                queue.task_done()

    workers = [
        asyncio.create_task(worker())
        for _ in range(config["workers"])
    ]

    try:
        for event in events:
            await queue.put(event)

        queue.shutdown()
        await queue.join()
    finally:
        await asyncio.gather(*workers)

    return {
        "counts": counts,
        "user_resources": dict(user_resources),
        "recent_events": list(recent_events),
    }


async def main():
    defaults = {
        "queue_size": 2,
        "recent_size": 3,
        "workers": 2,
    }

    environment = {
        "workers": 3,
    }

    command_line = {
        # 假设这是解析命令行参数得到的结果。
    }

    config = ChainMap(command_line, environment, defaults)

    events = [
        {"type": "login", "user": "alice", "resource": "/"},
        {"type": "download", "user": "alice", "resource": "/a.zip"},
        {"type": "login", "user": "bob", "resource": "/"},
        {"type": "login", "user": "alice", "resource": "/"},
        {"type": "logout", "user": "alice", "resource": "/"},
    ]

    result = await process_events(events, config)

    print(result["counts"])
    print(result["user_resources"])
    print(result["recent_events"])


asyncio.run(main())

这段代码中,每个容器承担的职责不同:

ChainMap
    决定配置来源和优先级

asyncio.Queue
    决定生产和消费如何等待,以及最大积压量

Counter
    把事件类型映射到累计次数

defaultdict(set)
    为新用户自动创建集合,并对资源去重

deque(maxlen=3)
    自动淘汰过旧事件,只保留最近窗口

这种组合的价值不在于“使用了很多容器”,而在于每个数据结构都明确表达了一个不变量:

配置:前层覆盖后层
队列:未超过容量才允许继续放入
计数:同类事件累加
用户资源:同一资源只保留一份
最近事件:长度最多为 3

当程序出错时,也可以根据不变量定位问题:

  • 最近事件少于预期:检查是否错误使用了 maxlen
  • 用户映射出现空用户:检查是否使用 user_resources[user] 做了意外查询;
  • join() 永远不返回:检查消费者是否遗漏 task_done()
  • 队列关闭后仍有生产者异常:检查生产者是否在 shutdown() 后继续 put()
  • 配置值不符合预期:检查映射链顺序,而不是只检查最终字典内容。

dequeCounterdefaultdictChainMap 和并发队列的共同点,是把常见控制逻辑下沉到数据结构中;它们的差异,则在于下沉的控制逻辑不同。选择容器时,先写出需要保持的状态不变量,再选择能够直接表达该不变量的类型,通常比先从方法名出发更可靠。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。