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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 02:50:40