Go 很容易启动并发:
go work()
困难的是回答另外几个问题:它什么时候结束,失败由谁处理,重复启动怎么办,服务退出时谁等待它。
在采集、MQTT 转发、定时任务和 WebSocket 服务中,我遇到的大部分并发问题并不是算法复杂,而是生命周期没有被设计。
一、每个 goroutine 都应该有归属
启动 goroutine 的代码,应当知道:
- 为什么启动;
- 如何停止;
- 错误发给谁;
- 退出时是否需要等待。
最危险的是“顺手异步”:
go func() {
client.Publish(...)
}()
调用方不知道发布是否成功,也无法限制并发数量。流量升高后,goroutine 和内存可能一起增长。
二、用 context 传播取消信号
长期运行的工作循环应监听 context.Context:
func run(ctx context.Context) error {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
if err := collect(); err != nil {
return err
}
}
}
}
context 适合传递取消、截止时间和请求范围信息,不适合装载任意业务参数。
三、Channel 表达事件,Mutex 保护状态
Channel:
更适合“传消息 / 发信号 / 编排 goroutine”
1. 通知初始化完成
Channel 很适合用来表达“一件事情已经发生”。
比如服务启动时,一个 goroutine 负责初始化连接,另一个 goroutine 等初始化完成后再继续执行:
func main() {
ready := make(chan struct{})
go func() {
// 模拟初始化:连接数据库、加载配置等
initDB()
// 通知初始化完成
close(ready)
}()
// 等待初始化完成
<-ready
fmt.Println("服务开始处理请求")
}
func initDB() {
time.Sleep(time.Second)
}
这里使用 close(ready),不是为了传递某个具体数据,而是为了广播一个状态:
初始化已经完成,等待方可以继续往下走了。
为什么这里适合用 close?
因为关闭后的 channel 再接收不会阻塞。也就是说:
<-ready
在 ready 被关闭前会阻塞等待;一旦 ready 被关闭,就会立即继续执行。
这种写法常见于“只通知一次”的场景,例如:
配置加载完成;
数据库连接完成;
缓存预热完成;
服务可以开始接收请求。
2. 传递任务
Channel 也常用于任务队列。生产者往 channel 里写任务,多个 worker 从 channel 里取任务处理。
type Job struct {
ID int
Data string
}
func worker(jobs <-chan Job, wg *sync.WaitGroup) {
defer wg.Done()
for job := range jobs {
fmt.Println("处理任务:", job.ID, job.Data)
}
}
func main() {
jobs := make(chan Job, 10)
var wg sync.WaitGroup
for i := 0; i < 3; i++ {
wg.Add(1)
go worker(jobs, &wg)
}
for i := 0; i < 5; i++ {
jobs <- Job{
ID: i,
Data: "task",
}
}
// 通知 worker:后面不会再有新任务了
close(jobs)
// 等待所有 worker 把任务处理完并退出
wg.Wait()
}
这里有两个关键点。
第一,close(jobs) 不会清空 channel 里已有的任务。
即使 jobs 里还有缓存任务,关闭后 worker 仍然可以继续读取,直到把已有任务全部取完。
第二,for job := range jobs 的退出条件是:
jobs 已经关闭,并且 jobs 里的数据已经全部读完。
所以 close(jobs) 的含义不是“立刻停止 worker”,而是:
不会再有新任务了,你们把剩下的任务处理完后就可以退出。
不过还要注意一点:close(jobs) 只负责通知,不负责等待。
如果 main 函数在 close(jobs) 后直接结束,那么整个程序会退出,其他 goroutine 也会被一起结束。为了保证 worker 真正处理完任务,需要使用 sync.WaitGroup 等待它们退出。
3. 汇聚错误
Channel 还适合收集多个 goroutine 的执行结果。
比如同时启动三个任务,然后统一接收它们返回的错误:
func main() {
errCh := make(chan error, 3)
go func() {
errCh <- taskA()
}()
go func() {
errCh <- taskB()
}()
go func() {
errCh <- taskC()
}()
for i := 0; i < 3; i++ {
if err := <-errCh; err != nil {
fmt.Println("任务失败:", err)
return
}
}
fmt.Println("全部任务成功")
}
func taskA() error { return nil }
func taskB() error { return errors.New("taskB failed") }
func taskC() error { return nil }
这里的核心是:
err := <-errCh
这行代码会阻塞。
如果当前还没有 goroutine 往 errCh 里发送结果,主 goroutine 就会停在这里等待。直到某个任务执行完成,并执行:
errCh <- taskA()
主 goroutine 才能继续往下走。
但是这个循环只接收三次:
for i := 0; i < 3; i++ {
因为前面正好启动了三个任务,每个任务都会发送一次结果。所以它的执行过程可以理解为:
第一次接收:等待第一个完成的任务;
第二次接收:等待第二个完成的任务;
第三次接收:等待第三个完成的任务;
三次接收结束后,循环结束,不再阻塞。
也就是说,不是 <-errCh 后面不阻塞了,而是代码已经不再继续接收了。
需要注意的是,这个例子里只要收到一个错误就会 return。如果某个任务先失败,main 会提前结束,其他任务可能还没来得及完成。
如果希望等待所有任务都结束后,再统一判断是否有错误,可以这样写:
func main() {
errCh := make(chan error, 3)
go func() {
errCh <- taskA()
}()
go func() {
errCh <- taskB()
}()
go func() {
errCh <- taskC()
}()
var firstErr error
for i := 0; i < 3; i++ {
if err := <-errCh; err != nil && firstErr == nil {
firstErr = err
}
}
if firstErr != nil {
fmt.Println("任务失败:", firstErr)
return
}
fmt.Println("全部任务成功")
}
这个版本的语义更清楚:
三个任务都等完;
只要有一个失败,最后就认为整体失败。
4. 广播退出
Channel 也可以用来通知多个 goroutine 一起退出。
func worker(stop <-chan struct{}, id int, wg *sync.WaitGroup) {
defer wg.Done()
for {
select {
case <-stop:
fmt.Println("worker 退出:", id)
return
default:
fmt.Println("worker 工作中:", id)
time.Sleep(time.Second)
}
}
}
func main() {
stop := make(chan struct{})
var wg sync.WaitGroup
go worker(stop, 1, &wg)
go worker(stop, 2, &wg)
go worker(stop, 3, &wg)
time.Sleep(3 * time.Second)
// close 可以广播给所有接收者
close(stop)
// 等待所有 worker 收到退出信号并结束
wg.Wait()
}
这里的重点是:
close(stop)
关闭 stop 后,所有正在等待:
<-stop
的 goroutine 都会收到信号。
因为关闭后的 channel 可以被所有接收者读到,而且读操作会立即返回,所以它天然适合做广播通知。
不过要注意,close(stop) 只是发出退出信号,不等于 worker 已经退出。
如果 worker 此时正在执行:
time.Sleep(time.Second)
它要等这一轮 sleep 结束后,下一次进入 select,才会发现 stop 已经关闭,然后退出。
所以实际代码里不应该依赖:
time.Sleep(time.Second)
来等待 goroutine 结束,而应该使用 sync.WaitGroup。
实际项目中,更常见的退出控制方式是 context.Context:
ctx, cancel := context.WithCancel(context.Background())
go func() {
for {
select {
case <-ctx.Done():
fmt.Println("worker 退出")
return
default:
fmt.Println("worker 工作中")
time.Sleep(time.Second)
}
}
}()
cancel()
context.Context 的底层思想和关闭 channel 很像:都是广播一个取消信号。区别是 context 更适合工程化场景,因为它可以携带超时、截止时间和取消原因,也能在函数调用链之间传递。
5. 实现有界队列
比如最多只允许队列里积压 100 个任务。
type Task struct {
ID int
}
func main() {
queue := make(chan Task, 100)
// 消费者
go func() {
for task := range queue {
fmt.Println("处理任务:", task.ID)
time.Sleep(500 * time.Millisecond)
}
}()
// 生产者
for i := 0; i < 1000; i++ {
queue <- Task{ID: i}
fmt.Println("提交任务:", i)
}
}
make(chan Task, 100) 表示缓冲区最多放 100 个任务。
如果消费者处理太慢,生产者写到第 101 个时就会阻塞。这个特性可以自然实现:
限流
背压
任务排队
防止无限堆积
Mutex:
更适合“保护共享内存状态”
1. 保护共享 map
Go 的普通 map 不是并发安全的,多个 goroutine 同时读写会出问题。
type Cache struct {
mu sync.Mutex
data map[string]string
}
func NewCache() *Cache {
return &Cache{
data: make(map[string]string),
}
}
func (c *Cache) Set(key, value string) {
c.mu.Lock()
defer c.mu.Unlock()
c.data[key] = value
}
func (c *Cache) Get(key string) string {
c.mu.Lock()
defer c.mu.Unlock()
return c.data[key]
}
这里 mu 的作用是保证同一时间只有一个 goroutine 能访问 data。
如果读多写少,也可以用 sync.RWMutex:
type Cache struct {
mu sync.RWMutex
data map[string]string
}
只读 data:RLock / RUnlock
修改 data:Lock / Unlock
2. 原子地读取和更新一组字段
比如设备状态里有多个字段,必须一起更新,不能更新一半被别人读到。
type DeviceStatus struct {
mu sync.Mutex
Online bool
LastSeenAt time.Time
ErrorCount int
}
func (d *DeviceStatus) MarkOnline() {
d.mu.Lock()
defer d.mu.Unlock()
d.Online = true
d.LastSeenAt = time.Now()
d.ErrorCount = 0
}
func (d *DeviceStatus) MarkOffline() {
d.mu.Lock()
defer d.mu.Unlock()
d.Online = false
d.LastSeenAt = time.Now()
d.ErrorCount++
}
这里不能用三个独立变量随便改,因为这些字段之间有一致性要求。
例如:
Online = false
LastSeenAt = 当前时间
ErrorCount + 1
这三个动作要么一起完成,要么都还没开始,不能被外部看到中间状态。
3. 管理客户端连接状态
比如维护 MQTT、WebSocket、TCP 客户端是否在线。
type ClientManager struct {
mu sync.Mutex
clients map[string]bool
}
func NewClientManager() *ClientManager {
return &ClientManager{
clients: make(map[string]bool),
}
}
func (m *ClientManager) SetOnline(clientID string) {
m.mu.Lock()
defer m.mu.Unlock()
m.clients[clientID] = true
}
func (m *ClientManager) SetOffline(clientID string) {
m.mu.Lock()
defer m.mu.Unlock()
m.clients[clientID] = false
}
func (m *ClientManager) IsOnline(clientID string) bool {
m.mu.Lock()
defer m.mu.Unlock()
return m.clients[clientID]
}
这种场景本质是:多个 goroutine 可能同时修改连接状态。
例如:
连接成功回调:SetOnline
断开回调:SetOffline
心跳检测:IsOnline
重连逻辑:读取状态后决定是否重连
所以适合用 Mutex 保护共享状态。
4. 维护设备级锁映射
比如同一台设备不能同时执行两个控制命令,不同设备之间可以并发。
type DeviceLocker struct {
mu sync.Mutex
locks map[string]*sync.Mutex
}
func NewDeviceLocker() *DeviceLocker {
return &DeviceLocker{
locks: make(map[string]*sync.Mutex),
}
}
func (d *DeviceLocker) GetLock(deviceID string) *sync.Mutex {
d.mu.Lock()
defer d.mu.Unlock()
lock, ok := d.locks[deviceID]
if !ok {
lock = &sync.Mutex{}
d.locks[deviceID] = lock
}
return lock
}
使用方式:
func controlDevice(locker *DeviceLocker, deviceID string) {
lock := locker.GetLock(deviceID)
lock.Lock()
defer lock.Unlock()
// 同一设备的控制命令串行执行
sendCommand(deviceID)
}
func sendCommand(deviceID string) {
fmt.Println("控制设备:", deviceID)
}
这样可以做到:
设备 A 的多个命令串行
设备 B 的多个命令串行
设备 A 和设备 B 之间可以并发
四、限制并发,而不是相信流量不会变大
批量采集或转发时,可以使用有界 worker pool:
jobs := make(chan Job, 100)
for i := 0; i < workerCount; i++ {
go worker(ctx, jobs)
}
队列满以后要明确策略:阻塞、丢弃、覆盖旧数据,还是返回错误。没有策略的无界并发,本质上是把流量压力转成内存压力。
五、错误不能消失在 goroutine 里
g, ctx := errgroup.WithContext(ctx)
g.Go(func() error { return serveHTTP(ctx) })
g.Go(func() error { return runCollector(ctx) })
g.Go(func() error { return runForwarder(ctx) })
return g.Wait()
任一核心组件失败时,取消其他组件并统一退出,比留下半个服务继续运行更容易维护。
非核心任务则可以独立重试,但必须记录失败次数和最终状态。
tips
除了普通单元测试,还应使用:
go test -race ./...
测试通过不一定证明没有并发问题,但可以把最常见的竞态提前暴露。