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

Golang多生产者单消费者公平作业调度方案设计咨询

调度器设计思路

核心结构拆分

  • 生产者信息管理器:给每个生产者单独维护一个作业缓冲区(最多存1个未消费作业),同时记录该生产者的权重、当前剩余调度配额,用互斥锁保证并发安全。
  • 调度逻辑核心:采用加权轮询算法,按权重给生产者分配调度次数,权重高的生产者每轮能获得更多执行机会;如果生产者修改了作业,会重置它的配额,确保本轮优先被调度。
  • 消费者协程:独立运行,循环从调度器获取下一个待处理作业,交给业务逻辑处理。
  • 简洁API:用Go风格的结构体方法封装所有操作,生产者只需要调用SubmitJob/CancelJob等方法,不用关心内部同步逻辑。

关键需求落地

  1. 异步协程模式:生产者完全独立运行在自己的goroutine里,调度器内部用互斥锁和channel处理并发,不用生产者手动做同步。
  2. 作业修改与位置保留:每个生产者的作业存在专属缓冲区,替换/取消作业直接操作这个缓冲区;如果修改后错过了当前轮次,调度器会重置该生产者的配额,让它在下一轮优先被选中,保证本轮提交权。
  3. 公平与权重支持:加权轮询算法既保证权重高的生产者有更多执行机会,又让同权重的生产者轮流调度,避免单个生产者抢占资源。
  4. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 17:45:28