如何通过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
报错原因
- 变量未初始化:
protocol仅在流表项存在ip_proto匹配字段时才会被赋值,若遍历所有流后未找到该字段,protocol未定义,直接返回会触发未绑定局部变量错误。 - 冗余代码:循环内的
m = {}未被使用,属于无效代码。 - 逻辑缺陷:原代码会覆盖之前找到的
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
相关产品推荐
相关产品推荐

