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

SimPy单进程服务多资源实现工作保持轮询调度问题求解

问题说明

使用SimPy实现带2个队列的工作保持轮询调度器时,出现2号队列数据包始终无法被调度的问题。调度器预期逻辑为:

  • 两个队列都有数据包时,按1号队列→2号队列的顺序交替服务
  • 其中一个队列为空时,直接服务非空队列的数据包(工作保持特性,不浪费链路带宽)
  • 所有完成服务的数据包统一发送到公共输出端口
    参考现有实现运行后,仅1号队列数据包可正常被服务,2号队列数据包始终得不到处理,原有实现代码如下:
class RoundRobinQueue(object):
    def __init__(self, env, rate, qlimit=None, limit_bytes=True):
        self.store = simpy.Store(env)
        self.store2 = simpy.Store(env)
        self.rate = rate
        self.env = env
        self.out = None
        self.packets_rec = 0
        self.packets_drop = 0
        self.qlimit = qlimit
        self.limit_bytes = limit_bytes
        self.byte_size = 0  # 当前队列总字节数
        self.busy = 0  # 标记是否正在传输数据包
        self.action = env.process(self.run())  # 启动调度进程
        self.trigger = 1

    def run(self):
        while True:
            if (self.trigger == 0 and len(self.store.items)>=0):
                self.trigger = 1
                msg = (yield self.store.get())
                self.byte_size -= msg.size
                self.busy = 1
                yield self.env.timeout(msg.size * 8.0 / self.rate)
                self.out.put(msg)
                self.busy = 0
            else:
                self.trigger = 1

            if (self.trigger == 1 and len(self.store2.items)>=0):
                self.trigger = 0
                msg2 = (yield self.store2.get())
                self.byte_size -= msg2.size
                self.busy = 1
                yield self.env.timeout(msg2.size * 8.0 / self.rate)
                self.out.put(msg2)
                self.busy = 0
            else:
                self.trigger = 0
错误原因

原有代码存在两个核心逻辑问题:

  1. 队列非空判断条件恒成立:len(self.store.items)>=0和len(self.store2.items)>=0永远为真(队列长度最小为0,不可能为负),导致每次循环都会直接进入第一个if分支,执行yield self.store.get()。如果此时1号队列为空,进程会永久阻塞在等待1号队列新数据包的位置,完全不会执行后续2号队列的判断逻辑,哪怕2号队列已经积压数据包也无法被处理。
  2. trigger状态切换逻辑混乱:没有处理双队列为空时的等待逻辑,也没有实现“优先服务轮询指向队列、空则跳转到非空队列”的工作保持逻辑,状态翻转完全和队列实际状态脱节。
修复方案

重构调度逻辑,核心改动点:

  • 用turn变量标记当前轮次优先服务的队列,0对应1号队列,1对应2号队列
  • 每次循环先检查优先队列是否有包,有则直接服务,服务后翻转轮次标记
  • 优先队列为空时检查另一个队列,有包则直接服务,服务后翻转轮次标记
  • 两个队列都为空时,同时等待两个队列的入包事件,任意队列收到包就取消另一个队列的等待事件,处理对应数据包,避免进程卡死在单个队列的等待上
  • 补充入队方法的丢包、字节计数逻辑(原代码缺失该部分)

修复后的完整代码:

import simpy

class RoundRobinQueue(object):
    def __init__(self, env, rate, qlimit=None, limit_bytes=True):
        self.store = simpy.Store(env)  # 1号队列
        self.store2 = simpy.Store(env) # 2号队列
        self.rate = rate
        self.env = env
        self.out = None
        self.packets_rec = 0
        self.packets_drop = 0
        self.qlimit = qlimit
        self.limit_bytes = limit_bytes
        self.byte_size = 0
        self.busy = 0
        self.turn = 0  # 0=下一次优先服务1号队列,1=优先服务2号队列
        self.action = env.process(self.run())

    def put(self, pkt):
        """数据包入队方法,需保证传入的pkt带qid属性,1表示入1号队列,2表示入2号队列"""
        self.packets_rec += 1
        # 队列长度检查,超过限制则丢包
        current_qlen = len(self.store.items) + len(self.store2.items)
        current_size = self.byte_size + pkt.size
        if self.qlimit:
            if (self.limit_bytes and current_size > self.qlimit) or (not self.limit_bytes and current_qlen +1 > self.qlimit):
                self.packets_drop +=1
                return
        # 按标记入对应队列
        if pkt.qid == 1:
            self.byte_size += pkt.size
            return self.store.put(pkt)
        else:
            self.byte_size += pkt.size
            return self.store2.put(pkt)

    def run(self):
        while True:
            msg = None
            # 先检查当前轮次指向的队列
            if self.turn == 0:
                if len(self.store.items) > 0:
                    msg = yield self.store.get()
                elif len(self.store2.items) >0:
                    msg = yield self.store2.get()
            else:
                if len(self.store2.items) >0:
                    msg = yield self.store2.get()
                elif len(self.store.items) >0:
                    msg = yield self.store.get()
            
            # 如果两个队列都为空,等待任意一个队列入包
            if msg is None:
                get1 = self.store.get()
                get2 = self.store2.get()
                # 等待任意一个队列有数据包
                res = yield self.env.any_of([get1, get2])
                # 取消另一个未完成的get请求,避免数据包被错误取出
                if get1 in res:
                    msg = res[get1]
                    if get2.triggered:
                        # 极端情况两个队列同时入包,把多取的包放回原队列
                        yield self.store2.put(res[get2])
                    else:
                        get2.cancel()
                else:
                    msg = res[get2]
                    if get1.triggered:
                        yield self.store.put(res[get1])
                    else:
                        get1.cancel()
            
            # 处理取到的数据包
            self.byte_size -= msg.size
            self.busy = 1
            # 按链路速率计算传输时延
            yield self.env.timeout(msg.size * 8.0 / self.rate)
            self.out.put(msg)
            self.busy = 0
            # 翻转轮次标记
            self.turn = 1 - self.turn
逻辑验证说明
  • 两个队列都有数据包时,turn标记每次服务后翻转,严格按1→2→1→2的顺序交替服务
  • 单队列有数据包时,会持续服务该队列数据包,不会空等轮次指向的空队列,满足工作保持特性
  • 双队列为空时进程挂起,不会空转消耗仿真资源,任意队列入包后立刻恢复调度
  • 不会出现进程永久阻塞在单个队列等待上的问题,两个队列的数据包都能被正常调度

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 11:33:20