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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:57:44