Learn
Go/18-project-concurrent-file-processor

项目:并发文件处理工具

本节用 Worker Pool 实现一个真实可用的小工具:扫描目录、统计每个文件的行数、找出 Top-N 最大文件。这是 Go 并发编程的经典实战。

1. 项目结构

fileinfo/
├── main.go      # 入口:扫描 → 派发 → 收集
├── scanner.go   # 扫描目录得到文件列表
├── counter.go   # 统计行数(耗时操作)
├── topn.go      # 用最小堆找 Top-N
└── main_test.go # 测试

2. scanner.go:扫描目录

package main
 
import (
    "os"
    "path/filepath"
    "strings"
)
 
// FileInfo 是 worker 处理的基本单位
type FileInfo struct {
    Path string
    Size int64
}
 
// ScanDir 递归扫描 root,返回所有文件的 FileInfo
// extensions: 过滤后缀(如 []string{".go", ".md"});为空表示不过滤
func ScanDir(root string, extensions []string) ([]FileInfo, error) {
    var out []FileInfo
    err := filepath.Walk(root, func(path string, info os.FileInfo, err error) error {
        if err != nil {
            return err
        }
        if info.IsDir() {
            return nil
        }
        if len(extensions) > 0 {
            ext := strings.ToLower(filepath.Ext(path))
            ok := false
            for _, e := range extensions {
                if e == ext {
                    ok = true
                    break
                }
            }
            if !ok {
                return nil
            }
        }
        out = append(out, FileInfo{Path: path, Size: info.Size()})
        return nil
    })
    return out, err
}
💡filepath.Walk 的回调

回调函数对每个文件/目录执行一次。返回 nil 继续遍历,返回 error 终止并把错误传给 Walk。

3. counter.go:行数统计

package main
 
import (
    "bufio"
    "os"
)
 
// LineCount 统计文件行数
// 大文件用 bufio.Scanner 一次读一行
func LineCount(path string) (int, error) {
    f, err := os.Open(path)
    if err != nil {
        return 0, err
    }
    defer f.Close()
 
    sc := bufio.NewScanner(f)
    // 单行最长 1MB(默认 64KB 太小)
    sc.Buffer(make([]byte, 64*1024), 1024*1024)
    n := 0
    for sc.Scan() {
        n++
    }
    return n, sc.Err()
}

4. topn.go:用最小堆找 Top-N

Go 标准库 container/heap 提供堆接口。Top-N 用最小堆最自然:堆顶是当前第 N 大,进来更小的就丢,进来更大的就替换。

package main
 
import "container/heap"
 
// FileStat 包含行数信息
type FileStat struct {
    Path  string
    Size  int64
    Lines int
}
 
// 最小堆:按 Lines 排序
type minHeap []FileStat
func (h minHeap) Len() int            { return len(h) }
func (h minHeap) Less(i, j int) bool  { return h[i].Lines < h[j].Lines }
func (h minHeap) Swap(i, j int)       { h[i], h[j] = h[j], h[i] }
func (h *minHeap) Push(x any)         { *h = append(*h, x.(FileStat)) }
func (h *minHeap) Pop() any {
    old := *h
    n := len(old)
    x := old[n-1]
    *h = old[:n-1]
    return x
}
 
// TopN 维护一个大小 <= n 的最小堆
type TopN struct {
    n   int
    h   minHeap
}
 
func NewTopN(n int) *TopN { return &TopN{n: n} }
 
// Add 尝试加入一条数据
func (t *TopN) Add(s FileStat) {
    if t.h.Len() < t.n {
        heap.Push(&t.h, s)
        return
    }
    if s.Lines > t.h[0].Lines {
        heap.Pop(&t.h)
        heap.Push(&t.h, s)
    }
}
 
