Golang多生产者单消费者公平作业调度方案设计咨询
调度器设计思路
核心结构拆分
- 生产者信息管理器:给每个生产者单独维护一个作业缓冲区(最多存1个未消费作业),同时记录该生产者的权重、当前剩余调度配额,用互斥锁保证并发安全。
- 调度逻辑核心:采用加权轮询算法,按权重给生产者分配调度次数,权重高的生产者每轮能获得更多执行机会;如果生产者修改了作业,会重置它的配额,确保本轮优先被调度。
- 消费者协程:独立运行,循环从调度器获取下一个待处理作业,交给业务逻辑处理。
- 简洁API:用Go风格的结构体方法封装所有操作,生产者只需要调用
SubmitJob/CancelJob等方法,不用关心内部同步逻辑。
关键需求落地
- 异步协程模式:生产者完全独立运行在自己的goroutine里,调度器内部用互斥锁和channel处理并发,不用生产者手动做同步。
- 作业修改与位置保留:每个生产者的作业存在专属缓冲区,替换/取消作业直接操作这个缓冲区;如果修改后错过了当前轮次,调度器会重置该生产者的配额,让它在下一轮优先被选中,保证本轮提交权。
- 公平与权重支持:加权轮询算法既保证权重高的生产者有更多执行机会,又让同权重的生产者轮流调度,避免单个生产者抢占资源。
- Go风格接口:用
NewScheduler创建实例,RegisterProducer注册生产者时指定权重,Start启动消费者,Stop优雅关闭,所有操作都是直观的方法调用,符合Go的惯用法。
代码示例
package main import ( "sync" "time" ) // Job 定义作业对象,可根据业务扩展字段 type Job struct { ID string Payload interface{} } // ProducerID 生产者唯一标识类型 type ProducerID string // Scheduler 调度器核心结构体 type Scheduler struct { mu sync.Mutex producers map[ProducerID]*producerInfo quitChan chan struct{} wg sync.WaitGroup } // producerInfo 存储单个生产者的调度信息 type producerInfo struct { weight int // 生产者权重 currentQuota int // 当前轮次剩余调度配额 jobBuffer *Job // 未被消费的作业缓冲区 } // NewScheduler 创建调度器实例 func NewScheduler() *Scheduler { return &Scheduler{ producers: make(map[ProducerID]*producerInfo), quitChan: make(chan struct{}), } } // RegisterProducer 注册生产者,指定权重 func (s *Scheduler) RegisterProducer(id ProducerID, weight int) { s.mu.Lock() defer s.mu.Unlock() s.producers[id] = &producerInfo{ weight: weight, currentQuota: weight, jobBuffer: nil, } } // SubmitJob 提交/替换未消费的作业,提交后重置本轮配额保证优先调度 func (s *Scheduler) SubmitJob(id ProducerID, job Job) { s.mu.Lock() defer s.mu.Unlock() prod, ok := s.producers[id] if !ok { return } prod.jobBuffer = &job prod.currentQuota = prod.weight // 重置配额,确保本轮优先被选中 } // CancelJob 取消未被消费的作业 func (s *Scheduler) CancelJob(id ProducerID) { s.mu.Lock() defer s.mu.Unlock() prod, ok := s.producers[id] if !ok { return } prod.jobBuffer = nil } // Start 启动消费者协程,传入作业处理函数 func (s *Scheduler) Start(handler func(Job)) { s.wg.Add(1) go func() { defer s.wg.Done() for { select { case <-s.quitChan: return default: if job := s.nextJob(); job != nil { handler(*job) } else { time.Sleep(10 * time.Millisecond) // 无作业时短暂休眠,避免空轮询 } } } }() } // Stop 停止调度器,等待消费者处理完剩余作业 func (s *Scheduler) Stop() { close(s.quitChan) s.wg.Wait() } // nextJob 按照加权轮询规则选取下一个待处理作业 func (s *Scheduler) nextJob() *Job { s.mu.Lock() defer s.mu.Unlock() for { hasAvailable := false for _, prod := range s.producers { if prod.currentQuota <= 0 { prod.currentQuota = prod.weight // 配额耗尽,重置进入下一轮 continue } if prod.jobBuffer != nil { // 取出作业,减少当前轮次配额 job := prod.jobBuffer prod.jobBuffer = nil prod.currentQuota-- return job } // 当前生产者无作业,消耗一个配额后继续检查下一个 prod.currentQuota-- hasAvailable = true } if !hasAvailable { return nil // 所有生产者都无作业 } } } // 示例使用 func main() { scheduler := NewScheduler() // 注册三个生产者,权重分别为3、2、1 scheduler.RegisterProducer("producer-1", 3) scheduler.RegisterProducer("producer-2", 2) scheduler.RegisterProducer("producer-3", 1) // 启动消费者,打印作业ID scheduler.Start(func(job Job) { println("处理作业:", job.ID) time.Sleep(500 * time.Millisecond) // 模拟作业处理耗时 }) var wg sync.WaitGroup wg.Add(3) // 生产者1:提交5个作业,中途替换第3个 go func() { defer wg.Done() for i := 0; i < 5; i++ { scheduler.SubmitJob("producer-1", Job{ID: "p1-job-" + string(rune('0'+i))}) time.Sleep(1 * time.Second) if i == 2 { scheduler.SubmitJob("producer-1", Job{ID: "p1-job-2-updated"}) } } }() // 生产者2:提交3个作业 go func() { defer wg.Done() for i := 0; i < 3; i++ { scheduler.SubmitJob("producer-2", Job{ID: "p2-job-" + string(rune('0'+i))}) time.Sleep(1.5 * time.Second) } }() // 生产者3:提交2个作业 go func() { defer wg.Done() for i := 0; i < 2; i++ { scheduler.SubmitJob("producer-3", Job{ID: "p3-job-" + string(rune('0'+i))}) time.Sleep(2 * time.Second) } }() wg.Wait() // 等待所有生产者提交完成 time.Sleep(5 * time.Second) // 等待剩余作业处理完毕 scheduler.Stop() }
对初步方案的分析
你提到的令牌桶方案确实存在明显局限性:
- 性能瓶颈:生产者数量较多时,令牌竞争会导致大量协程阻塞,每次提交都要等待令牌,降低并发效率。
- 灵活性不足:难以实现作业的替换/移除逻辑,也无法直接支持权重优先级调度。
- 不符合Go风格:Go更倾向于通过结构体封装状态和同步逻辑,让调用者只关注业务操作,而非手动管理令牌这类底层同步机制,你的方案会增加生产者的复杂度。
内容的提问来源于stack exchange,提问作者Roman Maltsev
相关产品推荐
相关产品推荐

