Go语言中共享csv.Writer实现多goroutine异步写入CSV文件
多Goroutine下安全共享csv.Writer写入CSV文件
需求与问题
需要执行数千次API调用,每个调用完成后立即将结果写入同一个CSV文件,而非等待所有调用结束后批量写入。但csv.Writer并非并发安全的,直接在多goroutine中共享使用会导致数据错乱、写入异常等问题,需要实现安全的共享写入机制。
现有代码示例:
package main import ( "encoding/csv" "os" ) type Row struct { Field1 string Field2 string } func main () { file, _ := os.Create("file.csv") w := csv.NewWriter(file) // 生成待写入的行数据 var rowsToWrite []Row // 希望用goroutine并行处理写入,但不确定线程安全写法 for _, r := range rowsToWrite { go func(row, writer) { err := writeToFile(row, writer) if err != nil { // 错误处理 } }(r, w) } } func writeToFile(row Row, writer ???) error { // 希望用共享writer追加写入CSV,维护文件写入位置 if err := w.Write(row); err != nil { return err } return nil }
解决方案
方法一:使用互斥锁(sync.Mutex)
通过互斥锁保证同一时间只有一个goroutine操作csv.Writer,避免并发写入冲突。
修改后的完整代码:
package main import ( "encoding/csv" "os" "sync" ) type Row struct { Field1 string Field2 string } func main() { file, err := os.Create("file.csv") if err != nil { panic(err) } defer file.Close() w := csv.NewWriter(file) defer w.Flush() var mu sync.Mutex var rowsToWrite []Row // 假设这里已经填充了rowsToWrite数据 var wg sync.WaitGroup wg.Add(len(rowsToWrite)) for _, r := range rowsToWrite { row := r // 避免循环变量捕获问题 go func() { defer wg.Done() err := writeToFile(row, w, &mu) if err != nil { // 根据实际场景处理错误,比如打印日志 println("写入失败:", err.Error()) } }() } wg.Wait() // 等待所有goroutine完成 } func writeToFile(row Row, writer *csv.Writer, mu *sync.Mutex) error { // 写入前加锁 mu.Lock() defer mu.Unlock() // 将Row转换为csv.Writer需要的字符串切片 record := []string{row.Field1, row.Field2} if err := writer.Write(record); err != nil { return err } // 可以在这里调用Flush,也可以在主goroutine最后统一Flush writer.Flush() return writer.Error() }
关键注意点:
- 必须用
sync.Mutex包裹所有对csv.Writer的写操作(包括Write和Flush) - 使用
sync.WaitGroup等待所有goroutine完成,避免程序提前退出 - 循环中要复制循环变量
r到局部变量row,否则所有goroutine会捕获到同一个循环变量的地址,导致数据错误 - 记得在最后调用
Flush,确保缓冲区数据写入文件
方法二:使用通道(Channel)实现写入队列
创建一个专门的写入goroutine,所有其他goroutine将需要写入的行发送到通道中,由这个单独的goroutine负责调用csv.Writer写入。这种方式避免了锁竞争,更适合高并发场景。
完整代码示例:
package main import ( "encoding/csv" "os" "sync" ) type Row struct { Field1 string Field2 string } func main() { file, err := os.Create("file.csv") if err != nil { panic(err) } defer file.Close() w := csv.NewWriter(file) defer w.Flush() // 创建通道,缓存大小可以根据并发量调整 rowChan := make(chan Row, 100) doneChan := make(chan struct{}) // 启动专门的写入goroutine go func() { for row := range rowChan { record := []string{row.Field1, row.Field2} if err := w.Write(record); err != nil { println("写入失败:", err.Error()) } w.Flush() } doneChan <- struct{}{} }() var rowsToWrite []Row // 假设这里已经填充了rowsToWrite数据 var wg sync.WaitGroup wg.Add(len(rowsToWrite)) for _, r := range rowsToWrite { row := r go func() { defer wg.Done() // 将行发送到通道 rowChan <- row }() } wg.Wait() // 等待所有goroutine发送完成 close(rowChan) // 关闭通道,通知写入goroutine结束 <-doneChan // 等待写入goroutine完成所有写入 }
关键注意点:
- 单独的写入goroutine是唯一操作
csv.Writer的goroutine,天然保证并发安全 - 通道的缓存大小可以根据实际并发数调整,避免goroutine阻塞
- 必须在所有数据发送完成后关闭通道,否则写入goroutine会一直阻塞在
range循环 - 用
doneChan确保写入goroutine完成所有剩余写入后再退出程序
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