// Result 按 Lines 降序返回
func (t *TopN) Result() []FileStat {
    out := make([]FileStat, len(t.h))
    copy(out, t.h)
    sort.Slice(out, func(i, j int) bool {
        return out[i].Lines > out[j].Lines
    })
    return out
}
ℹ️为什么最小堆
  • 完全排序:O(N log N)
  • Top-N:O(N log k),k 是 N 大小
  • 当 N = 100 万、k = 10 时,性能差 1000 倍

5. main.go:Worker Pool 串联三个模块

package main
 
import (
    "flag"
    "fmt"
    "log"
    "sync"
)
 
func main() {
    var (
        root = flag.String("root", ".", "要扫描的目录")
        ext  = flag.String("ext", ".go", "文件后缀,如 .go / .md")
        n    = flag.Int("top", 10, "Top-N 数量")
        work = flag.Int("workers", 4, "worker 数量")
    )
    flag.Parse()
 
    // 1) 扫描
    files, err := ScanDir(*root, []string{*ext})
    if err != nil {
        log.Fatal(err)
    }
    fmt.Printf("扫描到 %d 个 %s 文件\\n", len(files), *ext)
 
    // 2) worker pool 统计行数
    jobs := make(chan FileInfo)
    results := make(chan FileStat)
    var wg sync.WaitGroup
    for i := 0; i < *work; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for fi := range jobs {
                lines, err := LineCount(fi.Path)
                if err != nil {
                    log.Printf("错误 %s: %v", fi.Path, err)
                    continue
                }
                results <- FileStat{Path: fi.Path, Size: fi.Size(), Lines: lines}
            }
        }()
    }
 
    // 3) 派发
    go func() {
        for _, fi := range files {
            jobs <- fi
        }
        close(jobs)
    }()
 
    // 4) 收集到 TopN
    go func() {
        wg.Wait()
        close(results)
    }()
 
    top := NewTopN(*n)
    for stat := range results {
        top.Add(stat)
    }
 
    // 5) 打印结果
    fmt.Printf("\\nTop %d 行数最多:\\n", *n)
    for i, s := range top.Result() {
        fmt.Printf("  %2d. %5d 行  %s\\n", i+1, s.Lines, s.Path)
    }
}
💡三对 channel 配对
  • jobs:派发器 → worker
  • results:worker → 收集器
  • close(jobs) 由派发方调用;close(results) 由收集方在 wg.Wait 之后调用
  • 严格遵守:发送方关,接收方只读

6. 单文件可运行版(Playground)

沙盒不能读真实文件系统,我们用**内存中的"文件"**模拟整个流程。

单文件:Top-N 行数
package main
 
import (
    "container/heap"
    "fmt"
    "sort"
    "strings"
    "sync"
)
 
// ===== 内存"文件系统" =====
var fakeFiles = map[string]string{
    "/src/server.go":   strings.Repeat("line\\n", 850) + "end",
    "/src/handler.go":  strings.Repeat("// handler\\n", 120) + "done",
    "/src/db.go":       strings.Repeat("query\\n", 42),
    "/src/util.go":     strings.Repeat("util\\n", 8),
    "/src/main.go":     strings.Repeat("init\\n", 350),
    "/src/router.go":   strings.Repeat("route\\n", 67),
    "/src/middleware.go": strings.Repeat("mw\\n", 200),
}
 
// ScanDir 直接遍历 fakeFiles
func ScanDir() []string {
    paths := make([]string, 0, len(fakeFiles))
    for p := range fakeFiles { paths = append(paths, p) }
    sort.Strings(paths)
    return paths
}
 
// LineCount 数 \\n 个数
func LineCount(path string) int {
    return strings.Count(fakeFiles[path], "\\n")
}
 
