Go语言并发模式实战:goroutine池化与channel通信机制

goroutine与channel基础通信模型

Go语言的并发模型基于CSP(Communicating Sequential Processes)理论,goroutine是轻量级协程,channel是goroutine间通信的管道。与线程池模型不同,goroutine的创建开销约2KB栈空间,单机可轻松启动百万级goroutine。

启动goroutine使用go关键字,channel通过make创建:

package main

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

func worker(id int, jobs <-chan int, results chan<- int) {
    for job := range jobs {
        time.Sleep(time.Second) // 模拟处理
        results <- job * job
    }
}

func main() {
    jobs := make(chan int, 100)
    results := make(chan int, 100)

    // 启动3个worker goroutine
    for w := 1; w <= 3; w++ {
        go worker(w, jobs, results)
    }

    // 发送5个任务
    for j := 1; j <= 5; j++ {
        jobs <- j
    }
    close(jobs)

    // 收集结果
    for r := 1; r <= 5; r++ {
        fmt.Println("result:", <-results)
    }
}

jobs <-chan int是只读channel,results chan<- int是只写channel,类型限定防止误操作。close(jobs)关闭channel后,range循环在消费完剩余数据后自动退出。

worker pool模式与并发数控制

无限制创建goroutine会导致内存暴涨和调度开销增大。worker pool模式预创建固定数量的goroutine,通过channel分发任务,实现并发数控制和资源复用:

type Task struct {
    ID   int
    Data interface{}
}

type WorkerPool struct {
    tasks    chan Task
    results  chan error
    workerN  int
    wg       sync.WaitGroup
}

func NewWorkerPool(workerNum, queueSize int) *WorkerPool {
    return &WorkerPool{
        tasks:   make(chan Task, queueSize),
        results: make(chan error, queueSize),
        workerN: workerNum,
    }
}

func (p *WorkerPool) Start(handler func(Task) error) {
    for i := 0; i < p.workerN; i++ {
        p.wg.Add(1)
        go func() {
            defer p.wg.Done()
            for task := range p.tasks {
                err := handler(task)
                p.results <- err
            }
        }()
    }
}

func (p *WorkerPool) Submit(task Task) {
    p.tasks <- task
}

func (p *WorkerPool) Stop() {
    close(p.tasks)
    p.wg.Wait()
    close(p.results)
}

使用示例:

pool := NewWorkerPool(10, 100)
pool.Start(func(t Task) error {
    return processOrder(t.Data)
})

for i := 0; i < 1000; i++ {
    pool.Submit(Task{ID: i, Data: orderData})
}
pool.Stop()

sync.WaitGroup确保所有worker处理完成后才关闭results channel。queueSize缓冲区大小根据任务产生速度和消费速度比例设定,过小会导致Submit阻塞,过大占用额外内存。

context超时控制与goroutine取消传播

context包实现goroutine的取消信号传播和超时控制,是Go并发编程的标准模式:

func fetchWithTimeout(ctx context.Context, url string) ([]byte, error) {
    ctx, cancel := context.WithTimeout(ctx, 5*time.Second)
    defer cancel()

    req, err := http.NewRequestWithContext(ctx, "GET", url, nil)
    if err != nil {
        return nil, err
    }

    resp, err := http.DefaultClient.Do(req)
    if err != nil {
        return nil, err
    }
    defer resp.Body.Close()

    return io.ReadAll(resp.Body)
}

func batchFetch(ctx context.Context, urls []string) {
    results := make(chan []byte, len(urls))

    for _, url := range urls {
        go func(u string) {
            data, err := fetchWithTimeout(ctx, u)
            if err != nil {
                log.Printf("fetch %s failed: %v", u, err)
                return
            }
            results <- data
        }(url)
    }

    for i := 0; i < len(urls); i++ {
        select {
        case data := <-results:
            process(data)
        case <-ctx.Done():
            log.Println("batch cancelled:", ctx.Err())
            return
        }
    }
}

// 调用方
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
batchFetch(ctx, urls)

context.WithTimeout创建带超时的context,超时后ctx.Done()返回的channel被关闭,所有持有该context的goroutine收到取消信号。defer cancel()确保函数退出时释放context资源。http.NewRequestWithContext将context注入HTTP请求,超时时请求自动中断。

errgroup并发错误处理与编排

标准库sync不提供并发错误聚合能力,golang.org/x/sync/errgroup包填补了这一空白:

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

func parallelProcess(ctx context.Context, items []Item) error {
    g, ctx := errgroup.WithContext(ctx)
    g.SetLimit(5) // 最大并发数5

    results := make([]Result, len(items))

    for i, item := range items {
        i, item := i, item // 捕获循环变量

        g.Go(func() error {
            result, err := process(ctx, item)
            if err != nil {
                return fmt.Errorf("item %d: %w", i, err)
            }
            results[i] = result
            return nil
        })
    }

    if err := g.Wait(); err != nil {
        return err
    }
    return nil
}

g.SetLimit(5)限制最多5个goroutine同时运行,g.Go在达到限制时阻塞。errgroup.WithContext返回的context在任意goroutine返回错误时自动取消,其他goroutine通过ctx.Done()感知并提前退出。g.Wait()等待所有goroutine完成,返回第一个非nil错误。

channel高级模式:select多路复用与扇出扇入

select语句同时监听多个channel,实现多路复用:

func mergeChannels(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
}

// 使用
ch1 := produce(1, 2, 3)
ch2 := produce(4, 5, 6)
ch3 := produce(7, 8, 9)

merged := mergeChannels(ch1, ch2, ch3)
for v := range merged {
    fmt.Println(v) // 1-9 乱序输出
}

扇出(Fan-out)将一个channel分发给多个消费者,扇入(Fan-in)将多个channel合并为一个。配合select实现非阻塞操作:

select {
case msg := <-ch:
    handle(msg)
case <-time.After(2 * time.Second):
    log.Println("timeout")
default:
    // 无数据时立即执行
    log.Println("no message available")
}

default分支使select非阻塞——没有就绪的channel时立即执行default分支,实现轮询模式。goroutine泄露是Go并发编程中需要警惕的问题:如果goroutine阻塞在channel读写且没有外部引用该channel,goroutine将永远无法退出。防御方法:所有goroutine都接受context参数,确保有取消路径;使用runtime.NumGoroutine()监控goroutine数量;通过pprof分析goroutine调用栈:go tool pprof http://localhost:6060/debug/pprof/goroutine

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-bing-fa-mo-shi-shi-zhan-goroutine-chi-hua-yu/

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

相关推荐