求每个Worker独立控制SQL连接的Golang Worker Pool示例代码
可靠的Golang Worker Pool示例(每个Worker独立GORM连接)
以下是适配你需求的实现方案,每个Worker会初始化自己的GORM连接,持续处理任务直到所有任务完成:
核心代码实现
1. 定义任务结构与Worker函数
package main import ( "fmt" "sync" "gorm.io/driver/mysql" "gorm.io/gorm" ) // ImportTask 定义导入任务的结构,根据你的实际数据调整字段 type ImportTask struct { ID int Data string } // worker 每个Worker会创建独立的GORM连接,循环处理任务直到通道关闭 func worker(id int, tasks <-chan ImportTask, wg *sync.WaitGroup) { defer wg.Done() // 初始化当前Worker专属的GORM连接 dsn := "your_user:your_password@tcp(127.0.0.1:3306)/your_db?charset=utf8mb4&parseTime=True&loc=Local" db, err := gorm.Open(mysql.Open(dsn), &gorm.Config{}) if err != nil { fmt.Printf("Worker %d: DB连接失败 - %v\n", id, err) return } // 获取底层sql.DB实例,可配置连接池参数(可选) sqlDB, err := db.DB() if err != nil { fmt.Printf("Worker %d: 获取sql.DB失败 - %v\n", id, err) return } defer sqlDB.Close() // Worker退出时关闭连接 // 循环处理任务,直到tasks通道被关闭 for task := range tasks { fmt.Printf("Worker %d: 处理任务ID %d\n", id, task.ID) // 执行你的导入逻辑,这里示例为插入数据 if err := db.Create(&task).Error; err != nil { fmt.Printf("Worker %d: 任务ID %d处理失败 - %v\n", id, task.ID, err) continue } } fmt.Printf("Worker %d: 所有任务处理完成\n", id) }
2. 主函数启动Worker池与推送任务
func main() { const ( numWorkers = 16 // 你需要的Worker数量 numTasks = 10000 // 总任务数 ) // 创建带缓冲的任务通道,缓冲大小可根据任务量调整,避免推送任务时阻塞 tasks := make(chan ImportTask, numTasks) var wg sync.WaitGroup // 启动所有Worker for i := 1; i <= numWorkers; i++ { wg.Add(1) go worker(i, tasks, &wg) } // 推送所有任务到通道 for i := 1; i <= numTasks; i++ { tasks <- ImportTask{ ID: i, Data: fmt.Sprintf("import_data_%d", i), } } close(tasks) // 所有任务推送完成后关闭通道,Worker会自动结束循环 // 等待所有Worker完成任务 wg.Wait() fmt.Println("所有导入任务处理完成") }
关键注意点
- 独立DB连接:每个Worker初始化自己的GORM连接,彻底避免共享连接导致的并发写入段错误问题。
- 通道关闭时机:必须在所有任务推送完成后再关闭
tasks通道——如果你之前过早关闭通道,会导致Worker提前停止接收任务,这正是你遇到的「只处理1个任务就停止」的原因。 - 连接池优化:可以通过
sqlDB.SetMaxOpenConns()、sqlDB.SetMaxIdleConns()等方法配置连接池参数,根据你的MySQL配置和任务量调整,避免连接过多导致数据库压力过大。 - 错误处理扩展:示例中仅做了基础打印,你可以结合日志库记录错误,或者将失败任务放入重试队列,提升导入的可靠性。
内容的提问来源于stack exchange,提问作者DudiDude
相关产品推荐
相关产品推荐

