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

Go程序随Go协程/核心增加性能不升反降求助

多CPU服务器下Go协程处理性能下降的排查与优化方案

问题描述

我编写了一段数据处理代码,根据核心数启动n个Go协程处理任务,流程包含文件下载、正则匹配提取数据并写入CSV。在配备4个vCPU的小型服务器上,处理耗时约9分钟;但在配备192个vCPU、主频相近的AWS大型服务器上,仅处理环节(不含文件下载)耗时增至20分钟。恳请协助排查原因。

代码片段

package main

import (
    "encoding/csv"
    "errors"
    "fmt"
    "io"
    "io/ioutil"
    "log"
    "net/http"
    "os"
    "regexp"
    "strings"
    "time"

    "github.com/slyrz/warc"
)

var retrier = retry.NewRetrier(10, 100*time.Millisecond, 10*time.Second)

var regex = `regex`
var re = regexp.MustCompile(regex)

func extractData(input string) []string {
    matches := re.FindAllString(input, -1)
    return matches
}

func DownloadFile(url, filename string) error {
    err := retrier.Run(func() error {
        out, err := os.Create(filename)
        if err != nil {
            return fmt.Errorf("error creating file: %s", err)
        }

        defer out.Close()

        fmt.Println("Started downloading:", filename)

        resp, err := http.Get(url)
        if err != nil {
            os.Remove(filename)
            return fmt.Errorf("error fetching file: %s %s", filename, err)
        }

        defer resp.Body.Close()

        _, err = io.Copy(out, resp.Body)
        if err != nil {
            os.Remove(filename)
            return fmt.Errorf("error copying file to disk: %s %s", filename, err)
        }

        fmt.Println("Finished downloading:", filename)
        return nil
    })

    if err != nil {
        return err
    }

    return nil
}

const baseURL = "URL"

func GenerateOutname(input string) string {
    splitted := strings.Split(input, ".")
    return splitted[0]
}

func ProcessFile(filename string) {
    start := time.Now()
    count := 0

    f, err := os.Open(filename)
    if err != nil {
        panic(err)
    }

    defer f.Close()

    reader, err := warc.NewReader(f)
    if err != nil {
        panic(err)
    }

    defer reader.Close()

    output := GenerateOutname(filename)
    csvFile, err := os.Create("data/" + output + ".csv")

    if err != nil {
        fmt.Printf("failed creating file: %s", err)
    }

    defer csvFile.Close()

    csvwriter := csv.NewWriter(csvFile)

    for {
        record, err := reader.ReadRecord()
        if err != nil {
            if errors.Is(err, io.EOF) {
                break
            }
            fmt.Println("Error reading record", err)
            break
        }

        if record.Header.Get("warc-type") != "response" {
            continue
        }

        uri := record.Header.Get("WARC-Target-URI")
        date := record.Header.Get("WARC-Date")
        data, err := ioutil.ReadAll(record.Content)
        if err != nil {
            fmt.Println("Error reading from the reader:", err)
            return
        }

        resultString := string(data)
        matches := extractData(resultString)

        for _, data := range matches {
            err := csvwriter.Write([]string{date, data, uri})
            if err != nil {
                fmt.Println("error writing to csv", err)
            }
            count++
        }
    }

    csvwriter.Flush()

    elapsed := time.Since(start)

    fmt.Printf("%s - Took: %s\n", filename, elapsed)
    fmt.Printf("%s - Data: %d\n", filename, count)
}

func worker(id int, jobs <-chan string, results chan<- int) {
    for job := range jobs {
        remoteFile := job
        url := baseURL + remoteFile

        splitted := strings.Split(remoteFile, "/")
        localFile := splitted[len(splitted)-1]

        fmt.Printf("Started Job: Worker - %d Job - %s\n", id, localFile)

        err := DownloadFile(url, localFile)
        if err != nil {
            fmt.Printf("Failed to download: %s\n", localFile)
            results <- 1
            continue
        }

        ProcessFile(localFile)
        os.Remove(localFile)

        fmt.Printf("Completed Job: Worker - %d Job - %s\n", id, localFile)
        results <- 1
    }
}

