如何实现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 }
关键说明:
- 固定工作协程:程序启动时就创建2个协程,它们不会退出(除非任务通道被关闭),处理完一个任务后立刻等待下一个,实现无缝补位。
- 任务通道:无缓冲通道保证了任务提交会被阻塞,直到有空闲的工作协程接收,自然限制了并发数。
- 错误处理:把
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
相关产品推荐
相关产品推荐

