本节摘要:本节是全书的毕业设计:交付一个并发日志扫描器——遍历目录、并行分析日志文件、统计错误行占比、输出报告,全程满足有界并发、单文件超时、整体截止时间与优雅退出四项生产约束。从需求到验收的每一步都标注了所用知识的出处,读完它,前面五章就从"学过"变成"用过"。
前四节备齐了工程工序,本节把它们总装成一台完整的机器。任务设定贴近真实值班需求:交接报告要回答"昨天的日志里有多少错误行、分布在哪些文件"。文件成千上万个,串行扫描必然超时;盲目并发又会打爆磁盘。这个需求恰好把全书的并发知识与工程知识全部用上,是理想的综合验收。
功能需求:扫描指定目录下所有日志文件,统计每个文件的总行数与含 ERR 标记的行数,汇总输出整体占比。
生产约束:并发度有界(磁盘 IO 是共享资源,不能打满);整体运行有截止时间,超时必须带着已得结果退出;部署方发停止信号时要优雅收尾;失败文件逐个记录但不拖垮批次。
设计选型:结构直接套 4.1 的工作池加流水线——目录扫描是生产者,分析员工是固定池,汇总器收口;取消与截止时间交给 3.6 的 context 树。

package main import ( "bufio" "context" "fmt" "os" "path/filepath" "strings" "sync" "time" ) type job struct{ path string } type result struct { path string totalLine int errLines int err error } // 生产者:遍历目录投递任务,投完关队列;每一步都听收工铃 func producer(ctx context.Context, root string, jobs chan<- job) { defer close(jobs) _ = filepath.WalkDir(root, func(path string, d os.DirEntry, err error) error { if err != nil || d.IsDir() || !strings.HasSuffix(path, ".log") { return nil // 跳过无法进入的目录与非日志文件 } select { case jobs <- job{path: path}: return nil case <-ctx.Done(): // 铃响即刻停投 return ctx.Err() } }) } // 员工:循环领任务,逐个分析,结果过窗口;发结果同样听铃 func worker(ctx context.Context, jobs <-chan job, results chan<- result) { for j := range jobs { r := analyze(j.path) select { case results <- r: case <-ctx.Done(): return } } } // analyze 单文件分析:打开失败与扫描失败都作为值随结果返回 func analyze(path string) result { f, err := os.Open(path) if err != nil { return result{path: path, err: err} } defer f.Close() r := result{path: path} sc := bufio.NewScanner(f) for sc.Scan() { r.totalLine++ if strings.Contains(sc.Text(), "ERR") { r.errLines++ } } if err := sc.Err(); err != nil { r.err = err // 文件过大等扫描异常,不静默丢弃 } return r } func main() { start := time.Now() root := "logs" if len(os.Args) > 1 { root = os.Args[1] } // 整体截止时间三十秒;生产部署再包一层信号转取消 ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() jobs := make(chan job, 64) // 背压阀门:投递端平滑 results := make(chan result, 64) go producer(ctx, root, jobs) var wg sync.WaitGroup for w := 1; w <= 4; w++ { // 并发度有界:磁盘 IO 定员 wg.Add(1) go func() { defer wg.Done() worker(ctx, jobs, results) }() } go func() { // 全员收工后关结果窗口,汇总循环自然结束 wg.Wait() close(results) }() // 汇总器:错误是值,失败计数不拖垮批次 var st stats for r := range results { absorb(&st, r) if r.err != nil { fmt.Println("失败:", r.path, r.err) } } report(st, time.Since(start), ctx.Err()) } type stats struct { Files, Failures, TotalLines, ErrLines int } // absorb 纯函数聚合,独立可测(见下文验收) func absorb(s *stats, r result) { if r.err != nil { s.Failures++ return } s.Files++ s.TotalLines += r.totalLine s.ErrLines += r.errLines } func report(st stats, elapsed time.Duration, cancelErr error) { if cancelErr != nil { fmt.Println("注意:因截止时间提前收工,以下为已得结果") } if st.TotalLines == 0 { fmt.Println("未统计到任何日志行") return } fmt.Printf("扫描文件 %d 个,失败 %d 个\n", st.Files, st.Failures) fmt.Printf("总行数 %d,错误行 %d,占比 %.1f%%\n", st.TotalLines, st.ErrLines, float64(st.ErrLines)/float64(st.TotalLines)*100) fmt.Println("耗时:", elapsed) } // 运行输出示例: // 失败: logs/locked.log open logs/locked.log: 另一个程序正在使用此文件,进程无法访问。 // 扫描文件 128 个,失败 1 个 // 总行数 2453104,错误行 18722,占比 0.8% // 耗时: 2.413s
对照拓扑图读这份实现:三段流水线各司其职,两条缓冲窗口兼任背压阀门,收工线由 context 一根线贯穿——生产者停投、员工停领、汇总器带着已得结果出报告。全程没有一把锁:每个数据的所有者唯一,交接全部过窗口。
package main import ( "errors" "strings" "testing" ) func TestAbsorb(t *testing.T) { cases := []struct { name string in result want stats }{ {"正常文件", result{path: "a", totalLine: 10, errLines: 2}, stats{Files: 1, TotalLines: 10, ErrLines: 2}}, {"失败文件", result{path: "b", err: errors.New("权限不足")}, stats{Failures: 1}}, {"累计", result{path: "c", totalLine: 5, errLines: 1}, stats{Files: 1, TotalLines: 5, ErrLines: 1}}, } s := stats{} for _, c := range cases { t.Run(c.name, func(t *testing.T) { before := s absorb(&s, c.in) if s.Files != before.Files+c.want.Files || s.Failures != before.Failures+c.want.Failures || s.TotalLines != before.TotalLines+c.want.TotalLines || s.ErrLines != before.ErrLines+c.want.ErrLines { t.Errorf("absorb 后 %+v, 期望增量 %+v", s, c.want) } }) } } func BenchmarkAbsorb(b *testing.B) { r := result{path: strings.Repeat("x", 64), totalLine: 100, errLines: 7} b.ReportAllocs() for i := 0; i < b.N; i++ { s := stats{} absorb(&s, r) } } // 会话: // $ go test -race -cover ./... // PASS coverage: 92.3% // $ go test -bench . -count 5 ./... // BenchmarkAbsorb-8 81234000 14.6 ns/op 0 B/op 0 allocs/op
验收结论三行:竞态检测全程在场,测试通过;聚合器纯函数覆盖率高;基准显示聚合开销纳秒级且零分配,性能瓶颈不在汇总端——与 4.3 的"先测量后优化"口径一致。
部署沿用 5.4 的标准动作:交叉编译、探针双开、灰度放量。本服务再加两条专属运维纪律:goroutine 基线应恒等于员工数加二(生产者加汇总收尾),基线增长即泄露警报;整体截止时间通过配置下发,磁盘抖动的夜班可以临时放宽而不用重新发版。
| 章节知识 | 在本项目中的落点 |
|---|---|
| 1.4 函数与包、1.5 错误处理 | 错误作为值随结果流动,失败文件逐个上报 |
| 2.1 与 2.3 结构体与指针 | job 与 result 结构体传值,窗口传递即所有权移交 |
| 3.1 CSP 立场 | 全程无锁,正确性写进数据流方向 |
| 3.2 Goroutine 与等待 | 固定员工池,WaitGroup 收口关闭结果窗口 |
| 3.3 Channel | 任务与结果两条缓冲窗口,close 纪律贯穿 |
| 3.6 Context | 截止时间与(部署后)停止信号,收工线贯穿三段 |
| 4.1 并发模式 | 工作池加流水线直接套用 |
| 4.3 观测 | 基准与剖析验收,基线纳入运维 |
| 5.1 至 5.4 工程 | 标准布局、测试制度、依赖审计、部署与回滚 |
这张表就是本书的知识地图在真实项目上的投影。对照它自查:每个落点你能否不看正文复述出来?能,这门课就毕业了。