并发编程是 Go 语言最鲜明的标签也是它能在云原生时代脱颖而出的关键。Go 没有沿用传统的线程加锁模型而是提供了 Goroutine轻量级协程和 Channel通信管道这对组合拳。Goroutine 让你用 go 关键字就能启动一个并发任务Channel 则提供了安全的 Goroutine 间通信机制。本文深入讲解 Goroutine 的启动与调度、Channel 的两种类型无缓冲/有缓冲、select 多路复用、sync 包中的 WaitGroup 和 Mutex并通过生产者-消费者和工作池两个实战案例帮你掌握 Go 并发编程的核心技能。一、Goroutine轻量级并发单元1.1 什么是 GoroutineGoroutine 是 Go 运行时管理的轻量级线程。它的初始栈大小仅为 2KB操作系统线程通常为 1-2MB可以轻松创建成千上万个 Goroutine 而不会耗尽内存。启动一个 Goroutine 只需要在函数调用前加上 go 关键字packagemainimport(fmttime)funcsayHello(){fmt.Println(Hello from Goroutine!)}funcmain(){// 启动一个 Goroutine 执行 sayHellogosayHello()// 主 Goroutine 等待 1 秒否则子 Goroutine 可能还没执行完主程序就退出了time.Sleep(1*time.Second)fmt.Println(Main function finished)}1.2 Goroutine 的调度模型GMPGo 运行时的调度器采用 GMP 模型调度器的工作原理GOMAXPROCS 决定了 P 的数量默认等于 CPU 核心数。每个 P 维护一个本地 Goroutine 队列。M 从 P 的队列中获取 G 并执行。当一个 G 阻塞如等待 I/O时M 会与 P 分离P 会寻找或创建新的 M 继续执行其他 G。当 G 被阻塞时M 可以继续服务于其他 G从而实现高效的并发。设置 P 的数量通常不需要手动调整importruntimefuncmain(){// 获取 CPU 核心数fmt.Println(runtime.NumCPU())// 设置 P 的数量一般保持默认runtime.GOMAXPROCS(4)}1.3 sync.WaitGroup等待所有 Goroutine 完成使用 time.Sleep 等待不是一个好的做法。sync.WaitGroup 提供了更可靠的等待机制packagemainimport(fmtsync)funcworker(idint,wg*sync.WaitGroup){deferwg.Done()// 函数结束时通知 WaitGroup 该任务已完成fmt.Printf(Worker %d 开始工作\n,id)// 模拟耗时操作// time.Sleep(time.Second)fmt.Printf(Worker %d 完成\n,id)}funcmain(){varwg sync.WaitGroupfori:0;i5;i{wg.Add(1)// 增加一个等待计数goworker(i,wg)}wg.Wait()// 阻塞等待所有 Goroutine 完成fmt.Println(所有 Worker 已完成)}⚠️ 注意事项wg.Add(1) 必须在启动 Goroutine 之前调用不能在 Goroutine 内部调用否则可能因竞争条件导致计数错误。wg.Done() 等价于 wg.Add(-1)可以用 defer 确保一定执行。WaitGroup 必须通过指针传递值传递会拷贝计数器导致死锁。二、ChannelGoroutine 间的通信管道Go 并发模型的核心理念是不要通过共享内存来通信而要通过通信来共享内存。 Channel 正是实现这一理念的核心工具它是一个类型化的通信管道用于在 Goroutine 之间传递数据。2.1 无缓冲 Channel同步 Channel无缓冲 Channel 在发送和接收时都会阻塞直到另一端准备就绪。它天然实现了 Goroutine 之间的同步。funcmain(){// 创建一个无缓冲的 int 类型 Channelch:make(chanint)// 启动一个 Goroutine 发送数据gofunc(){fmt.Println(发送数据 42)ch-42// 发送阻塞直到有接收者fmt.Println(数据已发送完成)}()// 主 Goroutine 接收数据fmt.Println(等待接收数据...)value:-ch// 接收阻塞直到有发送者fmt.Printf(接收到数据: %d\n,value)// 输出顺序// 等待接收数据...// 发送数据 42// 接收到数据: 42// 数据已发送完成} 关键理解无缓冲 Channel 要求发送和接收同时就绪否则会阻塞等待。这提供了一种天然的同步机制。2.2 有缓冲 Channel异步 Channel有缓冲 Channel 在缓冲区未满时发送不阻塞在缓冲区非空时接收不阻塞。funcmain(){// 创建一个容量为 2 的缓冲 Channelch:make(chanint,2)// 发送两个数据不会阻塞ch-1ch-2fmt.Println(两个数据已发送缓冲区已满)// 再发送第三个数据会阻塞缓冲区已满// ch - 3 // 死锁// 接收数据fmt.Println(-ch)// 1fmt.Println(-ch)// 2}2.3 关闭 Channel 与遍历发送方可以主动关闭 Channel接收方可以通过 v, ok : -ch 检测 Channel 是否已关闭。更常见的是使用 for range 自动遍历直到 Channel 关闭。funcproducer(chchan-int){// 只发送 Channelfori:0;i5;i{ch-i}close(ch)// 发送完成后关闭}funcmain(){ch:make(chanint,3)goproducer(ch)// 使用 range 遍历 Channel直到它被关闭forv:rangech{fmt.Println(接收到:,v)}fmt.Println(Channel 已关闭遍历结束)}⚠️ Channel 操作规范只能由发送方关闭 Channel接收方不应关闭。关闭已关闭的 Channel 会引发 panic。向已关闭的 Channel 发送数据会引发 panic。从已关闭的 Channel 接收数据会返回零值v, ok : -ch 中 ok 为 false。2.4 Channel 的方向只读/只写在函数签名中可以限制 Channel 的方向提高类型安全// 只发送 ChannelfuncsendOnly(chchan-int,valueint){ch-value// -ch // 编译错误不能从只发送 Channel 接收}// 只接收 ChannelfuncreceiveOnly(ch-chanint){value:-ch// ch - 1 // 编译错误不能向只接收 Channel 发送fmt.Println(value)}funcmain(){ch:make(chanint)gosendOnly(ch,42)receiveOnly(ch)}三、select 多路复用select 语句用于同时等待多个 Channel 操作类似于 switch但每个 case 是一个 Channel 操作funcmain(){ch1:make(chanstring)ch2:make(chanstring)gofunc(){time.Sleep(2*time.Second)ch1-来自 ch1 的消息}()gofunc(){time.Sleep(1*time.Second)ch2-来自 ch2 的消息}()// 等待最先到达的消息select{casemsg1:-ch1:fmt.Println(msg1)casemsg2:-ch2:fmt.Println(msg2)case-time.After(3*time.Second):fmt.Println(超时 3 秒)}// 输出: 来自 ch2 的消息因为 ch2 先到达}3.1 select 的典型用途1超时控制funcdoWork()string{time.Sleep(2*time.Second)return工作完成}funcmain(){ch:make(chanstring)gofunc(){ch-doWork()}()select{caseresult:-ch:fmt.Println(result)case-time.After(1*time.Second):fmt.Println(超时工作未在 1 秒内完成)}// 输出: 超时工作未在 1 秒内完成}2非阻塞的 Channel 操作使用 default 分支实现非阻塞操作funcmain(){ch:make(chanint)select{casev:-ch:fmt.Println(接收到:,v)default:fmt.Println(没有数据可接收不阻塞)}select{casech-42:fmt.Println(发送成功)default:fmt.Println(缓冲区已满不阻塞)}}四、互斥锁保护共享资源虽然 Channel 是 Go 推荐的并发通信方式但某些场景下共享内存加锁仍然更合适如缓存、状态管理。sync.Mutex 提供了互斥锁sync.RWMutex 提供了读写锁允许并发读互斥写。typeCounterstruct{mu sync.Mutex valueint}func(c*Counter)Inc(){c.mu.Lock()c.valuec.mu.Unlock()}func(c*Counter)Value()int{c.mu.Lock()deferc.mu.Unlock()returnc.value}funcmain(){varwg sync.WaitGroup counter:Counter{}fori:0;i1000;i{wg.Add(1)gofunc(){deferwg.Done()counter.Inc()}()}wg.Wait()fmt.Println(counter.Value())// 1000安全}五、实战并发工作池Worker Pool工作池模式是一种经典的并发设计模式用于限制同时执行的并发任务数量避免资源耗尽packagemainimport(fmtsynctime)// 任务结构体typeTaskstruct{IDintDatastring}// 工作池typeWorkerPoolstruct{numWorkersinttaskQueuechanTask wg sync.WaitGroup}// 创建工作池funcNewWorkerPool(numWorkersint,queueSizeint)*WorkerPool{returnWorkerPool{numWorkers:numWorkers,taskQueue:make(chanTask,queueSize),}}// 启动工作池func(wp*WorkerPool)Start(){fori:0;iwp.numWorkers;i{wp.wg.Add(1)gowp.worker(i)}}// 工作 Goroutinefunc(wp*WorkerPool)worker(idint){deferwp.wg.Done()fortask:rangewp.taskQueue{fmt.Printf(Worker %d 开始处理任务 %d: %s\n,id,task.ID,task.Data)time.Sleep(100*time.Millisecond)// 模拟处理时间fmt.Printf(Worker %d 完成任务 %d\n,id,task.ID)}}// 提交任务func(wp*WorkerPool)Submit(task Task){wp.taskQueue-task}// 关闭工作池等待所有任务完成func(wp*WorkerPool)Stop(){close(wp.taskQueue)wp.wg.Wait()}funcmain(){// 创建 3 个工作 Goroutine队列容量 10pool:NewWorkerPool(3,10)pool.Start()// 提交 10 个任务fori:0;i10;i{pool.Submit(Task{ID:i,Data:fmt.Sprintf(data-%d,i)})}// 等待所有任务完成pool.Stop()fmt.Println(所有任务已完成)}六、小结Goroutinego 关键字启动轻量级协程栈大小仅 2KBGMP 调度模型高效管理。Channel类型化的通信管道无缓冲同步和有缓冲异步两种类型。select多路复用 Channel 操作支持超时控制和默认分支非阻塞。WaitGroup等待一组 Goroutine 完成Add、Done、Wait 三方法组合使用。Mutexsync.Mutex 互斥锁sync.RWMutex 读写锁保护共享资源。
第六篇:《并发编程核心:Goroutine 与 Channel》
并发编程是 Go 语言最鲜明的标签也是它能在云原生时代脱颖而出的关键。Go 没有沿用传统的线程加锁模型而是提供了 Goroutine轻量级协程和 Channel通信管道这对组合拳。Goroutine 让你用 go 关键字就能启动一个并发任务Channel 则提供了安全的 Goroutine 间通信机制。本文深入讲解 Goroutine 的启动与调度、Channel 的两种类型无缓冲/有缓冲、select 多路复用、sync 包中的 WaitGroup 和 Mutex并通过生产者-消费者和工作池两个实战案例帮你掌握 Go 并发编程的核心技能。一、Goroutine轻量级并发单元1.1 什么是 GoroutineGoroutine 是 Go 运行时管理的轻量级线程。它的初始栈大小仅为 2KB操作系统线程通常为 1-2MB可以轻松创建成千上万个 Goroutine 而不会耗尽内存。启动一个 Goroutine 只需要在函数调用前加上 go 关键字packagemainimport(fmttime)funcsayHello(){fmt.Println(Hello from Goroutine!)}funcmain(){// 启动一个 Goroutine 执行 sayHellogosayHello()// 主 Goroutine 等待 1 秒否则子 Goroutine 可能还没执行完主程序就退出了time.Sleep(1*time.Second)fmt.Println(Main function finished)}1.2 Goroutine 的调度模型GMPGo 运行时的调度器采用 GMP 模型调度器的工作原理GOMAXPROCS 决定了 P 的数量默认等于 CPU 核心数。每个 P 维护一个本地 Goroutine 队列。M 从 P 的队列中获取 G 并执行。当一个 G 阻塞如等待 I/O时M 会与 P 分离P 会寻找或创建新的 M 继续执行其他 G。当 G 被阻塞时M 可以继续服务于其他 G从而实现高效的并发。设置 P 的数量通常不需要手动调整importruntimefuncmain(){// 获取 CPU 核心数fmt.Println(runtime.NumCPU())// 设置 P 的数量一般保持默认runtime.GOMAXPROCS(4)}1.3 sync.WaitGroup等待所有 Goroutine 完成使用 time.Sleep 等待不是一个好的做法。sync.WaitGroup 提供了更可靠的等待机制packagemainimport(fmtsync)funcworker(idint,wg*sync.WaitGroup){deferwg.Done()// 函数结束时通知 WaitGroup 该任务已完成fmt.Printf(Worker %d 开始工作\n,id)// 模拟耗时操作// time.Sleep(time.Second)fmt.Printf(Worker %d 完成\n,id)}funcmain(){varwg sync.WaitGroupfori:0;i5;i{wg.Add(1)// 增加一个等待计数goworker(i,wg)}wg.Wait()// 阻塞等待所有 Goroutine 完成fmt.Println(所有 Worker 已完成)}⚠️ 注意事项wg.Add(1) 必须在启动 Goroutine 之前调用不能在 Goroutine 内部调用否则可能因竞争条件导致计数错误。wg.Done() 等价于 wg.Add(-1)可以用 defer 确保一定执行。WaitGroup 必须通过指针传递值传递会拷贝计数器导致死锁。二、ChannelGoroutine 间的通信管道Go 并发模型的核心理念是不要通过共享内存来通信而要通过通信来共享内存。 Channel 正是实现这一理念的核心工具它是一个类型化的通信管道用于在 Goroutine 之间传递数据。2.1 无缓冲 Channel同步 Channel无缓冲 Channel 在发送和接收时都会阻塞直到另一端准备就绪。它天然实现了 Goroutine 之间的同步。funcmain(){// 创建一个无缓冲的 int 类型 Channelch:make(chanint)// 启动一个 Goroutine 发送数据gofunc(){fmt.Println(发送数据 42)ch-42// 发送阻塞直到有接收者fmt.Println(数据已发送完成)}()// 主 Goroutine 接收数据fmt.Println(等待接收数据...)value:-ch// 接收阻塞直到有发送者fmt.Printf(接收到数据: %d\n,value)// 输出顺序// 等待接收数据...// 发送数据 42// 接收到数据: 42// 数据已发送完成} 关键理解无缓冲 Channel 要求发送和接收同时就绪否则会阻塞等待。这提供了一种天然的同步机制。2.2 有缓冲 Channel异步 Channel有缓冲 Channel 在缓冲区未满时发送不阻塞在缓冲区非空时接收不阻塞。funcmain(){// 创建一个容量为 2 的缓冲 Channelch:make(chanint,2)// 发送两个数据不会阻塞ch-1ch-2fmt.Println(两个数据已发送缓冲区已满)// 再发送第三个数据会阻塞缓冲区已满// ch - 3 // 死锁// 接收数据fmt.Println(-ch)// 1fmt.Println(-ch)// 2}2.3 关闭 Channel 与遍历发送方可以主动关闭 Channel接收方可以通过 v, ok : -ch 检测 Channel 是否已关闭。更常见的是使用 for range 自动遍历直到 Channel 关闭。funcproducer(chchan-int){// 只发送 Channelfori:0;i5;i{ch-i}close(ch)// 发送完成后关闭}funcmain(){ch:make(chanint,3)goproducer(ch)// 使用 range 遍历 Channel直到它被关闭forv:rangech{fmt.Println(接收到:,v)}fmt.Println(Channel 已关闭遍历结束)}⚠️ Channel 操作规范只能由发送方关闭 Channel接收方不应关闭。关闭已关闭的 Channel 会引发 panic。向已关闭的 Channel 发送数据会引发 panic。从已关闭的 Channel 接收数据会返回零值v, ok : -ch 中 ok 为 false。2.4 Channel 的方向只读/只写在函数签名中可以限制 Channel 的方向提高类型安全// 只发送 ChannelfuncsendOnly(chchan-int,valueint){ch-value// -ch // 编译错误不能从只发送 Channel 接收}// 只接收 ChannelfuncreceiveOnly(ch-chanint){value:-ch// ch - 1 // 编译错误不能向只接收 Channel 发送fmt.Println(value)}funcmain(){ch:make(chanint)gosendOnly(ch,42)receiveOnly(ch)}三、select 多路复用select 语句用于同时等待多个 Channel 操作类似于 switch但每个 case 是一个 Channel 操作funcmain(){ch1:make(chanstring)ch2:make(chanstring)gofunc(){time.Sleep(2*time.Second)ch1-来自 ch1 的消息}()gofunc(){time.Sleep(1*time.Second)ch2-来自 ch2 的消息}()// 等待最先到达的消息select{casemsg1:-ch1:fmt.Println(msg1)casemsg2:-ch2:fmt.Println(msg2)case-time.After(3*time.Second):fmt.Println(超时 3 秒)}// 输出: 来自 ch2 的消息因为 ch2 先到达}3.1 select 的典型用途1超时控制funcdoWork()string{time.Sleep(2*time.Second)return工作完成}funcmain(){ch:make(chanstring)gofunc(){ch-doWork()}()select{caseresult:-ch:fmt.Println(result)case-time.After(1*time.Second):fmt.Println(超时工作未在 1 秒内完成)}// 输出: 超时工作未在 1 秒内完成}2非阻塞的 Channel 操作使用 default 分支实现非阻塞操作funcmain(){ch:make(chanint)select{casev:-ch:fmt.Println(接收到:,v)default:fmt.Println(没有数据可接收不阻塞)}select{casech-42:fmt.Println(发送成功)default:fmt.Println(缓冲区已满不阻塞)}}四、互斥锁保护共享资源虽然 Channel 是 Go 推荐的并发通信方式但某些场景下共享内存加锁仍然更合适如缓存、状态管理。sync.Mutex 提供了互斥锁sync.RWMutex 提供了读写锁允许并发读互斥写。typeCounterstruct{mu sync.Mutex valueint}func(c*Counter)Inc(){c.mu.Lock()c.valuec.mu.Unlock()}func(c*Counter)Value()int{c.mu.Lock()deferc.mu.Unlock()returnc.value}funcmain(){varwg sync.WaitGroup counter:Counter{}fori:0;i1000;i{wg.Add(1)gofunc(){deferwg.Done()counter.Inc()}()}wg.Wait()fmt.Println(counter.Value())// 1000安全}五、实战并发工作池Worker Pool工作池模式是一种经典的并发设计模式用于限制同时执行的并发任务数量避免资源耗尽packagemainimport(fmtsynctime)// 任务结构体typeTaskstruct{IDintDatastring}// 工作池typeWorkerPoolstruct{numWorkersinttaskQueuechanTask wg sync.WaitGroup}// 创建工作池funcNewWorkerPool(numWorkersint,queueSizeint)*WorkerPool{returnWorkerPool{numWorkers:numWorkers,taskQueue:make(chanTask,queueSize),}}// 启动工作池func(wp*WorkerPool)Start(){fori:0;iwp.numWorkers;i{wp.wg.Add(1)gowp.worker(i)}}// 工作 Goroutinefunc(wp*WorkerPool)worker(idint){deferwp.wg.Done()fortask:rangewp.taskQueue{fmt.Printf(Worker %d 开始处理任务 %d: %s\n,id,task.ID,task.Data)time.Sleep(100*time.Millisecond)// 模拟处理时间fmt.Printf(Worker %d 完成任务 %d\n,id,task.ID)}}// 提交任务func(wp*WorkerPool)Submit(task Task){wp.taskQueue-task}// 关闭工作池等待所有任务完成func(wp*WorkerPool)Stop(){close(wp.taskQueue)wp.wg.Wait()}funcmain(){// 创建 3 个工作 Goroutine队列容量 10pool:NewWorkerPool(3,10)pool.Start()// 提交 10 个任务fori:0;i10;i{pool.Submit(Task{ID:i,Data:fmt.Sprintf(data-%d,i)})}// 等待所有任务完成pool.Stop()fmt.Println(所有任务已完成)}六、小结Goroutinego 关键字启动轻量级协程栈大小仅 2KBGMP 调度模型高效管理。Channel类型化的通信管道无缓冲同步和有缓冲异步两种类型。select多路复用 Channel 操作支持超时控制和默认分支非阻塞。WaitGroup等待一组 Goroutine 完成Add、Done、Wait 三方法组合使用。Mutexsync.Mutex 互斥锁sync.RWMutex 读写锁保护共享资源。