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/