实现 Worker Pool

🔴 困难

题目描述

实现一个 Worker Pool,支持动态调整 worker 数量。

参考答案

type WorkerPool struct {
    tasks    chan Task
    workers  int
    wg       sync.WaitGroup
    quit     chan struct{}
    mu       sync.RWMutex
}

type Task func()

func NewWorkerPool(workers int) *WorkerPool {
    return &WorkerPool{
        tasks:   make(chan Task, 100),
        workers: workers,
        quit:    make(chan struct{}),
    }
}

func (wp *WorkerPool) Start() {
    wp.mu.Lock()
    defer wp.mu.Unlock()
    
    for i := 0; i < wp.workers; i++ {
        wp.wg.Add(1)
        go wp.worker()
    }
}

func (wp *WorkerPool) worker() {
    defer wp.wg.Done()
    
    for {
        select {
        case task, ok := <-wp.tasks:
            if !ok {
                return
            }
            task()
        case <-wp.quit:
            return
        }
    }
}

func (wp *WorkerPool) Submit(task Task) {
    wp.tasks <- task
}

func (wp *WorkerPool) Stop() {
    close(wp.tasks)
    wp.wg.Wait()
}

func (wp *WorkerPool) Resize(n int) {
    wp.mu.Lock()
    defer wp.mu.Unlock()
    
    if n > wp.workers {
        // 增加 workers
        for i := 0; i < n-wp.workers; i++ {
            wp.wg.Add(1)
            go wp.worker()
        }
    } else if n < wp.workers {
        // 减少 workers
        for i := 0; i < wp.workers-n; i++ {
            wp.quit <- struct{}{}
        }
    }
    wp.workers = n
}

关键点

  1. 任务队列:缓冲 channel
  2. 动态调整:通过 quit channel 通知 worker 退出
  3. 优雅停止:关闭任务队列,等待所有 worker 完成