并发
Go 通过 goroutine 和 channel 提供简洁高效的并发支持。goroutine 是轻量级线程,调度由 Go 运行时管理,可高效运行成千上万个。
| 概念 | 说明 |
|---|---|
| Goroutine | 并发执行单位,使用 go 关键字启动,非阻塞 |
| Channel | goroutine 之间通信机制,使用 chan 创建,通过 <- 发送和接收 |
| WaitGroup | 等待多个 goroutine 完成 |
| Mutex | 互斥锁,保护共享资源 |
goroutine
使用 go 关键字开启一个新的 goroutine 执行函数:
go 函数名(参数列表)
package main
import (
"fmt"
"time"
)
func sayHello() {
for i := 0; i < 5; i++ {
fmt.Println("Hello")
time.Sleep(100 * time.Millisecond)
}
}
func main() {
go sayHello() // 启动 goroutine
for i := 0; i < 5; i++ {
fmt.Println("Main")
time.Sleep(100 * time.Millisecond)
}
}
输出顺序不固定,因为两个 goroutine 并发执行,没有先后顺序。
主函数退出时,尚未执行完的 goroutine 会被直接结束,所以需要等待 goroutine 完成
WaitGroup 等待 goroutine
sync.WaitGroup 用于等待多个 goroutine 完成,Add 增加计数器,Done 完成时调用,Wait 阻塞等待所有完成。
package main
import (
"fmt"
"sync"
)
func worker(id int, wg *sync.WaitGroup) {
defer wg.Done() // goroutine 完成时调用 Done()
fmt.Printf("Worker %d started\n", id)
fmt.Printf("Worker %d finished\n", id)
}
func main() {
var wg sync.WaitGroup
for i := 1; i <= 3; i++ {
wg.Add(1) // 增加计数器
go worker(i, &wg)
}
wg.Wait() // 等待所有 goroutine 完成
fmt.Println("All workers done")
}
通道(Channel)
通道用于 goroutine 之间的数据传递,使用 make 创建,使用 <- 操作符收发:
ch := make(chan int) // 创建通道,必须先创建才能使用
ch <- v // 把 v 发送到通道 ch
v := <-ch // 从 ch 接收数据并赋给 v
默认情况下通道是不带缓冲区的,发送端发送数据时必须同时有接收端接收,否则会阻塞
两个 goroutine 计算数字之和:
package main
import "fmt"
func sum(s []int, c chan int) {
sum := 0
for _, v := range s {
sum += v
}
c <- sum // 把 sum 发送到通道 c
}
func main() {
s := []int{7, 2, 8, -9, 4, 0}
c := make(chan int)
go sum(s[:len(s)/2], c)
go sum(s[len(s)/2:], c)
x, y := <-c, <-c // 从通道 c 中接收
fmt.Println(x, y, x+y) // -5 17 12
}
带缓冲 channel
通过 make 的第二个参数指定缓冲区大小,发送端数据可以放进缓冲区等待接收端获取,不需要立即同步。
无缓冲通道是同步的,发送方阻塞直到接收方取值;有缓冲通道在缓冲区满之前不会阻塞发送方
package main
import "fmt"
func main() {
// 带缓冲通道,缓冲区大小为 2
ch := make(chan int, 2)
// 可以同时发送两个数据,不用立刻同步读取
ch <- 1
ch <- 2
fmt.Println(<-ch) // 1
fmt.Println(<-ch) // 2
}
通道遵循先进先出原则,缓冲区满时发送会阻塞,直到前面的值被接收
关闭 channel
通过 close() 关闭通道,配合 range 遍历通道数据。如果通道不关闭,range 会一直阻塞等待新数据。
package main
import "fmt"
func fibonacci(n int, c chan int) {
x, y := 0, 1
for i := 0; i < n; i++ {
c <- x
x, y = y, x+y
}
close(c) // 发送完数据后关闭通道
}
func main() {
c := make(chan int, 10)
go fibonacci(cap(c), c)
// range 遍历通道数据,通道关闭且读完后结束
for i := range c {
fmt.Println(i)
}
}
关闭通道不会丢失里面的数据,只是读取完数据后不会一直阻塞等待新数据
select 多路选择
select 使得一个 goroutine 可以等待多个 channel 操作,会阻塞直到某个 case 可以继续执行。
package main
import "fmt"
func fibonacci(c, quit chan int) {
x, y := 0, 1
for {
select {
case c <- x: // 可以发送数据
x, y = y, x+y
case <-quit: // 收到 quit 信号
fmt.Println("quit")
return
}
}
}
func main() {
c := make(chan int)
quit := make(chan int)
go func() {
for i := 0; i < 10; i++ {
fmt.Println(<-c)
}
quit <- 0
}()
fibonacci(c, quit)
}
输出结果为 0 到 34 的斐波那契数列后输出 quit。
Mutex 互斥锁
多个 goroutine 同时访问共享变量会产生数据竞争,使用 sync.Mutex 保护共享资源。
package main
import (
"fmt"
"sync"
)
type Counter struct {
mu sync.Mutex
value int
}
func (c *Counter) Inc() {
c.mu.Lock() // 加锁
defer c.mu.Unlock() // 解锁,写在 defer 中保证一定释放
c.value++
}
func (c *Counter) Get() int {
c.mu.Lock()
defer c.mu.Unlock()
return c.value
}
func main() {
var wg sync.WaitGroup
c := Counter{}
for i := 0; i < 100; i++ {
wg.Add(1)
go func() {
defer wg.Done()
c.Inc() // 100 个 goroutine 同时自增
}()
}
wg.Wait()
fmt.Println("count:", c.Get()) // count: 100
}
并发实例:worker pool
使用 channel 配合 WaitGroup 实现一个任务分发的工作池:
package main
import (
"fmt"
"sync"
"time"
)
// worker 从 jobs 通道取任务,结果发送到 results 通道
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(50 * time.Millisecond)
results <- job * 2
}
}
func main() {
jobs := make(chan int, 10)
results := make(chan int, 10)
var wg sync.WaitGroup
// 启动 3 个 worker
for w := 1; w <= 3; w++ {
wg.Add(1)
go worker(w, jobs, results, &wg)
}
// 发送 5 个任务
for j := 1; j <= 5; j++ {
jobs <- j
}
close(jobs) // 任务发送完毕,关闭 jobs,worker 的 range 会退出
wg.Wait() // 等待所有 worker 完成
close(results)
// 读取结果
for r := range results {
fmt.Println("结果:", r)
}
}
常见问题:死锁(所有 goroutine 都在等待但无数据可用,解决:避免无限等待、正确关闭通道);数据竞争(多个 goroutine 同时访问同一变量,解决:使用 Mutex 或 Channel 同步)