Skip to content

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秒

你学到了什么

  1. errgroup:一个 goroutine 出错,其他全部通过 context 取消。
  2. 信号量限流chan struct{} 控制最大并发数——别把 LLM API 打爆。
  3. sync.Mutex:并发写入 slice,需要锁保护。

明天的预告

有些文档可能反复分析(比如每天跑一次)。明天加上分析缓存——已经分析过的不重跑。


Day 3 完成。50 份文档同时分析,耗时从 15 秒降到 3 秒。