Go: Go 并发模式实战

最后更新:2026-08-26

Go 并发不是靠库,而是靠模式——Fan-in、Fan-out、Pipeline 这些组合模式能解决 90% 的并发问题。

当你需要处理大量数据:从一个 generator 读取、多个 worker 并行处理、最后合并结果——怎么设计才优雅?这节课你将掌握 Go 社区最经典的并发模式。

1. 你将学到


2. 一个数据工程师的真实故事

(1) 痛点:单线程处理 1000 万条日志,跑了一夜

Fatima 是数据平台工程师,公司每天产生 1000 万条访问日志需要实时处理:

"第一版是一个 for loop:读一条→解析→写入数据库。上线第一天发现处理速度跟不上生产速度——日志队列每小时增长 200 万条。老板说'明天报表就看这批数据'。我加了更多服务器,但代码还是单线程的——CPU 只用了 5%。"

她当时的代码:

GO
// 坏代码:单线程,完全不用 CPU
func processLogs(logs []string) {
    for _, log := range logs {
        parsed := parse(log)    // CPU 密集型
        enriched := enrich(parsed) // CPU 密集型
        save(enriched)           // I/O 密集型
    }
    // 1000 万条日志:3 小时!
}

(2) Go 的并发解法:Pipeline + Fan-out

GO
// pipeline.go
package main

import (
    "fmt"
    "sync"
    "time"
)

type LogEntry struct {
    Raw     string
    Level   string
    Message string
    Time    time.Time
}

func main() {
    logs := make([]string, 100)
    for i := range logs {
        logs[i] = fmt.Sprintf("log-%d", i)
    }
    start := time.Now()

    // Pipeline: generator → parse → enrich → save
    rawCh := generator(logs)
    parsedCh := parseStage(rawCh, 4)   // 4 个 parse worker
    enrichedCh := enrichStage(parsedCh, 4) // 4 个 enrich worker
    saveStage(enrichedCh, 2) // 2 个 save worker

    fmt.Printf("处理完成,耗时: %v\n", time.Since(start))
}

// Stage 1: Generator(数据源)
func generator(logs []string) <-chan string {
    out := make(chan string, 100)
    go func() {
        defer close(out)
        for _, log := range logs {
            out <- log
        }
    }()
    return out
}

// Stage 2: Parse(Fan-out)
func parseStage(in <-chan string, workers int) <-chan LogEntry {
    out := make(chan LogEntry, 100)
    for i := 0; i < workers; i++ {
        go func() {
            for raw := range in {
                time.Sleep(1 * time.Millisecond) // 模拟解析
                out <- LogEntry{Raw: raw, Level: "INFO", Message: raw, Time: time.Now()}
            }
        }()
    }
    return out
}

// Stage 3: Enrich(Fan-out)
func enrichStage(in <-chan LogEntry, workers int) <-chan LogEntry {
    out := make(chan LogEntry, 100)
    for i := 0; i < workers; i++ {
        go func() {
            for entry := range in {
                time.Sleep(2 * time.Millisecond) // 模拟丰富数据
                out <- entry
            }
        }()
    }
    return out
}

// Stage 4: Save(Fan-in to 2 workers)
func saveStage(in <-chan LogEntry, workers int) {
    var wg sync.WaitGroup
    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for entry := range in {
                _ = entry // 模拟写入数据库
                time.Sleep(500 * time.Microsecond)
            }
        }()
    }
    wg.Wait()
}

(3) 收益:单线程 vs Pipeline

指标 单线程 Pipeline(4+4+2 workers) 提升
1000 万条处理时间 ~3 小时 ~12 分钟 15x
CPU 使用率 5% 85% 17x
代码复杂度 简单 中等
可扩展性 加服务器 加 worker(改一个数字)

3. Pipeline 流水线

▶ 示例:数字处理管道

GO
package main

import (
    "fmt"
)

// Stage 1: Generate
func generate(nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            out <- n
        }
    }()
    return out
}

