Go 并发模式实战:Worker Pool 与 Pipeline 的优雅实现

一、最简 Worker Pool

Go
package main

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

func worker(id int, jobs <-chan int, results chan<- int, wg *sync.WaitGroup) {
	defer wg.Done()
	for job := range jobs {
		fmt.Printf("worker %d 处理任务 %d\n", id, job)
		time.Sleep(100 * time.Millisecond)
		results <- job * 2
	}
}

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

	var wg sync.WaitGroup
	for w := 1; w <= 5; w++ {
		wg.Add(1)
		go worker(w, jobs, results, &wg)
	}

	for j := 1; j <= 20; j++ {
		jobs <- j
	}
	close(jobs)

	wg.Wait()
	close(results)

	for r := range results {
		fmt.Println("结果:", r)
	}
}

二、支持优雅退出

Go
func workerPool(ctx context.Context, jobs <-chan Job, workers int) {
	var wg sync.WaitGroup
	for i := 0; i < workers; i++ {
		wg.Add(1)
		go func(id int) {
			defer wg.Done()
			for {
				select {
				case <-ctx.Done():
					fmt.Printf("worker %d 收到退出信号\n", id)
					return
				case job, ok := <-jobs:
					if !ok {
						return
					}
					process(job)
				}
			}
		}(i)
	}
	wg.Wait()
}

三、Pipeline 模式

Go
func gen(nums ...int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for _, n := range nums {
			out <- n
		}
	}()
	return out
}

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
}

func main() {
	for n := range square(gen(1, 2, 3, 4, 5)) {
		fmt.Println(n)
	}
}

四、使用 errgroup 管理错误

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

func fetchAll(urls []string) ([]string, error) {
	g, ctx := errgroup.WithContext(context.Background())
	results := make([]string, len(urls))

	for i, url := range urls {
		i, url := i, url
		g.Go(func() error {
			req, _ := http.NewRequestWithContext(ctx, "GET", url, nil)
			resp, err := http.DefaultClient.Do(req)
			if err != nil {
				return err
			}
			defer resp.Body.Close()
			body, _ := io.ReadAll(resp.Body)
			results[i] = string(body)
			return nil
		})
	}

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

五、限制并发数

Go
// 用带缓冲 channel 做信号量
sem := make(chan struct{}, 10)
for _, task := range tasks {
	sem <- struct{}{}
	go func(t Task) {
		defer func() { <-sem }()
		process(t)
	}(task)
}

六、常见坑

  1. goroutine 泄漏:channel 没人读也没人写,goroutine 永久阻塞
  2. for 循环变量捕获:Go 1.22 前必须 i := i 重新声明
  3. channel 忘记 close:range 永不结束
  4. 向已关闭 channel 发送:panic,应由发送方负责 close
打赏作者 已有 0 人打赏,共 ¥0.00
我的打赏
Go语言爱好者
Go语言爱好者
Lv6 原创 1 粉丝 1114

Go 微服务实践者,云原生布道师

  • 1文章
  • 3961总阅读
  • 171获赞
  • 1114粉丝

评论(0)

💬

还没有评论,来说两句吧