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 ./...

测试通过不一定证明没有并发问题,但可以把最常见的竞态提前暴露。