// Stage 2: Square
func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            out <- n * n
        }
    }()
    return out
}

// Stage 3: Double
func double(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            out <- n * 2
        }
    }()
    return out
}

func main() {
    // Pipeline: generate → square → double
    for result := range double(square(generate(1, 2, 3, 4, 5))) {
        fmt.Println(result)
    }
    // 输出: 2, 8, 18, 32, 50  (n²×2)
}
▶ 试一试
100%
graph LR
    A[generate] -->|chan int| B[square]
    B -->|chan int| C[double]
    C -->|chan int| D[main]
💡 提示: Pipeline 的关键是每个 stage 返回 <-chan T(只读通道),接收 <-chan T 作为输入。每个 stage 在自己的 goroutine 中运行,通过 channel 连接——这就是 Go 的 CSP 模型。


4. Fan-out / Fan-in 扇出扇入

▶ 示例:Fan-out + Fan-in

GO 📖 仅展示
package main

import (
    "fmt"
    "sync"
)

// Stage 1: Generate
func generate(nums ...int) <-chan int {
    out := make(chan int, 10)
    go func() {
        defer close(out)
        for _, n := range nums {
            out <- n
        }
    }()
    return out
}

// Stage 2: Worker(Fan-out target)
func worker(id int, in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            result := n * n
            fmt.Printf("Worker %d: %d^2 = %d\n", id, n, result)
            out <- result
        }
    }()
    return out
}

// Fan-in: 合并多个 channel
func fanIn(channels ...<-chan int) <-chan int {
    out := make(chan int)
    var wg sync.WaitGroup

    for _, ch := range channels {
        wg.Add(1)
        go func(c <-chan int) {
            defer wg.Done()
            for v := range c {
                out <- v
            }
        }(ch)
    }

    go func() {
        wg.Wait()
        close(out)
    }()

    return out
}

func main() {
    // Fan-out: 一个 generator 分发到 3 个 worker
    in := generate(1, 2, 3, 4, 5, 6)
    w1 := worker(1, in)
    w2 := worker(2, in)
    w3 := worker(3, in)

    // Fan-in: 3 个 worker 的结果合并到一个 channel
    results := fanIn(w1, w2, w3)

    for result := range results {
        fmt.Printf("结果: %d\n", result)
    }
}
逻辑代码 55 行(超过 40 行限制,仅展示)
100%
graph LR
    G[Generate] -->|fan-out| W1[Worker 1]
    G -->|fan-out| W2[Worker 2]
    G -->|fan-out| W3[Worker 3]
    W1 -->|fan-in| R[Results]
    W2 -->|fan-in| R
    W3 -->|fan-in| R

5. Worker Pool 工作池

▶ 示例:Worker Pool

GO 📖 仅展示
package main

import (
    "fmt"
    "sync"
    "time"
)

type Job struct {
    ID      int
    Payload string
}

type Result struct {
    JobID   int
    Output  string
    Err     error
}

func worker(id int, jobs <-chan Job, results chan<- Result, wg *sync.WaitGroup) {
    defer wg.Done()
    for job := range jobs {
        fmt.Printf("Worker %d 处理 job %d\n", id, job.ID)
        time.Sleep(100 * time.Millisecond) // 模拟工作
        results <- Result{
            JobID:  job.ID,
            Output: fmt.Sprintf("processed by worker %d", id),
        }
    }
}

func main() {
    const numJobs = 20
    const numWorkers = 5

    jobs := make(chan Job, numJobs)
    results := make(chan Result, numJobs)

    // 启动 Worker Pool
    var wg sync.WaitGroup
    for i := 0; i < numWorkers; i++ {
        wg.Add(1)
        go worker(i, jobs, results, &wg)
    }

    // 发送 job
    for i := 0; i < numJobs; i++ {
        jobs <- Job{ID: i, Payload: fmt.Sprintf("data-%d", i)}
    }
    close(jobs)

    // 等待所有 worker 完成
    wg.Wait()
    close(results)

    // 收集结果
    for result := range results {
        fmt.Printf("Job %d → %s\n", result.JobID, result.Output)
    }
}
逻辑代码 46 行(超过 40 行限制,仅展示)

