Day 3:goroutine并发——50份文档同时分析
今天做什么
串行处理 50 份文档太慢了。今天用 goroutine + errgroup 让它们同时跑——50 份文档的耗时 = 最慢那份的耗时。
代码
go
package main
import (
"context"
"fmt"
"sync"
"time"
"golang.org/x/sync/errgroup"
)
// ===== 并发文档分析引擎 =====
type DocAnalyzer struct {
chain AnalysisChain
maxWorkers int
}
func NewAnalyzer(maxWorkers int) *DocAnalyzer {
return &DocAnalyzer{maxWorkers: maxWorkers}
}
func (a *DocAnalyzer) AnalyzeBatch(ctx context.Context,
docs []string) ([]AnalysisResult, error) {
// 用 errgroup 管理并发——一个出错全部取消
g, ctx := errgroup.WithContext(ctx)
// 信号量:限制并发数
sem := make(chan struct{}, a.maxWorkers)
results := make([]AnalysisResult, len(docs))
var mu sync.Mutex
for i, doc := range docs {
i, doc := i, doc // goroutine 闭包陷阱
g.Go(func() error {
// 获取信号量
select {
case sem <- struct{}{}:
defer func() { <-sem }()
case <-ctx.Done():
return ctx.Err()
}
// 分析单份文档
result, err := a.chain.Analyze(ctx, doc)
if err != nil {
return fmt.Errorf("文档 %d 分析失败: %w", i, err)
}
mu.Lock()
results[i] = result
mu.Unlock()
return nil
})
}
if err := g.Wait(); err != nil {
return nil, err
}
return results, nil
}
func main() {
ctx := context.Background()
analyzer := NewAnalyzer(5) // 最多5个并发
// 模拟50份文档
docs := generateDocs(50)
start := time.Now()
results, err := analyzer.AnalyzeBatch(ctx, docs)
elapsed := time.Since(start)
if err != nil {
fmt.Println("❌", err)
return
}
fmt.Printf("✅ 完成 %d 份文档分析,耗时 %v\n", len(results), elapsed)
fmt.Printf("平均每份: %v\n", elapsed/time.Duration(len(docs)))
// 串行对比:如果串行,耗时 ≈ 50 × 单份耗时
}性能对比
bash
# 串行:50份 × 300ms = 15秒
# 5并发:50份 × 300ms / 5 = 3秒
# 10并发:50份 × 300ms / 10 = 1.5秒你学到了什么
- errgroup:一个 goroutine 出错,其他全部通过 context 取消。
- 信号量限流:
chan struct{}控制最大并发数——别把 LLM API 打爆。 - sync.Mutex:并发写入 slice,需要锁保护。
明天的预告
有些文档可能反复分析(比如每天跑一次)。明天加上分析缓存——已经分析过的不重跑。
Day 3 完成。50 份文档同时分析,耗时从 15 秒降到 3 秒。

