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

使用eBPF/BCC采集NFS流量统计数据结果不准确问题求助

NFS流量统计程序与nethogs数据不符的问题分析与优化

问题描述

我正在研究eBPF,编写了一个用于统计NFS流量的程序,该程序可展示IP地址及对应收发数据量。目前程序可正常运行,但统计结果与nethogs工具的流量数据不符:我的程序统计的接收速率为nethogs的2倍,发送速率比nethogs统计值高20%-25%。相关代码如下:

#!/usr/bin/env python3

from bcc import BPF
import time
import socket
import struct

bpf_text = """
#include <uapi/linux/ptrace.h>
#include <linux/sched.h>
#include <linux/socket.h>
#include <linux/in.h>
#include <linux/tcp.h>
#include <linux/bpf.h>

#define AF_INET 2
#define NFS_PORT 2049

struct traffic {
    u64 sent;
    u64 recv;
};

struct sock_storage {
    struct sock *sk;
};

BPF_HASH(nfs_traffic, u32, struct traffic, 1024);           // key = server IPv4 (u32, network order)
BPF_HASH(call_args, u64, struct sock_storage, 4096);        // key = pid_tgid, temp storage for sk

static inline int is_nfs_tcp(struct sock *sk) {
    u16 family = 0;
    u16 dport = 0;
    u16 sport = 0;

    bpf_probe_read_kernel(&family, sizeof(family), &sk->__sk_common.skc_family);
    if (family != AF_INET)
        return 0;

    bpf_probe_read_kernel(&dport, sizeof(dport), &sk->__sk_common.skc_dport);
    bpf_probe_read_kernel(&sport, sizeof(sport), &sk->__sk_common.skc_num);
    if (ntohs(dport) != NFS_PORT && sport != NFS_PORT)
        return 0;

    return 1;
}

static inline u32 get_server_ip(struct sock *sk) {
    u32 daddr = 0;
    bpf_probe_read_kernel(&daddr, sizeof(daddr), &sk->__sk_common.skc_daddr);
    return daddr;
}

static inline void update_counter(u32 key, u64 bytes, int is_sent) {
    struct traffic *t = nfs_traffic.lookup(&key);

    if (!t) {
        struct traffic zero = {0};
        nfs_traffic.update(&key, &zero);
        t = nfs_traffic.lookup(&key);
        if (!t)
            return;
    }

    if (is_sent) {
        __sync_fetch_and_add(&t->sent, bytes);
        bpf_trace_printk ("Yo sent %llu", bytes);
    }
    else {
        __sync_fetch_and_add(&t->recv, bytes);
        bpf_trace_printk ("Yo recv %llu", bytes);
    }
}

int kprobe__tcp_sendmsg(struct pt_regs *ctx, struct sock *sk) {
    u64 pid_tgid = bpf_get_current_pid_tgid();

    struct sock_storage val = {};
    val.sk = sk;

    call_args.update(&pid_tgid, &val);
    return 0;
}

int kretprobe__tcp_sendmsg(struct pt_regs *ctx) {
    u64 pid_tgid = bpf_get_current_pid_tgid();
    struct sock_storage *val = call_args.lookup(&pid_tgid);
    if (!val)
        return 0;

    struct sock *sk = val->sk;
    s64 ret = PT_REGS_RC(ctx);   // signed – errors are negative

    call_args.delete(&pid_tgid);

    if (ret <= 0)
        return 0;

    if (!is_nfs_tcp(sk))
        return 0;

    u32 key = get_server_ip(sk);
    update_counter(key, (u64)ret, 1);   // sent = 1

    return 0;
}

int kprobe__tcp_cleanup_rbuf(struct pt_regs *ctx, struct sock *sk, int copied) {
    if (copied <= 0)
        return 0;

    if (!is_nfs_tcp(sk))
        return 0;

    u32 key = get_server_ip(sk);
    update_counter(key, (u64)copied, 0);  // recv

    return 0;
}
"""

# Load program
b = BPF(text=bpf_text)

# Attach probes
b.attach_kprobe(event="tcp_sendmsg", fn_name="kprobe__tcp_sendmsg")
b.attach_kretprobe(event="tcp_sendmsg", fn_name="kretprobe__tcp_sendmsg")

b.attach_kprobe(event="tcp_cleanup_rbuf", fn_name="kprobe__tcp_cleanup_rbuf")

print("Press Ctrl+C to stop\n")

def human(b):
    if b >= 1024*1024:
        return f"{b/(1024*1024):.2f} MB/s"
    if b >= 1024:
        return f"{b/1024:.1f} KB/s"
    return f"{b} B/s"

