Go语言的goroutine和channel为并发编程提供了语言级支持,但并发原语本身不等于并发设计。Worker Pool模式控制并发数量防止资源耗尽,Pipeline模式将任务分解为多阶段串行处理提升吞吐量,Fan-out/Fan-in模式实现任务分片与结果聚合。本文通过完整代码实现这三种核心并发模式,分析各模式适用场景和注意事项。
Worker Pool模式:固定并发数任务处理
Worker Pool是Go中最常用的并发模式,通过固定数量的worker goroutine消费任务通道中的任务,避免无限制创建goroutine导致的资源耗尽问题。
package main
import (
"context"
"fmt"
"sync"
"time"
)
type Task struct {
ID int
Input string
}
type Result struct {
TaskID int
Output string
Err error
}
func workerPool(ctx context.Context, tasks <-chan Task, results chan<- Result, workerCount int) {
var wg sync.WaitGroup
for i := 0; i < workerCount; i++ {
wg.Add(1)
go func(workerID int) {
defer wg.Done()
for {
select {
case <-ctx.Done():
return
case task, ok := <-tasks:
if !ok {
return
}
output, err := processTask(task)
results <- Result{
TaskID: task.ID,
Output: output,
Err: err,
}
}
}
}(i)
}
go func() {
wg.Wait()
close(results)
}()
}
func processTask(task Task) (string, error) {
time.Sleep(100 * time.Millisecond)
return fmt.Sprintf("processed-%s", task.Input), nil
}
func main() {
ctx := context.Background()
tasks := make(chan Task, 100)
results := make(chan Result, 100)
workerPool(ctx, tasks, results, 8)
// 发送任务
go func() {
for i := 0; i < 50; i++ {
tasks <- Task{ID: i, Input: fmt.Sprintf("task-%d", i)}
}
close(tasks)
}()
// 收集结果
for result := range results {
if result.Err != nil {
fmt.Printf("Task %d failed: %v\n", result.TaskID, result.Err)
} else {
fmt.Printf("Task %d: %s\n", result.TaskID, result.Output)
}
}
}
Worker Pool的关键设计点:tasks和results通道均使用缓冲通道,避免worker和生产者互相阻塞;context传递取消信号,支持优雅停机;使用sync.WaitGroup等待所有worker完成后关闭results通道,消费者通过range安全退出。
Pipeline模式:多阶段串行处理
Pipeline模式将复杂处理流程拆分为多个阶段,每个阶段一个或多个goroutine,通过通道连接。每个阶段的输出是下一阶段的输入,实现流水线并行处理。
package main
import (
"context"
"fmt"
)
// 阶段1:生成数据
func generate(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case <-ctx.Done():
return
case out <- n:
}
}
}()
return out
}
// 阶段2:平方运算
func square(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
select {
case <-ctx.Done():
return
case out <- n * n:
}
}
}()
return out
}
// 阶段3:格式化为字符串
func format(ctx context.Context, in <-chan int) <-chan string {
out := make(chan string)
go func() {
defer close(out)
for n := range in {
select {
case <-ctx.Done():
return
case out <- fmt.Sprintf("result: %d", n):
}
}
}()
return out
}
func main() {
ctx := context.Background()
// 构建Pipeline
nums := generate(ctx, 1, 2, 3, 4, 5)
squared := square(ctx, nums)
formatted := format(ctx, squared)
// 消费最终结果
for s := range formatted {
fmt.Println(s)
}
}
Pipeline每个阶段都是独立的goroutine,数据在各阶段间通过channel传递。各阶段可以有不同的处理速度,通过channel的缓冲区吸收速度差异。当上游阶段关闭channel时,下游阶段通过range自动感知并退出。
Fan-out/Fan-in模式:任务分片与结果聚合
Fan-out将一个阶段的输出分发到多个goroutine并行处理,Fan-in将多个goroutine的输出合并到一个通道。该模式适用于CPU密集型任务,通过并行处理提升整体吞吐量。
// Fan-out:启动多个square goroutine并行处理
func fanOut(ctx context.Context, in <-chan int, workerCount int) []<-chan int {
outs := make([]<-chan int, workerCount)
for i := 0; i < workerCount; i++ {
outs[i] = square(ctx, in)
}
return outs
}
// Fan-in:合并多个channel的输出
func fanIn(ctx context.Context, channels ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
multiplex := func(ch <-chan int) {
defer wg.Done()
for n := range ch {
select {
case <-ctx.Done():
return
case out <- n:
}
}
}
wg.Add(len(channels))
for _, ch := range channels {
go multiplex(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
ctx := context.Background()
// 生成数据 -> Fan-out到4个worker -> Fan-in合并结果
nums := generate(ctx, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
squaredChannels := fanOut(ctx, nums, 4)
merged := fanIn(ctx, squaredChannels...)
for n := range merged {
fmt.Println(n)
}
}
并发模式中的错误处理与取消传播
实际应用中Pipeline各阶段可能产生错误,需要统一收集并支持取消传播。使用errgroup包简化错误管理:
package main
import (
"context"
"fmt"
"strings"
"golang.org/x/sync/errgroup"
)
func pipelineWithErrorHandling(ctx context.Context) error {
g, ctx := errgroup.WithContext(ctx)
stage1 := make(chan string)
stage2 := make(chan string)
// 阶段1
g.Go(func() error {
defer close(stage1)
for i := 0; i < 10; i++ {
if i == 5 {
return fmt.Errorf("stage1 error at item %d", i)
}
select {
case <-ctx.Done():
return ctx.Err()
case stage1 <- fmt.Sprintf("item-%d", i):
}
}
return nil
})
// 阶段2
g.Go(func() error {
defer close(stage2)
for s := range stage1 {
processed := strings.ToUpper(s)
select {
case <-ctx.Done():
return ctx.Err()
case stage2 <- processed:
}
}
return nil
})
// 消费者
g.Go(func() error {
for s := range stage2 {
fmt.Println(s)
}
return nil
})
return g.Wait()
}
errgroup.WithContext创建的context在任一goroutine返回错误时自动取消,其他goroutine通过select监听ctx.Done()快速退出。g.Wait()等待所有goroutine结束后返回第一个非nil错误。这种模式确保错误发生后不会产生资源泄漏。
模式选型与组合策略
Worker Pool适合任务数量不确定、需要控制并发上限的场景,如HTTP请求处理、消息队列消费。Pipeline适合数据处理流程明确、各阶段计算量不同的场景,如ETL流水线、日志分析管道。Fan-out/Fan-in适合单一阶段计算密集、可并行处理的场景,如批量图像缩放、文档格式转换。
三种模式可以组合使用。一个完整的批处理系统可能:Worker Pool接收外部任务,每个任务进入Pipeline多阶段处理,Pipeline中计算密集阶段使用Fan-out/Fan-in并行加速。关键原则是保持每个goroutine的职责单一,通过channel传递数据而非共享内存,利用context实现统一的生命周期管理。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-bing-fa-mo-shi-shi-zhan-workerpool-yu-pipeline/