Go: Go 并发模式实战
最后更新:2026-08-26
Go 并发不是靠库,而是靠模式——Fan-in、Fan-out、Pipeline 这些组合模式能解决 90% 的并发问题。
当你需要处理大量数据:从一个 generator 读取、多个 worker 并行处理、最后合并结果——怎么设计才优雅?这节课你将掌握 Go 社区最经典的并发模式。
1. 你将学到
- Pipeline 流水线模式
- Fan-out / Fan-in 扇出扇入
- Worker Pool 工作池
- Or-Done 通道简化模式
- Cancellation Propagation 取消传播
golang.org/x/sync/errgroup错误组- 实战:实时数据处理管道
2. 一个数据工程师的真实故事
(1) 痛点:单线程处理 1000 万条日志,跑了一夜
Fatima 是数据平台工程师,公司每天产生 1000 万条访问日志需要实时处理:
"第一版是一个 for loop:读一条→解析→写入数据库。上线第一天发现处理速度跟不上生产速度——日志队列每小时增长 200 万条。老板说'明天报表就看这批数据'。我加了更多服务器,但代码还是单线程的——CPU 只用了 5%。"
她当时的代码:
// 坏代码:单线程,完全不用 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
// 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 流水线
▶ 示例:数字处理管道
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)
}
graph LR
A[generate] -->|chan int| B[square]
B -->|chan int| C[double]
C -->|chan int| D[main]
<-chan T(只读通道),接收 <-chan T 作为输入。每个 stage 在自己的 goroutine 中运行,通过 channel 连接——这就是 Go 的 CSP 模型。
4. Fan-out / Fan-in 扇出扇入
▶ 示例:Fan-out + Fan-in
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)
}
}
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
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)
}
}
(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 模式
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 退出
}
7. errgroup 错误传播
▶ 示例:errgroup
⚙️ 前置安装:运行
go get golang.org/x/sync/errgroup
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 超时
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. 完整示例:实时数据处理管道
// 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("管道处理完成")
}
expvar 或 prometheus 暴露每个 stage 的处理速度、队列长度和错误率。
❓ 常见问题
select 配合 done 通道即可。📖 小节
- Pipeline:数据流经过多个处理 stage,每个 stage 一个 goroutine
- Fan-out:一个 channel 分发到多个 worker(并行处理)
- Fan-in:多个 channel 合并到一个 channel(统一消费)
- Worker Pool:固定数量 worker 从 job 通道取任务
- Or-Done:安全消费不可控 channel,支持取消
- errgroup:管理 goroutine 组,支持错误收集和取消传播
- 组合模式:Pipeline 内部用 Worker Pool + Fan-out/Fan-in
📝 作业
-
基础题(难度⭐):实现一个 3 stage 数字处理 Pipeline:generate(1-100) → filter(只保留偶数) → sum(累加)。每个 stage 在独立 goroutine 中运行,通过 channel 连接。
-
进阶题(难度⭐⭐):实现一个并发的日志处理系统:10 个 parse worker(fan-out)+ 3 个 write worker(fan-in)+ errgroup 错误管理。模拟 1000 条日志,输出统计信息(总用时、处理条数、错误数)。
-
挑战题(难度⭐⭐⭐):设计并实现带负载均衡和动态扩缩容的 Worker Pool。要求:(1) 初始 5 个 worker;(2) 当 job 队列积压超过阈值时自动增加 worker(最多 20);(3) 当队列空闲时减少 worker(最少 2);(4) 用 Context 控制 worker 退出;(5) 用
-race验证无数据竞争。