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

Go UDP服务器调整读取缓冲区无效,仅处理少量数据包问题

问题描述

我编写了一个简单的UDP服务器,负责监听数据包、模拟业务操作等待50毫秒后将消息打印到终端,服务器代码如下:

package main

import (
    "fmt"
    "log"
    "log/slog"
    "net"
    "time"
)

func main() {
    udpAddr, err := net.ResolveUDPAddr("udp", "0.0.0.0:8080")
    if err != nil {
        log.Fatal(err)
    }

    conn, err := net.ListenUDP("udp", udpAddr)
    if err != nil {
        log.Fatal(err)
    }

    slog.Info("UDP server listening", "addr", udpAddr, "workers", numCpu)

    go startWorker(0, conn)

    // Block forever
    <-make(chan struct{})
}

func startWorker(id int, conn *net.UDPConn) {
    slog.Info("Starting worker", "id", id)

    i := 0
    for {
        buf := make([]byte, 1024)
        _, _, err := conn.ReadFromUDP(buf)
        if err != nil {
            slog.Error("failed to read UDP message", "err", err, "worker", id)
            continue
        }

        time.Sleep(50 * time.Millisecond)

        fmt.Print(i, " > ", string(buf))
        i += 1
    }
}

同时编写了一个Go客户端代码,仅向该服务器发送一次字符串消息:

package main

import (
    "log"
    "net"
)

func main() {
    udpAddr, err := net.ResolveUDPAddr("udp", "0.0.0.0:8080")
    if err != nil {
        log.Fatal(err)
    }

    conn, err := net.DialUDP("udp", nil, udpAddr)
    if err != nil {
        log.Fatal(err)
    }

    // Send a message to the server
    _, err = conn.Write([]byte("Hello UDP Server\n"))
    if err != nil {
        log.Fatal(err)
    }
}

为测试服务器的数据包处理能力,我在终端通过以下命令顺序运行客户端代码10000次:

time (for i in $(seq 10000); do ./client; done)

测试结果显示仅约400-500个数据包被处理,无错误日志,耗时如下:

real    0m12,538s
user    0m10,667s
sys     0m4,014s

我推测这是操作系统缓冲区限制导致的问题,因此在服务器代码中添加了设置1GB读取缓冲区的代码:

// 1 GB
if err = conn.SetReadBuffer(1024 * 1024 * 1024); err != nil {
    log.Fatal(err)
}

再次测试后,处理的数据包数量仅提升至600-700个。显然1GB缓冲区足以容纳10000条简单消息,请问该现象的原因是什么?如何解决?

注:我使用的是Ubuntu 22.04系统。


原因分析与解决方案

核心原因

  1. 单worker串行处理瓶颈:服务器仅启动1个worker goroutine,每次处理数据包时,完成ReadFromUDP后会阻塞50ms模拟业务逻辑。单worker每秒最多处理20个请求(1000ms/50ms),12秒左右最多处理240个;实际处理量略高是因为UDP缓冲区暂时堆积了部分数据包,但当缓冲区满后,后续数据包会被操作系统直接丢弃。
  2. 系统缓冲区上限限制:你设置的1GB缓冲区看似很大,但Ubuntu默认的net.core.rmem_max(全局UDP读取缓冲区上限)远小于1GB(默认通常为212992字节),SetReadBuffer的实际生效值不会超过这个系统限制,所以缓冲区还是会快速被填满导致丢包。
  3. UDP特性与客户端行为:UDP是无连接、不可靠协议,丢包不会有任何通知;同时客户端每次启动都要完成socket创建、发送、退出流程,大量客户端快速发送数据时,服务器单worker来不及处理,加剧了缓冲区溢出丢包的情况。

解决方案

1. 调整系统UDP缓冲区参数

先修改Ubuntu系统的UDP缓冲区上限,让SetReadBuffer的配置生效:

  • 临时生效(重启后失效):
    sudo sysctl -w net.core.rmem_max=1073741824
    sudo sysctl -w net.core.rmem_default=1073741824
    
  • 永久生效:
    编辑/etc/sysctl.conf,添加或修改以下内容:
    net.core.rmem_max=1073741824
    net.core.rmem_default=1073741824
    
    执行sudo sysctl -p让配置立即生效。

2. 多worker并行处理

将单worker改为多goroutine处理,利用CPU并行能力,避免业务阻塞导致的缓冲区堆积。修改服务器代码如下:

package main

import (
    "fmt"
    "log"
    "log/slog"
    "net"
    "runtime"
    "time"
)

func main() {
    numCpu := runtime.NumCPU()
    udpAddr, err := net.ResolveUDPAddr("udp", "0.0.0.0:8080")
    if err != nil {
        log.Fatal(err)
    }

    conn, err := net.ListenUDP("udp", udpAddr)
    if err != nil {
        log.Fatal(err)
    }

    // 设置读取缓冲区
    if err = conn.SetReadBuffer(1024 * 1024 * 1024); err != nil {
        log.Fatal(err)
    }

    slog.Info("UDP server listening", "addr", udpAddr, "workers", numCpu)

    // 启动与CPU核心数匹配的worker
    for i := 0; i < numCpu; i++ {
        go startWorker(i, conn)
    }

    // Block forever
    <-make(chan struct{})
}

func startWorker(id int, conn *net.UDPConn) {
    slog.Info("Starting worker", "id", id)

    i := 0
    for {
        buf := make([]byte, 1024)
        _, _, err := conn.ReadFromUDP(buf)
        if err != nil {
            slog.Error("failed to read UDP message", "err", err, "worker", id)
            continue
        }

        // 将业务逻辑放到独立goroutine,避免阻塞Read操作
        go func(msg []byte, count int) {
            time.Sleep(50 * time.Millisecond)
            fmt.Print(count, " > ", string(msg))
        }(buf, i)

        i += 1
    }
}

3. 优化客户端测试逻辑(可选)

避免每次启动新客户端进程,改为在单个进程内循环发送数据,减少进程启动开销:

package main

import (
    "fmt"
    "log"
    "net"
    "time"
)

func main() {
    udpAddr, err := net.ResolveUDPAddr("udp", "0.0.0.0:8080")
    if err != nil {
        log.Fatal(err)
    }

    conn, err := net.DialUDP("udp", nil, udpAddr)
    if err != nil {
        log.Fatal(err)
    }
    defer conn.Close()

    // 循环发送10000次消息
    for i := 0; i < 10000; i++ {
        _, err = conn.Write([]byte(fmt.Sprintf("Hello UDP Server #%d\n", i)))
        if err != nil {
            log.Fatal(err)
        }
        // 添加微小延迟,避免瞬间发送导致客户端自身缓冲区溢出
        time.Sleep(1 * time.Millisecond)
    }
}

4. 效果验证

调整后重新运行测试,服务器的处理能力会大幅提升,基本可以处理全部10000个数据包,耗时也会根据CPU核心数显著降低。


内容的提问来源于stack exchange,提问作者Ahmet Yazıcı

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 21:13:21