Go协程同步问题:处理S3 CSV时阻塞与数据竞争
问题分析
你的代码核心问题有两个:
for result := range output永久阻塞:因为output通道从未被关闭,main goroutine会一直等待新数据。- 直接关闭
output导致panic:如果在读取完CSV就关output,此时worker可能还在往通道里写数据,触发"向已关闭通道写入"的panic。 - 潜在数据竞争:如果
lingua.LanguageDetector不是并发安全的,多个worker共享同一个实例会引发数据竞争。
修复方案:正确同步worker与通道关闭
使用sync.WaitGroup跟踪所有worker的执行状态,等所有worker处理完任务并退出后,再安全关闭output通道。
修改后的完整代码:
package main import ( "encoding/csv" "io" "log" "sync" "time" "github.com/pemistahl/lingua-go" ) type DetectionResult struct { // 你的结果结构体定义 } func ReadCSV(bucket, key string) (*http.Response, error) { // 你的S3读取实现 } func NewDetector(languages []lingua.Language) lingua.LanguageDetector { // 你的Detector初始化实现 } func DetectText(detector lingua.LanguageDetector, text string) DetectionResult { // 你的检测逻辑实现 } func main() { resp, err := ReadCSV(bucket, key) if err != nil { log.Fatal(err) } defer resp.Body.Close() reader := csv.NewReader(resp.Body) detector := NewDetector(languages) var results []DetectionResult numWorkers := 4 input := make(chan string, numWorkers) output := make(chan DetectionResult, numWorkers) var wg sync.WaitGroup start := time.Now() // 启动worker,每个worker绑定WaitGroup for w := 1; w <= numWorkers; w++ { wg.Add(1) go worker(w, detector, input, output, &wg) } // 单独goroutine读取CSV并发送任务 go func() { defer close(input) // 读取完所有数据后关闭input通道 for { record, err := reader.Read() if err == io.EOF { break } if err != nil { log.Fatal(err) } text := record[0] input <- text } }() // 等待所有worker完成后关闭output通道 go func() { wg.Wait() close(output) }() // 收集结果 for result := range output { results = append(results, result) } elapsed := time.Since(start) log.Printf("Decoded %d lines of text in %s", len(results), elapsed) } func worker(id int, detector lingua.LanguageDetector, input chan string, output chan DetectionResult, wg *sync.WaitGroup) { defer wg.Done() // worker退出时标记完成 log.Printf("worker %d started\n", id) for t := range input { result := DetectText(detector, t) output <- result } log.Printf("worker %d finished\n", id) }
关键修改点:
- 引入
sync.WaitGroup:每个worker启动时调用wg.Add(1),退出时通过defer wg.Done()标记任务完成。 - 新增goroutine等待worker全部完成:
wg.Wait()会阻塞直到所有worker调用Done(),之后再关闭output通道,确保不会有worker往已关闭的通道写数据。 - 读取CSV的goroutine中
defer close(input):保证无论是否出错,input通道都会被关闭,worker能正常退出循环。
保持结果顺序的优化方案
如果需要结果顺序与原CSV一致,可以给每个任务添加序号,收集结果后按序号排序:
带序号的任务与结果实现:
// 定义带序号的任务结构体 type Task struct { Index int Text string } // 修改DetectionResult增加序号字段 type DetectionResult struct { Index int // 其他原有字段 } func main() { // ... 其他代码不变 ... input := make(chan Task, numWorkers) output := make(chan DetectionResult, numWorkers) // 读取CSV的goroutine修改为发送带序号的任务 go func() { defer close(input) index := 0 for { record, err := reader.Read() if err == io.EOF { break } if err != nil { log.Fatal(err) } text := record[0] input <- Task{Index: index, Text: text} index++ } }() // 收集结果后排序 var results []DetectionResult for result := range output { results = append(results, result) } // 按Index排序恢复原顺序 sort.Slice(results, func(i, j int) bool { return results[i].Index < results[j].Index }) // ... 后续代码不变 ... } func worker(id int, detector lingua.LanguageDetector, input chan Task, output chan DetectionResult, wg *sync.WaitGroup) { defer wg.Done() log.Printf("worker %d started\n", id) for task := range input { result := DetectText(detector, task.Text) result.Index = task.Index // 把序号赋值给结果 output <- result } log.Printf("worker %d finished\n", id) }
注意事项:
- 确认
lingua.LanguageDetector是否并发安全:如果官方文档说明它不是线程安全的,应该给每个worker创建独立的detector实例,避免数据竞争。 - 通道缓冲区大小:可以根据CSV行数调整,避免goroutine阻塞在发送任务阶段。
内容的提问来源于stack exchange,提问作者Or Shemtov
相关产品推荐
相关产品推荐

