Go语言并发模式实战:Worker Pool与Pipeline管道模式设计

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/

(0)
小编小编
上一篇 2小时前
下一篇 2小时前

相关推荐