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

如何用Python编写Ryu控制器应用获取交换机队列统计信息

问题:Ryu控制器获取队列统计信息时出现OFPT_ERROR错误

需要定期获取队列统计信息以计算队列长度,防止队列持续增长,编写了对应的Ryu控制器应用,但运行ryu-manager后出现错误。

第一种实现代码

from ryu.base import app_manager
from ryu.controller import ofp_event
from ryu.controller.handler import CONFIG_DISPATCHER, MAIN_DISPATCHER,DEAD_DISPATCHER
from ryu.controller.handler import set_ev_cls
from ryu.ofproto import ofproto_v1_3
from ryu.lib.packet import packet
from ryu.lib.packet import ethernet
from ryu.lib.packet import ether_types
from operator import attrgetter
from time import sleep
from ryu.app import simple_switch_13
from ryu.lib import hub
from ryu.ofproto.ofproto_v1_3_parser import OFPQueueStatsRequest
from ryu.ofproto.ofproto_v1_3_parser import OFPQueueStats

class simplequeues(simple_switch_13.SimpleSwitch13):
    def __init__(self, *args, **kwargs):
        super(simplequeues, self).__init__(*args, **kwargs)
        self.mac_to_port = {}
        self.datapaths = {}
        self.monitor_thread = hub.spawn(self._monitorqueues)
    @set_ev_cls(ofp_event.EventOFPStateChange,
                [MAIN_DISPATCHER, DEAD_DISPATCHER])
    def _state_change_handler(self, ev):
        datapath = ev.datapath
        if ev.state == MAIN_DISPATCHER:
            if datapath.id not in self.datapaths:
                self.logger.debug('register datapath: %016x', datapath.id)
                self.datapaths[datapath.id] = datapath
        elif ev.state == DEAD_DISPATCHER:
            if datapath.id in self.datapaths:
                self.logger.debug('unregister datapath: %016x', datapath.id)
                del self.datapaths[datapath.id]
    def _monitorqueues(self):
        while True:
            for dp in self.datapaths.values():
                self.send_queue_stats_request(dp)
            hub.sleep(30)

    def send_queue_stats_request(self, datapath):
        self.logger.debug('send stats request: %016x', datapath.id)
        ofp = datapath.ofproto
        ofp_parser = datapath.ofproto_parser
        req = ofp_parser.OFPQueueStatsRequest(datapath, ofp.OFPP_ANY,ofp.OFPQ_ALL)
        print('*'*20)
        datapath.send_msg(req)
        print('send*'*20)
    @set_ev_cls(ofp_event.EventOFPQueueStatsReply, MAIN_DISPATCHER)
    def queue_stats_reply_handler(self, ev):
        queues = []
        for stat in ev.msg.body:
            queues.append('port_no=%d queue_id=%d '
                            'tx_bytes=%d tx_packets=%d tx_errors=%d '
                            'duration_sec=%d duration_nsec=%d' %
                            (stat.port_no, stat.queue_id,
                            stat.tx_bytes, stat.tx_packets, stat.tx_errors,
                            stat.duration_sec, stat.duration_nsec))
        self.logger.debug('QueueStats: %s', queues)

第二种实现代码

@set_ev_cls(ofp_event.EventOFPStatsReply, MAIN_DISPATCHER)
def stats_reply_handler(self, ev):
    print("!!"*30,"start reply")
    msg = ev.msg
    ofp = msg.datapath.ofproto
    body = ev.msg.body
    print("!!"*30,"start reply condition")
    if msg.type == ofp.OFPST_QUEUE:
        print("!!"*30,"start request")
        self.queue_stats_reply_handler(body)

def queue_stats_reply_handler(self, body):
    print("start 2 reply")
    queues = []
    for stat in body:
        queues.append('port_no=%d queue_id=%d '
                        'tx_bytes=%d tx_packets=%d tx_errors=%d ' %
                        (stat.port_no, stat.queue_id,
                        stat.tx_bytes, stat.tx_packets, stat.tx_errors))
    self.logger.debug('QueueStats: %s', queues)

错误输出

EventOFPErrorMsg received.
version=0x4, msg_type=0x1, msg_len=0x28, xid=0x29c9db81
 `-- msg_type: OFPT_ERROR(1)
OFPErrorExperimenterMsg(type=0xffff, exp_type=0xa50, experimenter=0x4f4e4600, data=b'\x04\x12\x00\x18\x29\xc9\xdb\x81\x00\x05\x0f\xff\x00\x00\x00\x00\x00\x00\x0f\xff\x00\x00\x0f\xff')
 `-- data: version=0x4, msg_type=0x12, msg_len=0x18, xid=0x29c9db81
     `-- msg_type: OFPT_MULTIPART_REQUEST(18)

问题原因与解决办法

原因分析

错误中的OFPErrorExperimenterMsg表明交换机不支持通过OFPQ_ALL请求所有队列统计,或者请求格式不符合交换机预期。OpenFlow 1.3中,部分交换机(如OVS)在处理队列统计请求时,对OFPP_ANY/OFPQ_ALL的支持存在限制,需要更明确的请求参数。

修复步骤

  1. 修改队列统计请求方式
    不要直接使用ofp.OFPP_ANY和ofp.OFPQ_ALL,改为遍历交换机的所有物理端口,针对每个端口发送队列统计请求。可以先通过端口统计获取端口列表,再逐个端口请求:

    def send_queue_stats_request(self, datapath):
        self.logger.debug('send stats request: %016x', datapath.id)
        ofp = datapath.ofproto
        ofp_parser = datapath.ofproto_parser
        # 替换为实际拓扑中的物理端口号,或先获取端口列表
        target_ports = [1, 2, 3]
        for port_no in target_ports:
            req = ofp_parser.OFPQueueStatsRequest(datapath, port_no, ofp.OFPQ_ALL)
            datapath.send_msg(req)
    
  2. 检查交换机队列配置
    确保交换机已启用队列支持并配置了队列,以OVS为例:

    ovs-vsctl set port <port-name> qos=@newqos -- \
        --id=@newqos create qos type=linux-htb queues=0=@q0,1=@q1 -- \
        --id=@q0 create queue other-config:max-rate=100000000 -- \
        --id=@q1 create queue other-config:max-rate=200000000
    
  3. 统一回复处理逻辑
    推荐使用第一种实现中的EventOFPQueueStatsReply专用事件处理,而非第二种通用的EventOFPStatsReply,前者是Ryu针对队列统计回复的专属处理逻辑,准确性更高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 12:37:28