(1) Pipeline vs Worker Pool vs Fan-out/Fan-in

模式 核心思想 适用场景
Pipeline 数据流经过多个 stage,每个 stage 处理一步 有明确处理步骤(parse→transform→save)
Worker Pool 固定数量的 worker 从 job 通道取任务 批量任务处理(图片压缩、邮件发送)
Fan-out/Fan-in 分发到多个 worker 并行处理,再合并结果 无状态并行计算(数字运算、数据过滤)

6. Or-Done 通道

▶ 示例:Or-Done 模式

GO 📖 仅展示
package main

import (
    "fmt"
    "time"
)

// orDone 包装一个 channel,使其支持通过 done 通道取消
func orDone(done <-chan struct{}, c <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for {
            select {
            case <-done:
                return
            case v, ok := <-c:
                if !ok {
                    return
                }
                select {
                case out <- v:
                case <-done:
                    return
                }
            }
        }
    }()
    return out
}

func generateWithCancel(done <-chan struct{}, nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            select {
            case <-done:
                return
            case out <- n:
            }
        }
    }()
    return out
}

func main() {
    done := make(chan struct{})

    // 启动 generator(将在 3 秒后被取消)
    nums := generateWithCancel(done, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10)

    // 通过 orDone 安全消费
    for n := range orDone(done, nums) {
        fmt.Printf("处理: %d\n", n)
        if n >= 5 {
            close(done) // 取消所有操作
        }
    }

    fmt.Println("已取消")
    time.Sleep(100 * time.Millisecond) // 等待 goroutine 退出
}
逻辑代码 53 行(超过 40 行限制,仅展示)
🔥 易错: 不从 channel 中消费会导致 goroutine 泄漏——generator goroutine 会阻塞在发送上。Or-Done 模式确保在取消时不会泄漏 goroutine。 当上游发送者和下游消费者都需要感知取消时,用 Or-Done 包装 channel。


7. errgroup 错误传播

▶ 示例:errgroup

⚙️ 前置安装:运行 go get golang.org/x/sync/errgroup

GO
package main

import (
    "fmt"
    "time"

    "golang.org/x/sync/errgroup"
)

func main() {
    // 创建一个 errgroup
    // 第一个返回 error 的 goroutine 会取消其他 goroutine
    g := errgroup.Group{}

    // 启动 3 个 goroutine
    for i := 0; i < 3; i++ {
        id := i
        g.Go(func() error {
            return doWork(id)
        })
    }

    // 等待所有 goroutine 完成,获取第一个 error
    if err := g.Wait(); err != nil {
        fmt.Printf("任务失败: %v\n", err)
    } else {
        fmt.Println("全部成功")
    }
}

func doWork(id int) error {
    fmt.Printf("Worker %d 开始\n", id)
    time.Sleep(time.Duration(id+1) * 500 * time.Millisecond)

    if id == 1 {
        return fmt.Errorf("worker %d 出错", id)
    }

    fmt.Printf("Worker %d 完成\n", id)
    return nil
}
▶ 试一试

▶ 示例:errgroup + Context 超时

GO
package main

import (
    "context"
    "fmt"
    "time"

    "golang.org/x/sync/errgroup"
)

func main() {
    // 带 Context 的 errgroup
    ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
    defer cancel()

    g, ctx := errgroup.WithContext(ctx)

    // Worker 检查 Context 是否取消
    for i := 0; i < 5; i++ {
        id := i
        g.Go(func() error {
            select {
            case <-time.After(time.Duration(id+1) * time.Second):
                fmt.Printf("Worker %d 完成\n", id)
                return nil
            case <-ctx.Done():
                fmt.Printf("Worker %d 被取消: %v\n", id, ctx.Err())
                return ctx.Err()
            }
        })
    }

    if err := g.Wait(); err != nil {
        fmt.Printf("errgroup 停止: %v\n", err)
    }
}
▶ 试一试

