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
错误原因
原有代码存在两个核心逻辑问题:
- 队列非空判断条件恒成立:
len(self.store.items)>=0和len(self.store2.items)>=0永远为真(队列长度最小为0,不可能为负),导致每次循环都会直接进入第一个if分支,执行yield self.store.get()。如果此时1号队列为空,进程会永久阻塞在等待1号队列新数据包的位置,完全不会执行后续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
相关产品推荐
相关产品推荐

