无界队列一直涨

👁️ 3 人浏览 💬 0 人评论 ❤️ 添加收藏

本节的积压离散模型(纯算术,不靠真线程):

backlog(prod, cons, ticks)   无界队列:每 tick 进 prod 个、出 cons 个 → 每 tick 结束时队列长度(prod > cons 就一直涨)
bounded(n, cap, cons_every)  有界队列 cap:放 n 个,生产者每 tick 想放 1 个(满就等),消费者每 cons_every 个 tick 取 1 个 → (每 tick 队列长, 被挡下几次, 排空用了几 tick)
drops(n, cap, cons_every)    满了就丢(put_nowait)时,丢了几个

每个 tick 进 3 个、出 1 个,跑 5 个 tick,看队列长度怎么变:

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

tr = backlog(3, 1, 5)
print(",".join(str(x) for x in tr) + "/" + str(tr[-1]))
提交你的答案
请登录后提交答案。
去登录
代码编辑器
Ctrl + Enter 运行
本次输入:
输出:

                        
👩‍🏫
AI
💬 题目评论

全部评论