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

镜像端口抓包问题: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()

问题分析与解决方案

核心原因

  1. 代码逻辑错误:C的抓包代码中存在未定义变量self.byte_size,实际应使用self.size,导致接收完8字节大小后,后续包的判断逻辑失效,直接跳过处理,表现为丢包。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 14:39:55