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

如何实现Go语言最多两个协程持续无限并发执行,完成一个即启动下一个

解决Go协程并发补位的问题

你遇到的核心问题是原代码用了WaitGroup.Wait()强制等待所有协程完成才进入下一轮,没法实现“一个协程完成就立刻启动下一个”的动态补位效果。下面给你两种可行的解决方案,都能满足“同一时间最多2个协程运行,无限持续执行”的需求。

方案一:固定工作池模式(推荐,复用协程)

这种方式会提前创建好2个长期运行的工作协程,它们会持续监听任务通道,处理完一个任务后立刻接收下一个,完美实现补位逻辑:

package main

import (
	"sync"
	"time"

	"github.com/jmoiron/sqlx"
	// 根据你的数据库导入对应的驱动,比如PostgreSQL的_ "github.com/lib/pq"
)

const MAX_WORKERS = 2

func main() {
	// 初始化数据库连接(这里替换成你的DSN)
	db, err := sqlx.Connect("postgres", "host=localhost user=postgres dbname=test password=secret sslmode=disable")
	if err != nil {
		panic(err)
	}
	defer db.Close()

	// 创建无缓冲的任务通道,用来传递任务信号
	taskChan := make(chan struct{})

	// 启动固定数量的工作协程
	var workerWG sync.WaitGroup
	workerWG.Add(MAX_WORKERS)
	for workerID := 0; workerID < MAX_WORKERS; workerID++ {
		go func(id int) {
			defer workerWG.Done()
			// 持续监听任务通道,直到通道被关闭
			for {
				_, taskExists := <-taskChan
				if !taskExists {
					// 通道关闭,工作协程退出
					return
				}
				// 执行耗时的数据库操作
				if err := doDatabaseTask(db); err != nil {
					// 建议用日志代替panic,避免单个任务失败导致整个程序崩溃
					println("Worker", id, "执行任务失败:", err.Error())
				}
			}
		}(workerID)
	}

	// 无限往任务通道发送任务
	for {
		taskChan <- struct{}{}
		// 如果需要控制任务提交的频率,可以取消下面的注释
		// time.Sleep(500 * time.Millisecond)
	}

	// 若需要优雅退出,可在合适时机关闭通道并等待工作协程结束
	// close(taskChan)
	// workerWG.Wait()
}

func doDatabaseTask(db *sqlx.DB) error {
	// 模拟2-10秒的数据库调用(替换成你的实际业务SQL)
	sleepDuration := 2 + time.Duration(time.Now().UnixNano()%8)*time.Second
	_, err := db.Exec("SELECT pg_sleep($1)", sleepDuration.Seconds())
	return err
}

关键说明:

  1. 固定工作协程:程序启动时就创建2个协程,它们不会退出(除非任务通道被关闭),处理完一个任务后立刻等待下一个,实现无缝补位。
  2. 任务通道:无缓冲通道保证了任务提交会被阻塞,直到有空闲的工作协程接收,自然限制了并发数。
  3. 错误处理:把panic改成日志打印,避免单个任务失败导致整个程序挂掉,你可以根据业务需求调整错误处理逻辑。

方案二:信号量模式(代码更简洁)

如果你的任务逻辑比较简单,也可以用信号量来控制并发数,每次启动新协程前先获取信号量,完成后释放,这样也能保证同一时间最多2个协程运行:

package main

import (
	"time"

	"github.com/jmoiron/sqlx"
	// 导入数据库驱动
)

const MAX_CONCURRENT = 2

func main() {
	db, err := sqlx.Connect("postgres", "your-dsn-here")
	if err != nil {
		panic(err)
	}
	defer db.Close()

	// 用带缓冲的通道作为信号量,缓冲大小就是最大并发数
	semaphore := make(chan struct{}, MAX_CONCURRENT)

	// 无限提交任务
	for {
		// 获取信号量:如果通道已满,会阻塞直到有协程释放信号量
		semaphore <- struct{}{}
		// 启动新协程执行任务
		go func() {
			// 任务完成后释放信号量
			defer func() { <-semaphore }()
			if err := doDatabaseTask(db); err != nil {
				println("任务执行失败:", err.Error())
			}
		}()
		// 可选:控制任务提交频率
		// time.Sleep(300 * time.Millisecond)
	}
}

func doDatabaseTask(db *sqlx.DB) error {
	sleepDuration := 2 + time.Duration(time.Now().UnixNano()%8)*time.Second
	_, err := db.Exec("SELECT pg_sleep($1)", sleepDuration.Seconds())
	return err
}

关键说明:

  • 信号量通道的缓冲大小限制了同时运行的协程数量,每次启动协程前先往通道塞一个空结构体,满了就阻塞;协程完成后从通道取出一个,释放位置。
  • 这种方式不需要提前创建工作协程,代码更简洁,但协程会随任务创建和销毁,适合任务逻辑简单、性能要求不极端的场景。

对比原代码的改进点:

  • 去掉了每次循环重建通道和协程的逻辑,减少资源开销。
  • 不再依赖WaitGroup.Wait()阻塞整个循环,实现了“一个协程完成就立刻启动下一个”的动态补位。
  • 两种方案都能稳定控制最大并发数为2,同时无限持续执行任务。

内容的提问来源于stack exchange,提问作者Nicole Staline

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 22:17:40