【发布时间】:2021-12-24 20:33:48
【问题描述】:
我想知道并了解如何使用 go 执行过滤并发管道 在生产者/消费者计划中。
我已经编写了一个版本来检查一个值,如果没问题,就将它发送到一个频道 如果不是,则将该值发送到另一个通道。
在读取和处理完值之后,两个goroutine负责读取处理后的值 并将它们写入文件。这个版本运行正常。但是……
-
假设我不想要无效值。有没有办法改变 select 语句(或消费者 goroutine),这样只有 输出正确的值(即仅使用一个输出通道)。我尝试删除该 invalidValues 频道,但 我没有成功。
-
我尝试将 select 语句放在
if valid?;有一个分支,其中包含此版本中的完整语句和错误分支 只需等待完成的频道。通过这种方式,我可以丢弃无效值并使用一个通道,但这种方法也没有成功。
关于如何解决这个问题的任何想法?
- 此外,在这个方案中,我想知道为什么如果我省略从程序的 invalidValues 通道中删除值的 goroutine 没有完成?是不是通道需要清空,否则仍然被阻塞?有没有更优雅的方法来做一个范围 价值观?
谢谢!!
//Consumers
var wg sync.WaitGroup
wg.Add(Workers)
for i := 0; i < Workers; i++ {
// Deploy #Workers to read from the inputStream perform validation and output the valid results to one channel and the invalid to another
go func() {
for value := range inputStream {
var c *chan string
dataToWrite := value
if valid := checkValue(value); valid {
dataToWrite = value
c = &outputStream
} else {
c = &invalidValues
}
select {
case *c <- dataToWrite:
case <-done:
return
}
time.Sleep(time.Duration(5) * time.Second)
}
wg.Done()
}()
}
这里是完整版的代码
done := make(chan struct{})
defer close(done)
inputStream := make(chan string)
outputStream := make(chan string)
invalidValues := make(chan string)
//Producer reads a file with values and stores them in a channel
go func() {
count := 0
scanner := bufio.NewScanner(file)
for scanner.Scan() {
inputStream <- strings.TrimSpace(scanner.Text())
count = count + 1
}
close(inputStream)
}()
//Consumers
var wg sync.WaitGroup
wg.Add(Workers)
for i := 0; i < Workers; i++ {
// Deploy #Workers to read from the inputStream perform validation and output the valid results to one channel and the invalid to another
go func() {
for value := range inputStream {
var c *chan string
dataToWrite := value
if valid := checkValue(value); valid {
dataToWrite = value
c = &outputStream
} else {
c = &invalidValues
}
select {
case *c <- dataToWrite:
case <-done:
return
}
time.Sleep(time.Duration(5) * time.Second)
}
wg.Done()
}()
}
go func() {
wg.Wait()
close(outputStream)
close(invalidValues)
}()
//Write outputStream file
resultFile, err := os.Create("outputStream.txt")
if err != nil {
log.Fatal(err)
}
//Error file
errorFile, err := os.Create("errors.txt")
if err != nil {
log.Fatal(err)
}
//Create two goruotines for writing the outputStream file
var wg2 sync.WaitGroup
wg2.Add(2)
go func() {
//Write outputStream and error to files
for r := range outputStream {
_, err := resultFile.WriteString(r + "\n")
if err != nil {
log.Fatal(err)
}
}
resultFile.Close()
wg2.Done()
}()
go func() {
for r := range invalidValues {
_, err := errorFile.WriteString(r + "\n")
if err != nil {
log.Fatal(err)
}
}
errorFile.Close()
wg2.Done()
}()
wg2.Wait()
【问题讨论】:
标签: go concurrency channels