两个线程跑一遍
(每道题开头都有同一段:上面的有界队列。)
贯穿全条的有界队列(Condition 同步:满了 put 等、空了 get 等,等的时候把锁交出去):
box = Box(cap) 一个最多装 cap 个的队列
box.put(x) 满了就等到有空位,再放进去、叫醒等着取的
box.get() 空了就等到有东西,再取最前面的、叫醒等着放的
run_pc(cap, n) 一个生产者放 0..n-1,一个消费者取 n 个 → 消费到的顺序(FIFO,确定)
run_mpmc(cap, 生产者数, 消费者数, 每人几个) → (消费总数, 消费集合==生产集合) 多对多只问聚合量
run_poison(cap, n, 消费者数) 生产完放几颗毒丸(None),消费者见毒丸就退 → (消费到的真数据数, 是否都退出了)run_pc 就是「生产者线程放 0..n-1、消费者线程取 n 个」。容量 2、放 4 个,看消费顺序和个数:
import threading
class Box:
"""有界队列:满了 put 等、空了 get 等;等的时候把锁交出去(wait_for 放锁,醒来再拿回)。"""
def __init__(self, cap):
self.cap = cap
self.items = []
self.cond = threading.Condition()
def put(self, x):
with self.cond:
self.cond.wait_for(lambda: len(self.items) < self.cap)
self.items.append(x)
self.cond.notify_all()
def get(self):
with self.cond:
self.cond.wait_for(lambda: len(self.items) > 0)
x = self.items.pop(0)
self.cond.notify_all()
return x
def run_pc(cap, n):
"""一个生产者放 0..n-1,一个消费者取 n 个;交回消费者拿到的顺序(FIFO,确定)。"""
box = Box(cap)
got = []
def producer():
for i in range(n):
box.put(i)
def consumer():
for _ in range(n):
got.append(box.get())
p = threading.Thread(target=producer, daemon=True)
c = threading.Thread(target=consumer, daemon=True)
p.start(); c.start()
p.join(timeout=2); c.join(timeout=2)
return got
def run_mpmc(cap, nprod, ncons, per):
"""nprod 个生产者各放 per 个(编号 生产者号*100 + 序号),ncons 个消费者一起取,直到取满;
交回 (消费到的总数, 消费集合 == 生产集合)。只问聚合量——谁先谁后不确定,总数和集合确定。"""
box = Box(cap)
produced = set()
got = []
got_lock = threading.Lock()
total = nprod * per
def producer(pid):
for j in range(per):
item = pid * 100 + j
produced.add(item)
box.put(item)
def consumer():
while True:
with got_lock:
if len(got) >= total:
return
x = box.get()
with got_lock:
got.append(x)
ps = [threading.Thread(target=producer, args=(k,), daemon=True) for k in range(nprod)]
cs = [threading.Thread(target=consumer, daemon=True) for _ in range(ncons)]
for t in ps + cs:
t.start()
for t in ps:
t.join(timeout=2)
# 生产完了,给消费者补足唤醒(放哨兵占位让还在等的消费者醒来看到 len>=total)
for _ in range(ncons):
box.put(-1)
for t in cs:
t.join(timeout=2)
real = [x for x in got if x != -1]
return len(real), set(real) == produced
def run_poison(cap, n, ncons):
"""一个生产者放 0..n-1,放完再放 ncons 颗毒丸(None);ncons 个消费者见到毒丸就退出。
交回 (消费到的真数据个数, 是不是所有消费者都退出了)。毒丸让「优雅退出」变确定。"""
box = Box(cap)
got = []
got_lock = threading.Lock()
stopped = []
def producer():
for i in range(n):
box.put(i)
for _ in range(ncons):
box.put(None)
def consumer():
while True:
x = box.get()
if x is None:
stopped.append(1)
return
with got_lock:
got.append(x)
p = threading.Thread(target=producer, daemon=True)
cs = [threading.Thread(target=consumer, daemon=True) for _ in range(ncons)]
p.start()
for t in cs:
t.start()
p.join(timeout=2)
for t in cs:
t.join(timeout=2)
return len(got), len(stopped) == ncons
def run_poison_few(cap, n, ncons, npoison):
"""毒丸放少了:生产 0..n-1,只放 npoison 颗毒丸;ncons 个消费者。交回 (消费到的真数据数, 是否都退出了)。
npoison < ncons 时会有消费者永远等——join(timeout) 兜住,返回 False。"""
box = Box(cap)
got = []
got_lock = threading.Lock()
stopped = []
def producer():
for i in range(n):
box.put(i)
for _ in range(npoison):
box.put(None)
def consumer():
while True:
x = box.get()
if x is None:
stopped.append(1)
return
with got_lock:
got.append(x)
p = threading.Thread(target=producer, daemon=True)
cs = [threading.Thread(target=consumer, daemon=True) for _ in range(ncons)]
p.start()
for t in cs:
t.start()
p.join(timeout=2)
for t in cs:
t.join(timeout=1)
return len(got), len(stopped) == ncons
def backlog(prod, cons, ticks):
"""无界队列:每个 tick 进 prod 个、出 cons 个(出不能超过在队里的);交回每个 tick 结束时队列的长度。
生产快于消费(prod > cons)时会一直涨——积压。"""
q = 0
out = []
for _ in range(ticks):
q += prod
q -= min(cons, q)
out.append(q)
return out
def bounded(n, cap, cons_every):
"""有界队列 cap:要放 n 个,生产者每 tick 想放 1 个(满了就等,不算这一次),消费者每 cons_every 个 tick 取 1 个。
交回 (每 tick 结束时的队列长度, 生产者被挡下的次数, 全部排空用了多少 tick)。"""
q = 0
made = 0
blocked = 0
trace = []
t = 0
while made < n or q > 0:
t += 1
if (t % cons_every) == 0 and q > 0: # 消费者到点,取一个
q -= 1
if made < n: # 生产者想放一个
if q < cap:
q += 1; made += 1
else:
blocked += 1 # 满了,被背压挡下
trace.append(q)
if t > 100000:
break
return trace, blocked, t
def drops(n, cap, cons_every):
"""满了就丢(put_nowait):要放 n 个,满时这一个直接丢掉;消费者每 cons_every 个 tick 取一个。交回丢了几个。"""
q = 0
dropped = 0
for t in range(1, n * cons_every + cons_every + 1):
if (t % cons_every) == 0 and q > 0:
q -= 1
if n > 0:
if q < cap:
q += 1
else:
dropped += 1
n -= 1
return dropped
got = run_pc(2, 4)
print("".join(str(x) for x in got) + "/" + str(len(got)))
全部评论