镜像端口抓包问题:UDP报文超50kb时丢包(Python+Scapy)
端口镜像UDP抓包丢包问题
问题背景
A、B两台机器通过UDP套接字通信:A发送整数n给B,B返回经pickle序列化的(n×3) numpy数组。为适配大数组传输,B先发送数组序列化后的字节大小(8字节),再分1024字节块发送数据;A接收大小后持续收包直到数据完整,这一过程无丢包。
引入机器C,将B的端口镜像到C后,C用Scapy抓包:识别8字节包为数据大小,后续拼接数据包直到达到指定大小。但数组规模超过2000×3(约48KB)时,C出现丢包,A、B通信仍正常。
关键信息
- A与B通信无丢包
- 所有设备通过千兆以太网连接
- 用VS Code SSH访问所有设备
- 若在B的每个1024字节块发送后加
time.sleep(0.001),C抓大数组无丢包,但传输速度过慢 - 交换机管理日志无错误或丢包记录
- 测试数据:
- 2500×3数组:C收到59包中的57包
- 5000×3数组:C收到116包中的61包
- 10000×3数组:C收到235包中的79包
代码实现
运行在B的服务器代码
import numpy as np import pickle from struct import pack, unpack import socket import time class TestUDPServer(): def __init__(self): self.server_host = "192.168.0.101" self.server_port = 8080 self.start_server() self.main_loop() def receive(self): data, client_address = self.server_socket.recvfrom(1024) data_r = pickle.loads(data) return data_r, client_address def send(self, var, client_address): pickled_data = pickle.dumps(var) # Send the size of the message first size = len(pickled_data) self.server_socket.sendto(pack(">Q", size), client_address) # Send message in chunks chunk_size = 1024 for i in range(0, len(pickled_data), chunk_size): chunk = pickled_data[i:i + chunk_size] self.server_socket.sendto(chunk, client_address) #time.sleep(0.001) def start_server(self): self.server_socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.server_socket.bind((self.server_host, self.server_port)) print(f"UDP Server listening on {self.server_host}:{self.server_port}") def main_loop(self): while True: pc_size, client_address = self.receive() if pc_size>0: break pc = np.random.random((pc_size, 3)) print("received size:", pc_size) self.send(pc, client_address) if __name__ == "__main__": TestUDPServer()
运行在A的客户端代码
import numpy as np import pickle import socket from struct import pack, unpack class TestUDPClient(): def __init__(self): self.server_host = "192.168.0.101" self.server_port = 8080 self.start_client() self.await_trigger() def receive(self): # Receive the size of the message first size_data, _ = self.client_socket.recvfrom(8) size = unpack(">Q", size_data)[0] received_data = b"" chunk_size = 1024 # Receive message in chunks n_chunks = 0 while size > 0: n_chunks += 1 chunk, _ = self.client_socket.recvfrom(min(size, chunk_size)) if not chunk: break received_data += chunk size -= len(chunk) print("n_chunks", n_chunks) data_r = pickle.loads(received_data) return data_r def send(self, var): pickled_data = pickle.dumps(var) self.client_socket.sendto(pickled_data, (self.server_host, self.server_port)) def start_client(self): self.client_socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) print(f"UDP Client connected to {self.server_host}:{self.server_port}") def await_trigger(self): while True: size = int(input("Array length? (n x 3) : ")) # User inputs integer self.send(size) rec = self.receive() print("received shape:", rec.shape) self.client_socket.close() if __name__ == "__main__": TestUDPClient()
运行在C的抓包代码
from scapy.all import sniff, IP, TCP, UDP, ICMP from struct import pack, unpack import pickle import numpy as np import time class PortMirror(): def __init__(self): self.source_ip = '192.168.0.101' self.size = 0 self.n_packets = 0 self.data = b'' self.array = None sniff(iface='enp1s0', prn=self.packet_callback, store=0, filter='udp') def packet_callback(self, packet): if packet.haslayer('IP') and packet.haslayer('UDP'): # FIlter for IP/UDP src_ip = packet['IP'].src dst_ip = packet['IP'].dst src_port = packet['UDP'].sport dst_port = packet['UDP'].dport if src_port!=22 and dst_port!=22 and src_ip==self.source_ip: # Filter SSH and source ip if len(bytes(packet['Raw'].load))==8: self.size = unpack('>Q', bytes(packet['Raw'].load))[0] self.n_packets = 0 self.data = b'' elif len(self.data) < self.byte_size: self.n_packets += 1 print("n_packets", self.n_packets) self.data += bytes(packet['UDP'].payload) if len(self.data) == self.byte_size: self.array = pickle.loads(self.data) print("Success.", self.array.shape) self.n_packets = 0 self.byte_size = 0 if __name__ == "__main__": PortMirror()
问题分析与解决方案
核心原因
- 代码逻辑错误:C的抓包代码中存在未定义变量
self.byte_size,实际应使用self.size,导致接收完8字节大小后,后续包的判断逻辑失效,直接跳过处理,表现为丢包。 - Scapy用户态处理瓶颈:Scapy默认用用户态回调处理数据包,当B高速发送UDP包时,C的内核网卡接收队列会因用户态处理不及时溢出,而A作为合法接收方,应用层持续调用
recvfrom主动取数据,内核队列不会积压。
修复步骤
1. 修正C的代码逻辑错误
将C代码中所有self.byte_size替换为self.size:
# 修正后代码片段 elif len(self.data) < self.size: self.n_packets += 1 self.data += bytes(packet['UDP'].payload) if len(self.data) == self.size: self.array = pickle.loads(self.data) print("Success.", self.array.shape) self.n_packets = 0 self.size = 0
2. 优化Scapy抓包性能
- 缩小过滤范围:将抓包filter改为
udp and src port 8080,只处理B发送的目标UDP包,减少无关包的处理开销。 - 扩大缓冲区:在
sniff函数中增加buffer_size参数,扩大内核接收缓冲区:sniff(iface='enp1s0', prn=self.packet_callback, store=0, filter='udp and src port 8080', buffer_size=2**20) - 降低打印频率:移除每包打印
n_packets的逻辑,仅在数据接收完成后打印结果,减少IO开销。
3. 可选:优化B的发送策略
如果C的性能仍无法满足,可在B的发送逻辑中加入批量延迟,平衡传输速度与C的处理能力:
# 修改B的send函数循环 chunk_size = 1024 batch_size = 10 # 每发10个包延迟一次 for i in range(0, len(pickled_data), chunk_size): chunk = pickled_data[i:i + chunk_size] self.server_socket.sendto(chunk, client_address) if (i // chunk_size) % batch_size == batch_size - 1: time.sleep(0.0001) # 100微秒延迟,几乎不影响传输速度
内容的提问来源于stack exchange,提问作者wittyUsername
相关产品推荐
相关产品推荐

