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

Go实现数据分批并发更新SQL类数据库的方案

大型数据集的SQL并发批量更新方案

当需要更新SQL数据库中的大型数据切片时,单次执行UPDATE tablename SET uid=12345 WHERE url IN (...)这类语句容易引发锁表超时、数据库性能骤降等问题。通过拆分批次+并发执行小批量更新可以有效规避这些问题,以下是具体实现思路和优化后的代码示例:

核心思路

  1. 拆分大切片为小批次:将包含大量URL的列表拆分成固定大小的子列表(比如每100条一个批次),控制单条UPDATE语句的处理量。
  2. 工作池控制并发:通过Worker Pool限制并发更新的数量,避免并发过高压垮数据库连接池或导致数据库负载过高。
  3. 批量更新优先:每个工作任务处理一个批次的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 19:35:35