Go语言带信号关闭的生产者消费者模式问题咨询
基于Go语言的信号安全生产者消费者实现修复方案
针对你遇到的三个问题,以下是具体的修复思路和完整代码:
问题1:程序停止后消息丢失
原代码核心问题是收到停止信号后直接终止消费者的消息转发逻辑,未处理队列剩余消息,也没等待所有工作协程完成任务就退出。修复后确保:
- 消费者收到停止信号后,优先处理完队列中所有剩余消息
- 用统一的
WaitGroup等待生产者、消费者、所有工作协程全部完成
问题2:限制生产者批次生成消息
原生产者无限制持续发送消息,修复后实现:
- 每次生成固定10条消息
- 等待队列清空后,再生成下一批消息
- 等待过程中响应停止信号,避免阻塞无法退出
问题3:通知生产者停止生成新消息
原生产者没有停止机制,修复后:
- 给生产者传递全局上下文,实时监听停止信号
- 收到信号后立即停止生成新消息,并关闭消息队列,触发消费者收尾逻辑
完整修复代码
package main import ( "context" "fmt" "math/rand" "os" "os/signal" "sync" "syscall" "time" ) func main() { const nConsumers = 2 const batchSize = 10 // 消息队列,容量设为批次大小 in := make(chan int, batchSize) // 创建全局上下文,统一管理所有协程的停止信号 mainCtx, mainCancel := context.WithCancel(context.Background()) defer mainCancel() var wg sync.WaitGroup // 启动生产者 wg.Add(1) p := Producer{in: in, batchSize: batchSize} go p.Produce(mainCtx, &wg) // 启动消费者 c := Consumer{in: in, jobs: make(chan int, nConsumers)} wg.Add(1) go c.Consume(mainCtx, &wg) // 启动工作协程 for i := 1; i <= nConsumers; i++ { wg.Add(1) go c.Work(&wg, i) } // 监听系统终止信号 termChan := make(chan os.Signal, 1) signal.Notify(termChan, syscall.SIGINT, syscall.SIGTERM) <-termChan fmt.Println("\n收到停止信号,开始优雅关闭...") // 触发全局停止信号 mainCancel() // 等待所有协程完成任务 wg.Wait() fmt.Println("所有任务处理完成,程序退出") } type Producer struct { in chan int batchSize int } func (p Producer) Produce(ctx context.Context, wg *sync.WaitGroup) { defer wg.Done() defer close(p.in) // 生产者退出时关闭消息队列,告知消费者无新消息 task := 1 for { // 生成一批消息 for i := 0; i < p.batchSize; i++ { select { case p.in <- task: fmt.Printf("Send value %d\n", task) task++ case <-ctx.Done(): fmt.Println("生产者收到停止信号,停止生成新消息") return } time.Sleep(time.Millisecond * 500) } // 等待队列清空,同时响应停止信号 fmt.Println("等待队列清空,准备生成下一批消息") for { select { case <-ctx.Done(): fmt.Println("生产者收到停止信号,退出等待") return default: if len(p.in) == 0 { break } time.Sleep(time.Millisecond * 100) } } } } type Consumer struct { in chan int jobs chan int } func (c Consumer) Consume(ctx context.Context, wg *sync.WaitGroup) { defer wg.Done() defer close(c.jobs) // 所有消息处理完后关闭任务通道,告知工作协程退出 // 正常接收消息阶段 for { select { case job, ok := <-c.in: if !ok { fmt.Println("消息队列已关闭,消费者停止接收新消息") return } c.jobs <- job case <-ctx.Done(): fmt.Println("消费者收到停止信号,开始处理剩余消息") goto processRemaining } } // 处理剩余消息阶段 processRemaining: for job := range c.in { c.jobs <- job } fmt.Println("剩余消息处理完成,消费者退出") } func (c Consumer) Work(wg *sync.WaitGroup, i int) { defer wg.Done() for job := range c.jobs { fmt.Printf("Worker #%d start job %d\n", i, job) time.Sleep(time.Millisecond * time.Duration(3000+rand.Intn(3000))) fmt.Printf("Worker #%d finish job %d\n", i, job) } fmt.Printf("Worker #%d 无更多任务,退出\n", i) }
关键改动说明
- 全局上下文管理:使用
mainCtx统一控制所有协程的停止,避免分散的信号处理逻辑。 - 生产者批次控制:通过双层循环实现批次生成,每批发送完成后等待队列清空,同时监听停止信号。
- 消费者剩余消息处理:收到停止信号后,通过
goto跳转到剩余消息处理逻辑,确保队列中所有消息都被转发给工作协程。 - 完整协程等待:用一个
WaitGroup管理生产者、消费者和所有工作协程,确保所有任务完成后程序才退出。
内容的提问来源于stack exchange,提问作者d-velop
相关产品推荐
相关产品推荐