func ReadFiles() []string {
    file, err := os.Open("special.txt")

    if err != nil {
        log.Fatal("Error opening paths", err)
    }

    defer file.Close()

    reader := csv.NewReader(file)
    records, err := reader.ReadAll()

    if err != nil {
        fmt.Println("Error reading paths", err)
    }

    paths := []string{}

    for _, record := range records {
        paths = append(paths, record[0])
    }
    fmt.Println("Total paths:", len(paths))
    return paths
}

func main() {
    files := ReadFiles()
    paths := files

    const workers = 4 // depends on vCPUs
    jobs := make(chan string, len(paths))
    results := make(chan int, len(paths))

    for w := 1; w <= workers; w++ {
        go worker(w, jobs, results)
    }

    for _, path := range paths {
        jobs <- path
    }
    close(jobs)

    for a := 1; a <= len(paths); a++ {
        <-results
    }
}

核心排查点

  • 协程数硬编码,未利用多核资源
    代码中const workers = 4固定为4个协程,192核服务器的CPU资源完全闲置,反而因为单进程少量协程在多核环境下出现调度开销、缓存命中率下降等问题,导致处理速度变慢。

  • 磁盘IO成为瓶颈
    AWS大型服务器的存储如果是低IOPS的共享卷(如EBS gp2),当多个协程同时进行文件下载、读取、CSV写入操作时,IO等待时间会显著增加。此外,所有CSV文件都写入data/目录,可能引发文件系统锁竞争,进一步拖慢处理速度。

  • 内存与GC开销
    ProcessFile中用ioutil.ReadAll(record.Content)将整个WARC记录读入内存,若记录体积较大,会导致内存占用飙升,触发频繁的GC。在多核服务器上,GC的STW(Stop The World)开销会被放大,影响整体性能。

  • 正则匹配的潜在性能问题
    全局正则对象re虽然并发安全,但如果正则表达式存在高回溯逻辑,在多核环境下会因CPU缓存颠簸导致匹配效率下降。

优化方案

  1. 动态设置协程数
    导入runtime包,根据服务器核心数动态调整协程数量:

    func main() {
        files := ReadFiles()
        paths := files
    
        workers := runtime.NumCPU() // CPU密集型任务建议用核心数,IO密集型可乘2
        jobs := make(chan string, len(paths))
        results := make(chan int, len(paths))
    
        for w := 1; w <= workers; w++ {
            go worker(w, jobs, results)
        }
        // ... 剩余代码不变
    }
    
  2. 优化磁盘IO

    • 将临时下载文件放到内存文件系统(如/tmp,AWS服务器默认配置tmpfs),减少磁盘读写:
      localFile := filepath.Join(os.TempDir(), splitted[len(splitted)-1])
      
    • 使用高性能存储卷(如IO2/IO3)存放data/目录,提升IOPS;按文件前缀拆分CSV存储目录,避免单目录文件过多引发竞争。
  3. 减少内存占用
    避免一次性读取整个WARC记录,直接流式处理record.Content:

    // 替换ioutil.ReadAll,直接从reader匹配正则
    matches := re.FindAllReader(record.Content, -1)
    // 注意:需要自行处理reader到字符串的转换,或使用支持流式匹配的方法
    
  4. 正则表达式优化
    用go tool pprof分析正则匹配的耗时热点,简化正则逻辑,减少回溯(比如避免贪婪匹配、减少不必要的捕获组)。

  5. 性能分析验证

    • 生成CPU分析报告:
      go run -cpuprofile cpu.prof your_program.go
      go tool pprof cpu.prof
      
    • 生成内存分析报告:
      go run -memprofile mem.prof your_program.go
      go tool pprof mem.prof
      
    • 用htop或iostat监控服务器的CPU、内存、磁盘IO使用率,确认瓶颈所在。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 03:35:57