如何在缓冲读取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 } } }
关键细节说明
- 残留缓冲区
leftover:确保半条记录不会丢失,每次新读取的内容都会先和残留内容拼接,保证行的完整性。 - 标准库
csv.Reader:底层自动处理行边界、带换行的字段等特殊情况,比自己实现的分块逻辑更稳定,几乎不会踩坑。 - 提前解析批次大小:把
strconv.Atoi(args[10])移到循环外,避免重复解析浪费性能。
内容的提问来源于stack exchange,提问作者Gabriel
相关产品推荐
相关产品推荐

