如何用Go高效将4GB大型CSV数据导入MySQL?求优化代码方案
如何高效将4GB大型CSV导入MySQL(Go语言实现方案)
嘿,我之前刚好处理过类似的大CSV导入MySQL的需求,你的代码速度慢主要是因为逐行处理+单条插入的开销太大——频繁的IO操作和数据库单条提交会严重拖慢效率。下面我来拆解优化思路,再给你一套完整的实现代码。
一、你的现有代码的性能瓶颈
- 逐行读取CSV并单条发送到channel,没有利用批量读取的优势,磁盘IO开销大
- 如果后续是单条插入数据库,会导致大量的数据库连接请求和事务提交,这是性能慢的核心原因
csv.NewReader默认的缓冲较小,读取大文件时会频繁触发磁盘读取
二、核心优化思路
- CSV读取层面:增大读取缓冲,批量读取行数据,减少磁盘IO次数;跳过CSV表头(如果有的话)避免无效解析
- 数据库操作层面:用批量插入+事务包裹减少提交次数;预编译SQL语句复用执行计划;合理配置数据库连接池
- 并发层面:采用生产者-消费者模式,多个goroutine并行处理插入任务,充分利用CPU和数据库资源
三、完整代码实现
1. 定义结构体与依赖
假设你的CSV对应User结构体,这里用csvutil库(性能比常见的gocsv更优)来做CSV解析,你可以先安装它:
go get github.com/jszwec/csvutil
结构体定义:
package main import ( "bufio" "database/sql" "fmt" "io" "log" "os" "sync" "time" "github.com/jszwec/csvutil" _ "github.com/go-sql-driver/mysql" ) type User struct { ID string `csv:"id"` Name string `csv:"name"` Email string `csv:"email"` // 按需添加你的其他字段 }
2. 优化后的CSV批量读取函数
这个函数会批量读取CSV行,达到指定批量大小后发送到channel,避免逐行处理的开销:
func usersFileLoader(filename string, batchSize int, channel chan []User) { defer close(channel) file, err := os.Open(filename) if err != nil { log.Fatalf("打开文件失败: %v", err) } defer file.Close() // 增大读取缓冲到1MB,提升大文件读取速度 reader := csvutil.NewDecoder(bufio.NewReaderSize(file, 1024*1024)) // 跳过CSV表头(如果你的文件有表头的话,没有就注释掉这行) var header []string if err := reader.DecodeHeader(&header); err != nil && err != io.EOF { log.Fatalf("解析表头失败: %v", err) } var batch []User for { var user User err := reader.Decode(&user) if err == io.EOF { // 发送最后一批剩余数据 if len(batch) > 0 { channel <- batch } break } if err != nil { log.Printf("解析行失败,跳过该行: %v", err) continue } batch = append(batch, user) // 达到批量大小就发送到channel if len(batch) >= batchSize { channel <- batch batch = nil // 重置批量切片 } } }
3. 数据库批量插入函数
用事务包裹批量插入,预编译SQL语句,减少数据库开销:
func batchInsertUsers(db *sql.DB, users []User) error { // 开启事务 tx, err := db.Begin() if err != nil { return fmt.Errorf("开启事务失败: %v", err) } defer func() { if r := recover(); r != nil { _ = tx.Rollback() } }() // 预编译插入语句,复用执行计划提升性能 stmt, err := tx.Prepare(`INSERT INTO users (id, name, email) VALUES (?, ?, ?)`) if err != nil { _ = tx.Rollback() return fmt.Errorf("预编译语句失败: %v", err) } defer stmt.Close() // 批量执行插入 for _, user := range users { _, err := stmt.Exec(user.ID, user.Name, user.Email) if err != nil { _ = tx.Rollback() return fmt.Errorf("插入记录失败: %v", err) } } // 提交事务 if err := tx.Commit(); err != nil { return fmt.Errorf("提交事务失败: %v", err) } return nil }
4. 主函数:整合生产者与消费者
配置数据库连接池,启动多个消费者goroutine并行处理插入:
func main() { // 初始化MySQL连接,替换成你的数据库信息 dbDSN := "username:password@tcp(127.0.0.1:3306)/your_db?parseTime=true&charset=utf8mb4" db, err := sql.Open("mysql", dbDSN) if err != nil { log.Fatalf("连接数据库失败: %v", err) } defer db.Close() // 配置连接池参数,根据你的服务器性能调整 db.SetMaxOpenConns(20) // 最大打开连接数 db.SetMaxIdleConns(10) // 最大空闲连接数 db.SetConnMaxLifetime(time.Minute * 3) // 连接最大存活时间 batchSize := 1000 // 每批处理的记录数,建议1000-5000之间 workerCount := 3 // 消费者goroutine数量,不要超过连接池最大连接数 userChan := make(chan []User, 5) // 带缓冲的channel,避免生产者阻塞 // 启动生产者:读取CSV并发送批量数据 go usersFileLoader("large_users.csv", batchSize, userChan) // 启动消费者:并行处理批量插入 var wg sync.WaitGroup for i := 0; i < workerCount; i++ { wg.Add(1) go func(workerID int) { defer wg.Done() for batch := range userChan { log.Printf("Worker %d 开始处理 %d 条记录", workerID, len(batch)) if err := batchInsertUsers(db, batch); err != nil { log.Printf("Worker %d 批量插入失败: %v", workerID, err) } else { log.Printf("Worker %d 成功插入 %d 条记录", workerID, len(batch)) } } }(i+1) } wg.Wait() log.Println("所有CSV数据已成功导入MySQL!") }
四、额外的性能优化建议
- 数据库端优化:
- 导入前执行
SET UNIQUE_CHECKS=0; SET FOREIGN_KEY_CHECKS=0;关闭唯一键和外键检查,导入完成后再设回1 - 增大
innodb_buffer_pool_size(建议设为服务器内存的50%-70%),让更多数据在内存中处理 - 关闭MySQL自动提交(我们的代码用了事务,所以已经避免了这个问题)
- 导入前执行
- 内存控制:如果服务器内存有限,可以适当减小
batchSize,避免内存占用过高 - 错误处理:生产环境可以记录错误行的内容和行号,方便后续排查问题
- CSV解析库:
csvutil的性能比很多其他CSV库更优,如果你用的是自己实现的Unmarshal,建议替换成它提升解析速度
内容的提问来源于stack exchange,提问作者Biplav Pokharel
相关产品推荐
相关产品推荐

