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

如何在缓冲读取CSV文件时避免记录被拆分?

解决CSV分块读取时记录被拆分的问题

这个问题我之前处理过好多次了——固定大小缓冲区读CSV最容易踩这个坑!核心问题就是你每次直接把整个buffer拿去解析,完全没考虑到缓冲区末尾可能是半条记录。下面给你一套可行的解决方案,直接改代码就能用:

核心思路

  • 维护一个残留字符串缓冲区,专门保存每次读取后不完整的行
  • 每次读取新块后,先和残留缓冲区的内容拼接
  • 找到拼接后内容的最后一个换行符,只解析换行符之前的完整内容
  • 把换行符之后的半行内容存回残留缓冲区,留到下一次处理
  • 当文件读到EOF时,再把残留缓冲区里剩下的最后一行(此时应该是完整的)解析掉

修改后的自定义分块处理代码

func insertData(args []string, csvFile *os.File, tableToInsert IndexRecord) {
    log.Printf("Table %v.%v found. Inserting data in database. Batches of %v", tableToInsert.Schema, tableToInsert.Name, args[10])
    buffer := make([]byte, 2048000)
    var leftover []byte // 保存不完整的行
    batchSize, err := strconv.Atoi(args[10])
    if err != nil {
        log.Fatalf("Invalid batch size: %v", err)
    }

    for {
        n, err := csvFile.Read(buffer)
        if n > 0 {
            // 拼接残留内容和新读取的字节
            combined := append(leftover, buffer[:n]...)
            // 找到最后一个换行符的位置
            lastNewline := bytes.LastIndex(combined, []byte("\n"))
            if lastNewline != -1 {
                // 解析完整的部分(到最后一个换行符)
                _, csvRecords := ParseCSVBytes(combined[:lastNewline+1])
                InsertRecords(tableToInsert, csvRecords, batchSize)
                // 保存剩余的不完整部分
                leftover = combined[lastNewline+1:]
            } else {
                // 这次读取的内容里没有完整的换行,全部存到残留区
                leftover = combined
            }
        }

        if err != nil {
            if err != io.EOF {
                fmt.Println("Error reading file:", err)
                break
            }
            // 处理文件末尾剩余的最后一行
            if len(leftover) > 0 {
                _, csvRecords := ParseCSVBytes(leftover)
                InsertRecords(tableToInsert, csvRecords, batchSize)
            }
            break
        }
    }
}

更可靠的标准库方案(推荐)

如果你的CSV存在带换行符的字段(比如引号包裹的多行内容),上面的简单换行符判断会失效。这种情况下直接用Go标准库的csv.Reader更省心——它本身支持按行读取,能自动处理所有CSV特殊规则,还可以通过BufSize控制缓冲区大小:

func insertData(args []string, csvFile *os.File, tableToInsert IndexRecord) {
    log.Printf("Table %v.%v found. Inserting data in database. Batches of %v", tableToInsert.Schema, tableToInsert.Name, args[10])
    batchSize, err := strconv.Atoi(args[10])
    if err != nil {
        log.Fatalf("Invalid batch size: %v", err)
    }

    reader := csv.NewReader(csvFile)
    reader.BufSize = 2048000 // 设置和你之前一致的缓冲区大小
    var batch [][]string

    for {
        record, err := reader.Read()
        if err != nil {
            if err != io.EOF {
                fmt.Println("Error reading CSV:", err)
            }
            // 处理最后一批剩余的记录
            if len(batch) > 0 {
                InsertRecords(tableToInsert, batch, batchSize)
            }
            break
        }
        batch = append(batch, record)
        // 当批次达到指定大小,插入数据库并清空批次
        if len(batch) >= batchSize {
            InsertRecords(tableToInsert, batch, batchSize)
            batch = nil
        }
    }
}

关键细节说明

  1. 残留缓冲区leftover:确保半条记录不会丢失,每次新读取的内容都会先和残留内容拼接,保证行的完整性。
  2. 标准库csv.Reader:底层自动处理行边界、带换行的字段等特殊情况,比自己实现的分块逻辑更稳定,几乎不会踩坑。
  3. 提前解析批次大小:把strconv.Atoi(args[10])移到循环外,避免重复解析浪费性能。

内容的提问来源于stack exchange,提问作者Gabriel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 15:42:50