
本文介绍如何通过 goroutine 并发模型(Extract/Transform/Load 分阶段协程化)结合 bufio.Scanner 和文件偏移控制,提升大型 CSV 文件的解析与入库性能。
本文介绍如何通过 goroutine 并发模型(extract/transform/load 分阶段协程化)结合 `bufio.scanner` 和文件偏移控制,提升大型 csv 文件的解析与入库性能。
在处理大型 CSV 文件(如每行约 300 字节、总大小达 GB 级)时,单线程顺序读取 + 解析 + 写入数据库往往成为性能瓶颈。Go 语言天然支持轻量级并发,合理使用 goroutines 可显著提升 ETL 整体吞吐量——但关键在于职责分离与流水线协同,而非盲目增加 goroutine 数量。
推荐采用三阶段流水线式并发设计:
Miller (mlr) 是一个命令行工具,用于查询、整形和重新格式化名称索引数据,如 CSV、TSV、JSON 和 JSON Lines。它将 awk、sed、cut、join 和 sort 的功能整合到一个专为结构化数据处理而构建的单一工具中。
- Extract(E):单个 goroutine 负责顺序、带缓冲地读取文件(使用 bufio.Scanner 或 bufio.Reader),按行产出字符串,通过 channel(如 chan string)发送给 Transform 阶段;
- Transform(T):启动多个 goroutine(例如 runtime.NumCPU() 个),从 channel 接收原始行,执行 strings.Split、字段校验、结构体映射、JSON 序列化等 CPU 密集操作,再将结果(如 []byte 或自定义 struct)发送至 Load 阶段;
- Load(L):1–2 个 goroutine 负责批量插入数据库(如使用 sql.DB.Prepare() + stmt.Exec() 批量提交),避免高频单条事务开销。
// 示例:简化版 ETL 流水线(省略错误处理)
func runETLPipeline(filename string) {
lines := make(chan string, 1024)
records := make(chan map[string]interface{}, 1024)
done := make(chan struct{})
// Extract: 单 goroutine 读取
go func() {
f, _ := os.Open(filename)
defer f.Close()
scanner := bufio.NewScanner(f)
for scanner.Scan() {
lines <- scanner.Text()
}
close(lines)
}()
// Transform: 多 goroutine 并行解析
for i := 0; i < runtime.NumCPU(); i++ {
go func() {
for line := range lines {
fields := strings.Split(line, ",")
record := map[string]interface{}{
"id": fields[0],
"name": fields[1],
"age": fields[2],
}
records <- record
}
}()
}
go func() { close(records) }() // 所有 T 完成后关闭 records
// Load: 单 goroutine 批量写入(可扩展为带 buffer 的批量器)
go func() {
db, _ := sql.Open("sqlite3", "./etl.db")
stmt, _ := db.Prepare("INSERT INTO users(id,name,age) VALUES(?,?,?)")
defer stmt.Close()
for record := range records {
stmt.Exec(record["id"], record["name"], record["age"])
}
close(done)
}()
<-done
}⚠️ 注意事项:
- bufio.Scanner 不支持随机位置读取(即无法“从第 N 字节开始扫描”),因其内部依赖逐字节状态机判断换行符;若需分片并行读取(如多 goroutine 各自处理文件某一段),应改用 os.Seek() + bufio.Reader.ReadBytes('\n') 手动实现行边界对齐,并确保起始偏移落在完整行开头(需向前搜索最近的 \n);
- channel 缓冲区大小需权衡内存占用与背压控制,建议设为 1024–4096;
- Transform 阶段若涉及复杂计算或外部调用,可进一步引入 errgroup.Group 统一错误传播;
- 实际生产中应添加日志、指标(如每秒处理行数)、超时控制及优雅退出逻辑。
总结:并发 ETL 的核心不是“越多 goroutine 越快”,而是让 I/O、CPU、DB 等不同资源瓶颈环节并行运转,形成稳定高效的流水线。合理使用缓冲通道与固定数量的工作 goroutine,配合 Go 原生调度器,即可在普通服务器上轻松应对百万级 CSV 行的高效处理。

















