如何构建循环数据流管道?Go语言实现死锁排查与扩展问题
问题描述
核心需求
假设有函数F,逻辑如下:
- 从输入获取值存入
a - 输出
a+10 - 从输入获取值存入
b - 输出
a+b - 停止执行
需要实现两个模块A、B,二者输出互为对方输入,形成循环管道。初始值1传入模块A,最终获取模块B的输出(预期结果为33)。
运行流程要求
- A接收1存入
a,输出11 - B接收11存入
a,输出21 - A接收21存入
b,输出1+21=22 - B接收22存入
b,输出11+22=33
遇到的问题
本人尝试了如下Go代码,但出现死锁:
package main import ( "fmt" "sync" ) func my_function(in, out chan int) { go func() { a := <-in out <- a + 10 b := <-in out <- a + b close(out) }() } func main() { input_a := make(chan int, 2) output_a := make(chan int, 2) input_a <- 1 var wg sync.WaitGroup wg.Add(2) go my_function(input_a, output_a) go my_function(output_a, input_a) go func() { wg.Wait() close(input_a) }() for x := range input_a { fmt.Println("Result:", x) } }
扩展需求
- 能否扩展至5个模块组成的循环管道并获取最后一个模块的输出?
- 此需求用于解决Advent Of Code 2019第7题,若有更优实现方案也请告知。
问题解答
死锁原因分析
代码死锁主要有两个核心问题:
my_function内部额外启动goroutine后,外层goroutine直接退出,导致wg的计数从未被减少——你调用wg.Add(2)但没有任何地方执行wg.Done(),最终wg.Wait()永久阻塞,input_a无法被关闭,主goroutine的range循环也一直等待。- 循环管道的关闭时机逻辑错误,模块执行完成后没有正确通知等待组。
修复后的双模块代码
修改后的代码会正确跟踪goroutine完成状态,避免死锁:
package main import ( "fmt" "sync" ) func myFunction(in, out chan int, wg *sync.WaitGroup) { defer wg.Done() // 第一步:获取a并输出a+10 a := <-in out <- a + 10 // 第二步:获取b并输出a+b b := <-in out <- a + b close(out) } func main() { // 创建无缓冲通道,按流程同步执行即可 chanAB := make(chan int) chanBA := make(chan int) var wg sync.WaitGroup wg.Add(2) // 启动两个模块goroutine go myFunction(chanAB, chanBA, &wg) go myFunction(chanBA, chanAB, &wg) // 传入初始值1到模块A的输入通道 chanAB <- 1 // 等待所有模块完成后关闭通道 go func() { wg.Wait() close(chanAB) }() // 读取最终结果(模块B的最后输出会进入chanAB) var result int for val := range chanAB { result = val } fmt.Println("最终结果:", result) // 输出33 }
扩展到5个模块的实现
要扩展到N个模块(比如5个),可以用切片管理所有通道,形成环形管道:
package main import ( "fmt" "sync" ) func module(in, out chan int, wg *sync.WaitGroup) { defer wg.Done() a := <-in out <- a + 10 b := <-in out <- a + b close(out) } func main() { moduleCount := 5 chans := make([]chan int, moduleCount) for i := range chans { chans[i] = make(chan int) } var wg sync.WaitGroup wg.Add(moduleCount) // 启动每个模块,形成环形:模块i的输入是chans[i],输出是chans[(i+1)%moduleCount] for i := 0; i < moduleCount; i++ { inChan := chans[i] outChan := chans[(i+1)%moduleCount] go module(inChan, outChan, &wg) } // 传入初始值到第一个模块的输入通道 chans[0] <- 1 // 等待所有模块完成后关闭第一个通道(接收最终输出) go func() { wg.Wait() close(chans[0]) }() // 读取最终结果 var result int for val := range chans[0] { result = val } fmt.Println("5模块最终结果:", result) }
Advent Of Code 2019第7题的更优方案
针对AoC2019第7题(放大器电路问题),除了通道模拟管道,还有更贴合题意的实现思路:
- 用结构体封装放大器状态:每个放大器保存当前程序指针、输入输出队列,避免通道同步的额外开销。
- 同步执行而非异步goroutine:放大器的执行是严格顺序依赖的(必须等前一个输出后才能继续),同步执行更简单,也不会出现死锁。
- 闭包封装状态:用闭包为每个放大器绑定独立的输入输出逻辑,简化代码结构。
核心逻辑示例:
func createAmplifier(program []int) func(int) int { pc := 0 inputChan := make(chan int, 2) outputChan := make(chan int) go func() { // 此处实现Intcode虚拟机逻辑,读取inputChan输入,输出到outputChan // 省略具体Intcode实现,按AoC2019第7题要求编写 }() return func(input int) int { inputChan <- input return <-outputChan } }
之后按顺序调用每个放大器的函数,形成循环即可。这种方式更直观,也更容易调试。
内容的提问来源于stack exchange,提问作者Be Chiller Too
相关产品推荐
相关产品推荐