try:
    while True:
        time.sleep(1.0)

        table = b["nfs_traffic"]
        items = table.items()

        if not items:
            continue

        ts = time.strftime("%H:%M:%S")
        print(f"[{ts}]")

        for k, v in items:
            ip = socket.inet_ntoa(struct.pack("I", k.value))
            sent = v.sent
            recv = v.recv

            print ("sent: ", sent, " recv: ", recv, "\n");

            if sent == 0 and recv == 0:
                continue

            print(f"  {ip:15}   sent {human(sent):>10}   recv {human(recv):>10}")

        table.clear()

except KeyboardInterrupt:
    print("\nStopped.")
finally:
    b.cleanup()

问题分析与优化方案

1. 发送流量偏高(+20%-25%):统计了TCP重传数据

  • 核心原因:tcp_sendmsg的返回值是内核成功放入发送队列的字节数,包含应用层首次发送的数据和TCP协议层自动重传的数据。而nethogs统计的是应用层实际发起的发送请求量,不包含重传的冗余数据,因此统计值会偏高。
  • 优化方案:
    • 改用TCP tracepoint:tracepoint:tcp:tcp_sendmsg提供了data_len参数,直接对应应用层发送的字节数,不包含重传。修改BPF代码中的探针点,替换kprobe/kretprobe为tracepoint;
    • 示例代码片段(替换原发送相关探针):
      TRACEPOINT_PROBE(tcp, tcp_sendmsg) {
          struct sock *sk = args->sk;
          u64 data_len = args->data_len;
      
          if (data_len <= 0)
              return 0;
          if (!is_nfs_tcp(sk))
              return 0;
      
          u32 key = get_server_ip(sk);
          update_counter(key, data_len, 1);
          return 0;
      }
      
    • Python侧无需再attach kprobe/kretprobe,改为加载tracepoint即可。

2. 接收流量翻倍:方向判断错误或探针点选择不当

  • 核心原因:
    • 方向逻辑模糊:原is_nfs_tcp函数只要源端口或目标端口是2049就统计,导致本地作为NFS服务器时,接收客户端请求的流量和服务器发送响应的流量都被计入同一统计(如果IP判断逻辑混淆);
    • 探针点不准确:tcp_cleanup_rbuf用于清理接收缓冲区,某些场景下(如NFS缓存)可能被多次触发,导致同一数据被重复统计;该函数的copied参数也可能包含非应用层数据。
  • 优化方案:
    • 明确流量方向:如果仅作为NFS客户端统计与服务器的流量,修改is_nfs_tcp,只保留目标端口为2049(发送请求)或源端口为2049(接收响应)的判断,避免无关流量被统计:
      static inline int is_nfs_client_tcp(struct sock *sk) {
          u16 family = 0;
          u16 dport = 0;
      
          bpf_probe_read_kernel(&family, sizeof(family), &sk->__sk_common.skc_family);
          if (family != AF_INET)
              return 0;
      
          bpf_probe_read_kernel(&dport, sizeof(dport), &sk->__sk_common.skc_dport);
          // 仅统计客户端向NFS服务器发送的请求(目标端口2049)
          if (ntohs(dport) != NFS_PORT)
              return 0;
      
          return 1;
      }
      
    • 更换接收探针点:改用tracepoint:tcp:tcp_recvmsg或kretprobe__tcp_recvmsg,获取应用层实际读取的字节数。示例(kretprobe版本):
      int kretprobe__tcp_recvmsg(struct pt_regs *ctx) {
          struct sock *sk = (struct sock *)PT_REGS_PARM1(ctx);
          s64 ret = PT_REGS_RC(ctx);
      
          if (ret <= 0)
              return 0;
          if (!is_nfs_tcp(sk))
              return 0;
      
          u32 key = get_server_ip(sk);
          update_counter(key, (u64)ret, 0);
          return 0;
      }
      

3. 其他潜在优化点

  • 并发安全优化:创建traffic结构体时,原代码的update+lookup存在竞态,可改用lookup_or_try_init简化逻辑:
    static inline void update_counter(u32 key, u64 bytes, int is_sent) {
        struct traffic *t = nfs_traffic.lookup_or_try_init(&key, &(struct traffic){0});
        if (!t)
            return;
    
        if (is_sent) {
            __sync_fetch_and_add(&t->sent, bytes);
        } else {
            __sync_fetch_and_add(&t->recv, bytes);
        }
    }
    
  • 流量统计逻辑优化:原代码每次清空哈希表,可能遗漏未处理的流量。改为维护上一次的统计值,计算每秒增量:
    • 在Python侧保存上一轮的sent/recv数据,当前值减去上一轮值即为每秒速率,避免清空哈希表。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 13:35:54