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

如何通过Ryu控制器采集SDN交换机流表IP协议并写入CSV

Ryu控制器采集流表IP协议并保存到CSV的问题解决

问题描述

我尝试用Ryu控制器从SDN交换机采集流表统计信息,需要实现通过OFPFlowStatsRequest和回复处理器采集每个流的ip_protocol并保存到CSV文件的功能。我写了一段提取协议的代码,但运行报错。

我的代码

def _protocol(self, dpid, flows):
       
        for flow in flows:
           m = {}
           for i in flow.match.items():
               key = list(i)[0]  # match key
               val = list(i)[1]  # match value
               if key == "ip_proto":
                   protocol = val
        return protocol

报错信息

File "/home/ai/SDN/daapp.py", line 780, in flow_stats_reply_handler
    protocol  = self._protocol(dpid, gflows[dpid])
  File "/home/ai/SDN/daapp.py", line 458, in _protocol
    return protocol
UnboundLocalError: local variable 'protocol' referenced before assignment

报错原因

  1. 变量未初始化:protocol仅在流表项存在ip_proto匹配字段时才会被赋值,若遍历所有流后未找到该字段,protocol未定义,直接返回会触发未绑定局部变量错误。
  2. 冗余代码:循环内的m = {}未被使用,属于无效代码。
  3. 逻辑缺陷:原代码会覆盖之前找到的ip_proto值,最终仅返回最后一个匹配项的协议,不符合“采集每个流的ip_protocol”的需求。

修复后的完整代码

from ryu.base import app_manager
from ryu.controller import ofp_event
from ryu.controller.handler import CONFIG_DISPATCHER, MAIN_DISPATCHER
from ryu.controller.handler import set_ev_cls
from ryu.ofproto import ofproto_v1_3
import csv
import time

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

    def __init__(self, *args, **kwargs):
        super(FlowStatsCollector, self).__init__(*args, **kwargs)
        self.datapaths = {}
        # 初始化CSV文件,写入表头
        self._init_csv()

    def _init_csv(self):
        with open('flow_ip_protocols.csv', 'w', newline='') as f:
            writer = csv.writer(f)
            writer.writerow(['dpid', 'flow_id', 'ip_protocol', 'collect_time'])

    @set_ev_cls(ofp_event.EventOFPSwitchFeatures, CONFIG_DISPATCHER)
    def switch_features_handler(self, ev):
        datapath = ev.msg.datapath
        self.datapaths[datapath.id] = datapath
        # 发送首次流统计请求,并设置10秒定时请求
        self.send_flow_stats_request(datapath)
        self.logger.info("Switch %d connected, start collecting flow stats", datapath.id)

    def send_flow_stats_request(self, datapath):
        ofproto = datapath.ofproto
        parser = datapath.ofproto_parser

        req = parser.OFPFlowStatsRequest(datapath, 0, ofproto.OFPTT_ALL,
                                        ofproto.OFPP_ANY, ofproto.OFPG_ANY, 0, 0)
        datapath.send_msg(req)
        self.logger.debug("Sent flow stats request to switch %d", datapath.id)
        # 每10秒重复发送请求
        self.send_event_after(10, self.send_flow_stats_request, datapath)

    @set_ev_cls(ofp_event.EventOFPFlowStatsReply, MAIN_DISPATCHER)
    def flow_stats_reply_handler(self, ev):
        flows = ev.msg.body
        datapath = ev.msg.datapath
        collect_time = time.strftime("%Y-%m-%d %H:%M:%S")

        # 遍历每个流表项,提取ip_protocol
        for flow_idx, flow in enumerate(flows):
            ip_proto = None
            for match_key, match_val in flow.match.items():
                if match_key == 'ip_proto':
                    ip_proto = match_val
                    break  # 找到后停止遍历当前流的匹配字段
            # 追加写入CSV文件
            with open('flow_ip_protocols.csv', 'a', newline='') as f:
                writer = csv.writer(f)
                writer.writerow([datapath.id, flow_idx, ip_proto, collect_time])
        self.logger.info("Collected %d flows from switch %d", len(flows), datapath.id)

代码说明

  • 初始化时创建flow_ip_protocols.csv并写入表头,包含交换机ID、流序号、IP协议、采集时间四个字段。
  • 交换机连接后自动发送流统计请求,并设置每10秒重复请求一次(可根据需求修改定时时间)。
  • 处理流统计回复时,逐个遍历流表项,提取ip_proto字段(不存在则设为None),并将数据追加写入CSV文件,避免覆盖历史数据。
  • 修复了原代码中变量未定义的问题,同时满足“采集每个流的ip_protocol”的需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 12:45:49