(2) sync.WaitGroup vs errgroup

特性 sync.WaitGroup errgroup
错误收集 不支持 ✅ 返回第一个 error
取消传播 不支持 ✅ WithContext 自动取消
goroutine 管理 手动 Add/Done ✅ Go 方法自动管理
超时控制 手动实现 ✅ 结合 Context 轻松实现

8. 完整示例:实时数据处理管道

GO
// data_pipeline.go
package main

import (
    "context"
    "fmt"
    "math/rand"
    "sync"
    "time"
)

// ---------- 数据类型 ----------

type DataPoint struct {
    ID        int
    Value     float64
    Timestamp time.Time
}

type ProcessedData struct {
    DataPoint
    Normalized bool
    Anomaly    bool
    Category   string
}

// ---------- Pipeline Stages ----------

// Stage 1: Source(数据源)
func source(ctx context.Context, count int) <-chan DataPoint {
    out := make(chan DataPoint, 100)
    go func() {
        defer close(out)
        for i := 0; i < count; i++ {
            select {
            case <-ctx.Done():
                return
            case out <- DataPoint{
                ID:        i,
                Value:     rand.Float64() * 100,
                Timestamp: time.Now(),
            }:
            }
        }
    }()
    return out
}

// Stage 2: Normalize(归一化,Fan-out)
func normalize(ctx context.Context, in <-chan DataPoint, workers int) <-chan ProcessedData {
    out := make(chan ProcessedData, 100)
    var wg sync.WaitGroup

    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for dp := range orDone(ctx.Done(), in) {
                select {
                case <-ctx.Done():
                    return
                case out <- ProcessedData{
                    DataPoint:  dp,
                    Normalized: true,
                    Category:   categorize(dp.Value),
                }:
                }
            }
        }()
    }

    go func() {
        wg.Wait()
        close(out)
    }()

    return out
}

// Stage 3: Detect Anomaly(异常检测,Fan-out)
func detectAnomaly(ctx context.Context, in <-chan ProcessedData, workers int) <-chan ProcessedData {
    out := make(chan ProcessedData, 100)
    var wg sync.WaitGroup

    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for pd := range orDone(ctx.Done(), in) {
                pd.Anomaly = pd.Value > 90 || pd.Value < 10
                select {
                case <-ctx.Done():
                    return
                case out <- pd:
                }
            }
        }()
    }

    go func() {
        wg.Wait()
        close(out)
    }()

    return out
}

// Stage 4: Sink(输出,Fan-in)
func sink(ctx context.Context, in <-chan ProcessedData) {
    var normalCount, anomalyCount int
    for pd := range orDone(ctx.Done(), in) {
        if pd.Anomaly {
            anomalyCount++
            fmt.Printf("[异常] ID=%d, Value=%.2f, Category=%s\n",
                pd.ID, pd.Value, pd.Category)
        } else {
            normalCount++
        }
    }
    fmt.Printf("\n统计: 正常=%d, 异常=%d\n", normalCount, anomalyCount)
}

// 辅助函数
func orDone(done <-chan struct{}, in <-chan DataPoint) <-chan DataPoint {
    out := make(chan DataPoint)
    go func() {
        defer close(out)
        for {
            select {
            case <-done:
                return
            case v, ok := <-in:
                if !ok {
                    return
                }
                select {
                case out <- v:
                case <-done:
                    return
                }
            }
        }
    }()
    return out
}