// ===== Top-N 最小堆 =====
type FileStat struct {
    Path  string
    Lines int
}
type minHeap []FileStat
func (h minHeap) Len() int            { return len(h) }
func (h minHeap) Less(i, j int) bool  { return h[i].Lines < h[j].Lines }
func (h minHeap) Swap(i, j int)       { h[i], h[j] = h[j], h[i] }
func (h *minHeap) Push(x any)         { *h = append(*h, x.(FileStat)) }
func (h *minHeap) Pop() any {
    old := *h; n := len(old)
    x := old[n-1]; *h = old[:n-1]; return x
}
type TopN struct{ n int; h minHeap }
func NewTopN(n int) *TopN { return &TopN{n: n} }
func (t *TopN) Add(s FileStat) {
    if t.h.Len() < t.n { heap.Push(&t.h, s); return }
    if s.Lines > t.h[0].Lines {
        heap.Pop(&t.h); heap.Push(&t.h, s)
    }
}
func (t *TopN) Result() []FileStat {
    out := make([]FileStat, len(t.h))
    copy(out, t.h)
    sort.Slice(out, func(i, j int) bool { return out[i].Lines > out[j].Lines })
    return out
}
 
// ===== 演示 =====
func main() {
    const workers = 3
    const topK = 3
 
    paths := ScanDir()
    fmt.Printf("扫描到 %d 个文件\\n", len(paths))
    for _, p := range paths {
        fmt.Printf("  %s\\n", p)
    }
 
    jobs := make(chan string, len(paths))
    results := make(chan FileStat, len(paths))
 
    // 启动 workers
    var wg sync.WaitGroup
    for w := 1; w <= workers; w++ {
        wg.Add(1)
        go func(id int) {
            defer wg.Done()
            for p := range jobs {
                fmt.Printf("  [worker %d] 统计 %s\\n", id, p)
                results <- FileStat{Path: p, Lines: LineCount(p)}
            }
        }(w)
    }
 
    // 派发
    for _, p := range paths { jobs <- p }
    close(jobs)
    go func() { wg.Wait(); close(results) }()
 
    // 收集到 TopN
    top := NewTopN(topK)
    for s := range results { top.Add(s) }
 
    fmt.Printf("\\nTop %d 行数最多:\\n", topK)
    for i, s := range top.Result() {
        fmt.Printf("  %d. %4d 行  %s\\n", i+1, s.Lines, s.Path)
    }
}
💡输出顺序
  • 文件顺序:map 遍历是随机的——ScanDir 用 sort.Strings 稳定了
  • worker 打印:取决于调度——展示的顺序仅是参考
  • Top-N 排序:明确按 Lines 降序

7. 性能分析:顺序 vs 并发

顺序 vs 并发 性能对比
package main
 
import (
    "fmt"
    "strings"
    "time"
)
 
var fakeFiles = map[string]string{
    "a.go": strings.Repeat("x\\n", 2000),
    "b.go": strings.Repeat("x\\n", 1500),
    "c.go": strings.Repeat("x\\n", 1800),
    "d.go": strings.Repeat("x\\n", 1200),
    "e.go": strings.Repeat("x\\n", 2500),
    "f.go": strings.Repeat("x\\n", 900),
}
 
func count(p string) int { return strings.Count(fakeFiles[p], "\\n") }
 
func runSequential() int {
    total := 0
    for p := range fakeFiles {
        total += count(p)
    }
    return total
}
 
func runParallel() int {
    ch := make(chan int, len(fakeFiles))
    for p := range fakeFiles {
        p := p
        go func() {
            ch <- count(p)
        }()
    }
    total := 0
    for i := 0; i < len(fakeFiles); i++ {
        total += <-ch
    }
    return total
}
 
func bench(name string, fn func() int) time.Duration {
    start := time.Now()
    got := fn()
    dur := time.Since(start)
    fmt.Printf("  %-12s -> total=%d  (%v)\\n", name, got, dur)
    return dur
}
 
func main() {
    // 多次跑取稳定值
    var seqTotal, parTotal time.Duration
    for i := 0; i < 3; i++ {
        seqTotal += bench("Sequential", runSequential)
        parTotal += bench("Parallel",   runParallel)
        fmt.Println()
    }
    fmt.Printf("平均:Sequential=%v  Parallel=%v\\n",
        seqTotal/3, parTotal/3)
}
⚠️并发不是免费的
  • 数据规模小时,并发反而慢(goroutine 启动有开销)
  • CPU 密集型任务:workers ≤ CPU 核数
  • IO 密集型任务(读文件、HTTP):workers 可以远超 CPU 核数

