项目:并发文件处理工具
本节用 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:派发器 → workerresults:worker → 收集器close(jobs)由派发方调用;close(results)由收集方在 wg.Wait 之后调用- 严格遵守:发送方关,接收方只读
6. 单文件可运行版(Playground)
沙盒不能读真实文件系统,我们用**内存中的"文件"**模拟整个流程。
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 并发
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/,可分析性能
🎯 练习
// 任务:基于上面的 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。