func categorize(value float64) string {
    switch {
    case value < 30:
        return "low"
    case value < 70:
        return "medium"
    default:
        return "high"
    }
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
    defer cancel()

    fmt.Println("启动实时数据处理管道...")

    // 构建 Pipeline
    data := source(ctx, 1000)
    normalized := normalize(ctx, data, 5)    // 5 个归一化 worker
    detected := detectAnomaly(ctx, normalized, 3) // 3 个异常检测 worker
    sink(ctx, detected)

    fmt.Println("管道处理完成")
}
💡 提示: 在实际生产环境中,Pipeline 的每个 stage 应该有自己的缓冲区大小、重试逻辑和监控指标。可以用 expvarprometheus 暴露每个 stage 的处理速度、队列长度和错误率。


❓ 常见问题

Q Pipeline 和 Worker Pool 有什么区别?
A Pipeline 是数据流经过多个处理阶段,每个阶段处理一步(如解析→转换→保存),数据形状可以变。Worker Pool 是固定数量 worker 从同一个 job 通道取任务,数据形状不变(如图片压缩)。
Q Fan-out 和 Fan-in 的核心思想是什么?
A Fan-out 是把一个 channel 的数据分发到多个 worker 并行处理(提高吞吐量)。Fan-in 是把多个 channel 的结果合并到一个 channel(统一消费)。两者通常配对使用。
Q 什么时候必须用 Or-Done?
A 当你从无法控制 goroutine 生命周期的 channel 读取时(比如第三方库返回的 channel)。Or-Done 确保在调用者取消时,读取 goroutine 不会阻塞在 channel 操作上。简单的场景用 select 配合 done 通道即可。
Q errgroup 和 sync.WaitGroup 怎么选?
A 需要错误收集或取消传播时用 errgroup。只是等待 goroutine 完成时用 sync.WaitGroup(更轻量)。errgroup 是 WaitGroup 的功能超集,但会额外引入一个 goroutine 来管理取消。
Q Pipeline 中 channel 的缓冲区大小怎么定?
A 基本原则是让生产者不会阻塞在消费者上。缓冲区过大浪费内存,过小增加 goroutine 切换。经验值:缓冲区 = 生产速度 × 期望延迟。实际生产中通过压测确定。
Q goroutine 泄漏怎么防止?
A 三条规则:(1) 有创建 goroutine 的地方就有退出机制(channel 关闭或 done 信号);(2) 用 errgroup/WaitGroup 管理生命周期;(3) 使用 Or-Done 包装不可控的 channel。永远不要创建不知道自己何时退出的 goroutine。
Q 这些模式可以组合使用吗?
A 可以,而且应该组合。例如:Pipeline 中每个 stage 内部可以用 Worker Pool + Fan-out/Fan-in。errgroup 可以作为 Pipeline 编排器统一管理取消和错误。

📖 小节


📝 作业

  1. 基础题(难度⭐):实现一个 3 stage 数字处理 Pipeline:generate(1-100) → filter(只保留偶数) → sum(累加)。每个 stage 在独立 goroutine 中运行,通过 channel 连接。

  2. 进阶题(难度⭐⭐):实现一个并发的日志处理系统:10 个 parse worker(fan-out)+ 3 个 write worker(fan-in)+ errgroup 错误管理。模拟 1000 条日志,输出统计信息(总用时、处理条数、错误数)。

  3. 挑战题(难度⭐⭐⭐):设计并实现带负载均衡和动态扩缩容的 Worker Pool。要求:(1) 初始 5 个 worker;(2) 当 job 队列积压超过阈值时自动增加 worker(最多 20);(3) 当队列空闲时减少 worker(最少 2);(4) 用 Context 控制 worker 退出;(5) 用 -race 验证无数据竞争。

Web-Tutorial.com

Web-Tutorial 技术团队

由多位开发者共同维护的编程教程平台。每篇教程由对应领域的开发者编写和审核,确保内容准确可靠。如发现任何问题,欢迎向我们反馈。

100%

🙏 帮我们做得更好

我们是刚上线的编程教程站,几个人的小团队,精力有限。页面虽经检查,难免还有疏漏——链接失效、排版错乱、内容有误、语言生硬……

如果您发现了,麻烦告诉我们,我们会在收到反馈后第一时间进行修复,再次感谢您的光临 🙏