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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 05:50:27