You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Go Worker Pool停止触发空指针Panic及方法优化问询

Go Worker Pool 停止时的空指针Panic修复方案

问题根源分析

  1. 空指针来源:p.workers切片中存在nil元素,或是遍历workers时并发修改了切片(比如其他goroutine移除worker),导致遍历到无效的Worker指针。
  2. Recover未生效:原代码的recover可能仅覆盖了worker的任务处理goroutine,而Pool.Stop()方法中的循环逻辑未被recover包裹,Panic发生在未被保护的同步代码块中。
  3. 无法确认停止状态:原逻辑未等待worker完成退出流程,无法保证worker完全停止后再执行后续操作。

修复后的核心代码

Pool与Worker结构定义

import (
    "sync"
    "log"
)

type Worker struct {
    id       int
    taskChan chan func()
    wg       *sync.WaitGroup
    stopChan chan struct{}
}

type Pool struct {
    workers []*Worker
    mu      sync.Mutex
    wg      sync.WaitGroup
}

重写Worker.Stop()方法

func (w *Worker) Stop() {
    if w == nil {
        return // 提前拦截nil worker,避免空指针触发
    }
    select {
    case <-w.stopChan:
        // 已停止过,直接返回
        return
    default:
        close(w.stopChan)
        // 标记当前worker退出完成
        w.wg.Done()
    }
}

重写Pool.Stop()方法

func (p *Pool) Stop() {
    p.mu.Lock()
    defer p.mu.Unlock()

    // 遍历并停止所有worker,跳过nil元素
    for _, worker := range p.workers {
        if worker == nil {
            continue
        }
        worker.Stop()
    }

    // 等待所有worker完全退出
    p.wg.Wait()

    // 清空workers切片,避免后续操作访问无效指针
    p.workers = nil
    log.Println("所有worker已完全停止")
}

配套Worker启动逻辑

确保worker初始化时正确绑定WaitGroup与停止通道:

func NewWorker(id int, poolWg *sync.WaitGroup) *Worker {
    w := &Worker{
        id:       id,
        taskChan: make(chan func()),
        stopChan: make(chan struct{}),
        wg:       &sync.WaitGroup{},
    }
    poolWg.Add(1)
    w.wg.Add(1)
    go w.run()
    return w
}

func (w *Worker) run() {
    defer func() {
        if err := recover(); err != nil {
            log.Printf("worker %d 执行Panic: %v", w.id, err)
        }
        w.wg.Done()
    }()

    for {
        select {
        case task, ok := <-w.taskChan:
            if !ok {
                return
            }
            task()
        case <-w.stopChan:
            // 处理完剩余任务后退出
            for {
                select {
                case task, ok := <-w.taskChan:
                    if !ok {
                        return
                    }
                    task()
                default:
                    // 任务队列已空,退出goroutine
                    return
                }
            }
        }
    }
}

测试用例

import (
    "testing"
    "time"
)

func TestPoolStop(t *testing.T) {
    pool := &Pool{}
    // 创建3个worker并分配任务
    for i := 0; i < 3; i++ {
        worker := NewWorker(i, &pool.wg)
        pool.mu.Lock()
        pool.workers = append(pool.workers, worker)
        pool.mu.Unlock()
        // 发送耗时任务
        worker.taskChan <- func() {
            time.Sleep(100 * time.Millisecond)
        }
    }

    // 调用停止方法
    pool.Stop()

    // 验证停止状态
    if pool.workers != nil {
        t.Error("停止后workers切片未清空")
    }
}

关键修复点总结

  • 空指针防护:Pool.Stop()遍历跳过nil元素,Worker.Stop()开头判断自身是否为nil。
  • 并发安全:用互斥锁mu保护p.workers的遍历与修改,避免并发读写异常。
  • 停止状态确认:通过sync.WaitGroup跟踪所有worker生命周期,Pool.Stop()中调用Wait()确保全部worker退出。
  • Panic捕获:在worker的run()方法中添加defer recover(),避免任务Panic影响整个Pool。

内容的提问来源于stack exchange,提问作者alex

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.05 03:25:53