【发布时间】:2022-11-05 07:08:32
【问题描述】:
我试图处理从 AWS S3 读取的 CSV 文件,对于每一行我想激活 worker 函数来做一些工作并返回结果
理想情况下,我希望将结果作为原始 CSV 排序,但这不是必需的,由于某种原因,当我运行此代码时,我会遇到奇怪的数据竞争和这一行:
for result := range output {
results = append(results, result)
}
永远阻塞
我尝试使用也不起作用的 WaitGroup,关闭 output 频道也导致我出现“试图将某些内容放入关闭的频道”的错误
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)
start := time.Now()
for w := 1; w < numWorkers+1; w++ {
go worker(w, detector, input, output)
}
go func() {
for {
record, err := reader.Read()
if err == io.EOF {
close(input)
break
}
if err != nil {
log.Fatal(err)
}
text := record[0]
input <- text
}
}()
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) {
log.Printf("worker %d started\n", id)
for t := range input {
result := DetectText(detector, t)
output <- result
}
log.Printf("worker %d finished\n", id)
}
尝试处理 CSV(最好按顺序),并使用对 worker 的函数调用结果来丰富它
尝试设置WaitGroup,尝试在完成读取(EOF)时关闭输出通道 - 导致错误
【问题讨论】:
-
通过编写 4 个文本检测器(工作人员),我从您那里读到了一个隐含的声明,即检测比从 CSV 读取慢 4 倍。那正确吗?我可以看到,也许你正试图检测一种给定未知输入的语言,所以要进行大量的猜测和检查?
标签: csv go concurrency channels routines