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

如何在Ryu控制器中实现每5个h1→h2流仅将第5个上报至控制器?

Ryu控制器实现每第5个指定流数据包上报控制器的解决方案

问题根源

你当前的方案没有下发针对性的流表项,导致所有h1(10.0.0.1)->h2(10.0.0.3)的数据包都会触发packet_in上报控制器;但如果直接下发永久流表项,又会切断控制器与后续数据包的联系,无法完成计数。核心矛盾是既要让大部分数据包在交换机层面直接处理,又要让控制器能精准捕获每第5个数据包。

解决思路

  1. 按流维度维护计数:不再用全局count,而是针对(交换机ID, 源IP, 目的IP)的组合维护计数,避免不同流的计数混乱。
  2. 下发带数据包数量限制的临时流表项:当处理非第5个数据包时,下发流表项让交换机直接丢弃后续若干个同流数据包,直到第5个数据包触发packet_in,再重置计数并更新流表项。
  3. 完善数据包处理流程:补充IPv4协议判断、packet_out消息发送等必要逻辑,避免空指针或交换机无响应问题。

修改后的完整代码

from ryu.base import app_manager
from ryu.controller import ofp_event
from ryu.controller.handler import MAIN_DISPATCHER, CONFIG_DISPATCHER
from ryu.controller.handler import set_ev_cls
from ryu.ofproto import ofproto_v1_3
from ryu.lib.packet import packet, ipv4

class CountPacketApp(app_manager.RyuApp):
    OFP_VERSIONS = [ofproto_v1_3.OFP_VERSION]

    def __init__(self, *args, **kwargs):
        super(CountPacketApp, self).__init__(*args, **kwargs)
        self.ip_to_port = {}
        # 按(交换机ID, 源IP, 目的IP)维护流计数
        self.flow_counts = {}

    @set_ev_cls(ofp_event.EventOFPSwitchFeatures, CONFIG_DISPATCHER)
    def switch_features_handler(self, event):
        datapath = event.msg.datapath
        ofproto = datapath.ofproto
        parser = datapath.ofproto_parser

        # 下发默认流表,未匹配的数据包上报控制器
        match = parser.OFPMatch()
        actions = [parser.OFPActionOutput(ofproto.OFPP_CONTROLLER, ofproto.OFPCML_NO_BUFFER)]
        self.add_flow(datapath, 0, match, actions)

    def add_flow(self, datapath, priority, match, actions, buffer_id=None, packet_count=None):
        ofproto = datapath.ofproto
        parser = datapath.ofproto_parser

        inst = [parser.OFPInstructionActions(ofproto.OFPIT_APPLY_ACTIONS, actions)]
        if buffer_id:
            mod = parser.OFPFlowMod(datapath=datapath, buffer_id=buffer_id,
                                    priority=priority, match=match,
                                    instructions=inst, packet_count=packet_count)
        else:
            mod = parser.OFPFlowMod(datapath=datapath, priority=priority,
                                    match=match, instructions=inst, packet_count=packet_count)
        datapath.send_msg(mod)

    @set_ev_cls(ofp_event.EventOFPPacketIn, MAIN_DISPATCHER)
    def packet_in_handler(self, event):
        msg = event.msg
        datapath = msg.datapath
        ofproto = datapath.ofproto
        ofp_parser = datapath.ofproto_parser
        dpid = datapath.id
        pkt = packet.Packet(msg.data)

        # 解析IPv4协议,非IPv4包按正常流程转发
        ip_pkt = pkt.get_protocols(ipv4.ipv4)
        if not ip_pkt:
            self.normal_forward(datapath, msg, ofp_parser, ofproto)
            return

        src_ip = ip_pkt[0].src
        dst_ip = ip_pkt[0].dst

        # 更新IP与端口的映射
        self.ip_to_port.setdefault(dpid, {})
        self.ip_to_port[dpid][src_ip] = msg.in_port

        # 处理h1->h2的指定流
        if src_ip == "10.0.0.1" and dst_ip == "10.0.0.3":
            flow_key = (dpid, src_ip, dst_ip)
            # 获取当前流的计数,初始为0
            current_count = self.flow_counts.get(flow_key, 0)
            current_count += 1
            self.flow_counts[flow_key] = current_count

            if current_count % 5 == 0:
                # 第5个数据包,上报控制器
                print(f"第{current_count}个数据包,上报控制器")
                actions = [ofp_parser.OFPActionOutput(ofproto.OFPP_CONTROLLER, ofproto.OFPCML_NO_BUFFER)]
                # 下发流表项,让接下来4个同流数据包直接丢弃
                match = ofp_parser.OFPMatch(eth_type=0x0800, ipv4_src=src_ip, ipv4_dst=dst_ip)
                self.add_flow(datapath, 100, match, [], packet_count=4)
            else:
                # 非第5个数据包,直接丢弃
                actions = []
                # 如果是当前周期的第1个包,下发流表项处理剩余3个包
                if current_count == 1:
                    match = ofp_parser.OFPMatch(eth_type=0x0800, ipv4_src=src_ip, ipv4_dst=dst_ip)
                    self.add_flow(datapath, 100, match, [], packet_count=3)
            
            # 发送packet_out处理当前数据包
            out = ofp_parser.OFPPacketOut(datapath=datapath, buffer_id=msg.buffer_id,
                                        in_port=msg.in_port, actions=actions, data=msg.data)
            datapath.send_msg(out)
        else:
            # 其他流按正常流程转发
            self.normal_forward(datapath, msg, ofp_parser, ofproto)

    def normal_forward(self, datapath, msg, ofp_parser, ofproto):
        """正常转发非指定流的数据包"""
        dpid = datapath.id
        pkt = packet.Packet(msg.data)
        ip_pkt = pkt.get_protocols(ipv4.ipv4)
        if ip_pkt:
            src_ip = ip_pkt[0].src
            dst_ip = ip_pkt[0].dst
            self.ip_to_port.setdefault(dpid, {})
            self.ip_to_port[dpid][src_ip] = msg.in_port

            if dst_ip in self.ip_to_port[dpid]:
                out_port = self.ip_to_port[dpid][dst_ip]
            else:
                out_port = ofproto.OFPP_FLOOD
            
            actions = [ofp_parser.OFPActionOutput(out_port)]
            # 下发流表项优化后续转发
            match = ofp_parser.OFPMatch(eth_type=0x0800, ipv4_src=src_ip, ipv4_dst=dst_ip)
            if msg.buffer_id != ofproto.OFP_NO_BUFFER:
                self.add_flow(datapath, 1, match, actions, msg.buffer_id)
                return
            else:
                self.add_flow(datapath, 1, match, actions)
            
            out = ofp_parser.OFPPacketOut(datapath=datapath, buffer_id=msg.buffer_id,
                                        in_port=msg.in_port, actions=actions, data=msg.data)
            datapath.send_msg(out)

关键部分说明

  1. 流计数维护:用self.flow_counts字典,以(dpid, src_ip, dst_ip)为键,确保每个流的计数独立。
  2. 临时流表项:通过packet_count参数指定流表项处理的数据包数量,达到数量后流表项自动失效,触发下一次packet_in。例如第1个包处理后,下发处理3个包的流表项,覆盖第2-4个包;第5个包处理后,下发处理4个包的流表项,覆盖第6-9个包。
  3. 正常转发逻辑:抽离出normal_forward方法处理非指定流的数据包,同时下发流表项优化后续转发性能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 17:50:35