如何高效提取PCAP文件的TCP流并分析iRTT、重传率等参数?
用Scapy实现TCP流提取与参数分析
针对你的需求,以下是基于Scapy的解决方案,解决内存占用问题的同时,实现TCP流拆分和iRTT、重传率等参数计算:
一、高效拆分TCP流(解决内存错误)
Scapy的PcapReader可以逐包读取PCAP文件,无需一次性加载所有数据包到内存,完美解决Pyshark遇到的内存溢出问题。同时通过排序源/目的IP+端口生成流标识,确保双向通信的数据包归为同一个流。
from scapy.all import PcapReader, IP, TCP from collections import defaultdict def get_tcp_stream_key(pkt): """生成唯一的TCP流标识,确保双向包归为同一流""" src_endpoint = (pkt[IP].src, pkt[TCP].sport) dst_endpoint = (pkt[IP].dst, pkt[TCP].dport) # 排序端点,让A->B和B->A的包属于同一个流 return tuple(sorted((src_endpoint, dst_endpoint))) def split_tcp_streams(pcap_path): """逐包读取PCAP,拆分TCP流""" streams = defaultdict(list) with PcapReader(pcap_path) as reader: for pkt in reader: if IP in pkt and TCP in pkt: stream_key = get_tcp_stream_key(pkt) streams[stream_key].append(pkt) return streams
二、计算TCP流核心参数
1. 参数定义与实现逻辑
- iRTT(初始RTT):TCP连接建立阶段,SYN包到SYN-ACK包的时间差
- 重传率:重传数据包数量 / 流内总TCP包数量 × 100%,重传包通过跟踪已发送的序列号范围识别(接近Wireshark的
tcp.analysis.retransmission逻辑)
2. 完整分析代码
def analyze_tcp_stream(stream_pkts): """分析单个TCP流的iRTT、重传率等参数""" # 初始化流的双向数据跟踪 first_pkt = stream_pkts[0] src1, sport1 = first_pkt[IP].src, first_pkt[TCP].sport dst1, dport1 = first_pkt[IP].dst, first_pkt[TCP].dport syn_pkt = None syn_ack_pkt = None seq_tracker = {"dir1": set(), "dir2": set()} # 记录两个方向已发送的序列号 retrans_count = 0 total_pkts = len(stream_pkts) for pkt in stream_pkts: # 区分数据包方向 is_dir1 = (pkt[IP].src == src1 and pkt[TCP].sport == sport1 and pkt[IP].dst == dst1 and pkt[TCP].dport == dport1) tracker_key = "dir1" if is_dir1 else "dir2" # 捕获SYN和SYN-ACK包 if pkt[TCP].flags.S and not pkt[TCP].flags.A: syn_pkt = pkt elif pkt[TCP].flags.S and pkt[TCP].flags.A: syn_ack_pkt = pkt # 计算当前包的序列号范围 seq = pkt[TCP].seq payload_len = len(pkt[TCP].payload) seq_range = range(seq, seq + payload_len) if payload_len > 0 else {seq} # 判断是否为重传包:序列号已被之前的包发送过 if any(s in seq_tracker[tracker_key] for s in seq_range): retrans_count += 1 # 更新已发送的序列号集合 seq_tracker[tracker_key].update(seq_range) # 计算iRTT irt = syn_ack_pkt.time - syn_pkt.time if (syn_pkt and syn_ack_pkt) else None # 计算重传率 retrans_rate = (retrans_count / total_pkts) * 100 if total_pkts > 0 else 0 return { "stream_identifier": f"{src1}:{sport1} <-> {dst1}:{dport1}", "iRTT": round(irt, 6) if irt else None, "retransmission_count": retrans_count, "total_tcp_packets": total_pkts, "retransmission_rate": round(retrans_rate, 2) }
三、调用示例
if __name__ == "__main__": pcap_file = "test.pcapng" tcp_streams = split_tcp_streams(pcap_file) print(f"共识别到 {len(tcp_streams)} 个TCP流\n") for idx, (_, pkts) in enumerate(tcp_streams.items(), 1): analysis_result = analyze_tcp_stream(pkts) print(f"=== 流 {idx} 分析结果 ===") print(f"流标识: {analysis_result['stream_identifier']}") if analysis_result['iRTT']: print(f"初始RTT(iRTT): {analysis_result['iRTT']} 秒") else: print("未找到SYN/SYN-ACK包,无法计算iRTT") print(f"重传包数量: {analysis_result['retransmission_count']}") print(f"总TCP包数量: {analysis_result['total_tcp_packets']}") print(f"重传率: {analysis_result['retransmission_rate']}%\n")
四、内存优化进阶(超大PCAP适用)
如果PCAP文件极大(GB级),可以直接边读包边分析,无需保存整个流的数据包,进一步降低内存占用:
def analyze_streams_online(pcap_path): """在线分析TCP流,不保存完整数据包""" streams = defaultdict(lambda: { "syn_pkt": None, "syn_ack_pkt": None, "seq_tracker": {"dir1": set(), "dir2": set()}, "retrans_count": 0, "total_pkts": 0, "direction": None }) with PcapReader(pcap_path) as reader: for pkt in reader: if IP not in pkt or TCP not in pkt: continue src, sport, dst, dport = pkt[IP].src, pkt[TCP].sport, pkt[IP].dst, pkt[TCP].dport stream_key = tuple(sorted(((src, sport), (dst, dport)))) stream_data = streams[stream_key] # 初始化流方向信息 if not stream_data["direction"]: stream_data["direction"] = (src, sport, dst, dport) src1, sport1, dst1, dport1 = stream_data["direction"] stream_data["total_pkts"] += 1 is_dir1 = (src == src1 and sport == sport1 and dst == dst1 and dport == dport1) tracker_key = "dir1" if is_dir1 else "dir2" # 捕获SYN/SYN-ACK if pkt[TCP].flags.S and not pkt[TCP].flags.A: stream_data["syn_pkt"] = pkt elif pkt[TCP].flags.S and pkt[TCP].flags.A: stream_data["syn_ack_pkt"] = pkt # 检测重传 seq = pkt[TCP].seq payload_len = len(pkt[TCP].payload) seq_range = range(seq, seq + payload_len) if payload_len > 0 else {seq} if any(s in stream_data["seq_tracker"][tracker_key] for s in seq_range): stream_data["retrans_count"] += 1 stream_data["seq_tracker"][tracker_key].update(seq_range) # 整理结果 results = [] for _, data in streams.items(): irt = data["syn_ack_pkt"].time - data["syn_pkt"].time if (data["syn_pkt"] and data["syn_ack_pkt"]) else None retrans_rate = (data["retrans_count"] / data["total_pkts"]) * 100 if data["total_pkts"] > 0 else 0 src1, sport1, dst1, dport1 = data["direction"] results.append({ "stream_identifier": f"{src1}:{sport1} <-> {dst1}:{dport1}", "iRTT": round(irt, 6) if irt else None, "retransmission_count": data["retrans_count"], "total_tcp_packets": data["total_pkts"], "retransmission_rate": round(retrans_rate, 2) }) return results
内容的提问来源于stack exchange,提问作者Mohammad.J
相关产品推荐
相关产品推荐

