Go实现数据分批并发更新SQL类数据库的方案
大型数据集的SQL并发批量更新方案
当需要更新SQL数据库中的大型数据切片时,单次执行UPDATE tablename SET uid=12345 WHERE url IN (...)这类语句容易引发锁表超时、数据库性能骤降等问题。通过拆分批次+并发执行小批量更新可以有效规避这些问题,以下是具体实现思路和优化后的代码示例:
核心思路
- 拆分大切片为小批次:将包含大量URL的列表拆分成固定大小的子列表(比如每100条一个批次),控制单条UPDATE语句的处理量。
- 工作池控制并发:通过Worker Pool限制并发更新的数量,避免并发过高压垮数据库连接池或导致数据库负载过高。
- 批量更新优先:每个工作任务处理一个批次的URL,执行批量UPDATE,比单条更新效率提升显著。
优化后的Go实现代码
package main import ( "database/sql" "fmt" "sync" _ "github.com/go-sql-driver/mysql" ) // BatchTask 每个任务处理一个批次的URL更新 type BatchTask struct { URLs []string UID string } // WorkerPool 工作池,控制并发更新的数量 type WorkerPool struct { numWorkers int tasksChan chan BatchTask wg sync.WaitGroup } // NewWorkerPool 创建指定并发数的工作池 func NewWorkerPool(numWorkers int, taskBuffer int) *WorkerPool { return &WorkerPool{ numWorkers: numWorkers, tasksChan: make(chan BatchTask, taskBuffer), } } // Start 启动工作池,开始处理任务 func (wp *WorkerPool) Start(db *sql.DB) { for i := 0; i < wp.numWorkers; i++ { wp.wg.Add(1) go func(workerID int) { defer wp.wg.Done() for task := range wp.tasksChan { if err := executeBatchUpdate(db, task.UID, task.URLs); err != nil { fmt.Printf("Worker %d 更新失败: %v\n", workerID, err) // 可根据需求添加重试逻辑,注意幂等性 continue } fmt.Printf("Worker %d 完成批次更新,共处理 %d 条URL\n", workerID, len(task.URLs)) } }(i) } } // AddBatch 添加一个批次更新任务 func (wp *WorkerPool) AddBatch(urls []string, uid string) { wp.tasksChan <- BatchTask{ URLs: urls, UID: uid, } } // Wait 等待所有任务完成并关闭工作池 func (wp *WorkerPool) Wait() { close(wp.tasksChan) wp.wg.Wait() } // executeBatchUpdate 执行单批次SQL更新 func executeBatchUpdate(db *sql.DB, uid string, urls []string) error { if len(urls) == 0 { return nil } // 构建IN子句占位符,比如(?, ?, ?) placeholders := make([]string, len(urls)) args := make([]interface{}, len(urls)+1) args[0] = uid for i := range urls { placeholders[i] = "?" args[i+1] = urls[i] } query := fmt.Sprintf("UPDATE tablename SET uid = ? WHERE url IN (%s)", fmt.Sprintf("%s", placeholders)) _, err := db.Exec(query, args...) return err } // 示例调用 func main() { // 初始化数据库连接(根据实际配置修改) db, err := sql.Open("mysql", "user:password@tcp(127.0.0.1:3306)/dbname") if err != nil { panic(err) } defer db.Close() // 模拟大型URL切片 largeURLList := generateLargeURLList(1000) // 批次大小设置为100 batchSize := 100 // 创建并发数为3的工作池 pool := NewWorkerPool(3, 10) pool.Start(db) // 拆分大切片为小批次并添加到工作池 for i := 0; i < len(largeURLList); i += batchSize { end := i + batchSize if end > len(largeURLList) { end = len(largeURLList) } pool.AddBatch(largeURLList[i:end], "12345") } // 等待所有任务完成 pool.Wait() fmt.Println("所有批次更新完成") } // generateLargeURLList 生成模拟的大型URL列表 func generateLargeURLList(count int) []string { urls := make([]string, count) for i := 0; i < count; i++ { urls[i] = fmt.Sprintf("https://example.com/%d", i) } return urls }
关键注意事项
- 批次大小调整:根据数据库的
max_allowed_packet配置和实际性能,调整批次大小(一般50-200条较为合适),避免单条SQL语句过长导致报错。 - 并发数控制:工作池的并发数不要超过数据库连接池的最大连接数(可通过
db.SetMaxOpenConns()设置),防止连接耗尽。 - 错误处理与重试:添加错误捕获和重试逻辑时,必须确保UPDATE操作是幂等的(重复执行不会改变最终结果),避免数据异常。
- 事务与隔离级别:如果业务需要,可在批次更新中加入事务,但要注意事务时长,避免长时间锁表。
内容的提问来源于stack exchange,提问作者Johnny
相关产品推荐
相关产品推荐

