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

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)
}

关键改动说明

  1. 全局上下文管理:使用mainCtx统一控制所有协程的停止,避免分散的信号处理逻辑。
  2. 生产者批次控制:通过双层循环实现批次生成,每批发送完成后等待队列清空,同时监听停止信号。
  3. 消费者剩余消息处理:收到停止信号后,通过goto跳转到剩余消息处理逻辑,确保队列中所有消息都被转发给工作协程。
  4. 完整协程等待:用一个WaitGroup管理生产者、消费者和所有工作协程,确保所有任务完成后程序才退出。

内容的提问来源于stack exchange,提问作者d-velop

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 17:02:47