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

如何构建循环数据流管道?Go语言实现死锁排查与扩展问题

问题描述

核心需求

假设有函数F,逻辑如下:

  1. 从输入获取值存入a
  2. 输出a+10
  3. 从输入获取值存入b
  4. 输出a+b
  5. 停止执行

需要实现两个模块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)
    }
}

扩展需求

  1. 能否扩展至5个模块组成的循环管道并获取最后一个模块的输出?
  2. 此需求用于解决Advent Of Code 2019第7题,若有更优实现方案也请告知。

问题解答

死锁原因分析

代码死锁主要有两个核心问题:

  1. my_function内部额外启动goroutine后,外层goroutine直接退出,导致wg的计数从未被减少——你调用wg.Add(2)但没有任何地方执行wg.Done(),最终wg.Wait()永久阻塞,input_a无法被关闭,主goroutine的range循环也一直等待。
  2. 循环管道的关闭时机逻辑错误,模块执行完成后没有正确通知等待组。

修复后的双模块代码

修改后的代码会正确跟踪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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 22:25:57