8. 真实部署的注意事项

// 实际项目里 main 应该有信号处理
func main() {
    ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt)
    defer cancel()
    // ... 业务逻辑
    select {
    case <-ctx.Done():
        fmt.Println("收到 Ctrl+C,正在退出…")
    }
}
💡生产检查清单
  • ✅ flag 让工具可配置(路径、并发数、Top-N)
  • ✅ log.Printf 报错到 stderr,stdout 给结果(方便管道)
  • ✅ 信号监听支持优雅退出
  • ✅ 用 errgroup 替代裸 WaitGroup(自动传播错误)
  • ✅ pprof 暴露 /debug/pprof/,可分析性能

🎯 练习

Top-N 最大文件
// 任务:基于上面的 TopN,实现"找 Top-N 最大的文件(按 size)"
// 1. FileStat 加 Size 字段
// 2. minHeap 改为按 Size 排序
// 3. 演示 fakeFiles 改成 [{"path": "a", "size": 1024}, ...] 形式
// 4. 打印 Top 3 最大文件
 
package main
 
import (
    "container/heap"
    "fmt"
    "sort"
)
 
type FileStat struct {
    Path string
    Size int
}
type maxHeap []FileStat // 改最大堆:堆顶是当前最大
func (h maxHeap) Len() int            { return len(h) }
func (h maxHeap) Less(i, j int) bool  { return h[i].Size > h[j].Size } // 反转
func (h maxHeap) Swap(i, j int)       { h[i], h[j] = h[j], h[i] }
func (h *maxHeap) Push(x any)         { *h = append(*h, x.(FileStat)) }
func (h *maxHeap) Pop() any {
    old := *h; n := len(old); x := old[n-1]; *h = old[:n-1]; return x
}
type TopN struct{ n int; h maxHeap }
func NewTopN(n int) *TopN { return &TopN{n: n} }
func (t *TopN) Add(s FileStat) {
    if t.h.Len() < t.n { heap.Push(&t.h, s); return }
    if s.Size > t.h[0].Size {
        heap.Pop(&t.h); heap.Push(&t.h, s)
    }
}
func (t *TopN) Result() []FileStat {
    out := make([]FileStat, len(t.h))
    copy(out, t.h)
    sort.Slice(out, func(i, j int) bool { return out[i].Size > out[j].Size })
    return out
}
 
func main() {
    files := []FileStat{
        {"/var/log/syslog", 8 * 1024 * 1024},
        {"/var/log/auth",   256 * 1024},
        {"/etc/passwd",     1024},
        {"/tmp/cache",      64 * 1024 * 1024},
        {"/home/user/big",  512 * 1024 * 1024},
        {"/home/user/med",  2 * 1024 * 1024},
        {"/boot/vmlinuz",   12 * 1024 * 1024},
    }
    top := NewTopN(3)
    for _, f := range files { top.Add(f) }
    fmt.Println("Top 3 最大文件:")
    for i, f := range top.Result() {
        fmt.Printf("  %d. %d 字节  %s\\n", i+1, f.Size, f.Path)
    }
}

小结

  • ✅ Worker Pool 三对 channel:jobs、results、收尾 close
  • ✅ container/heap 实现 Top-N:O(N log k) 优于全排序
  • ✅ 最小堆找 Top-N 大的(堆顶是当前 N 大阈值)
  • ✅ CPU 密集任务:workers ≤ NumCPU
  • ✅ IO 密集任务(读文件、HTTP):workers 可以是 NumCPU * 2 ~ * 10
  • ✅ bufio.Scanner 配大 buffer 处理长行
  • ✅ 真实部署要 flag、log、信号处理、pprof
  • ✅ 沙盒里用内存数据模拟整个流程

🎉 Go 18 章完结!你已经能用 Go 写出生产级的 CLI 工具和 REST API。