Go Academy · 课时

通道方向与流水线模式

只读通道、只写通道与流水线

第 4 / 4 课12 个步骤

通道方向与流水线模式 是 CoddyKit 上的免费 Go Academy 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Go Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Go Academy 课程共包含 4 节课。

单向通道类型

Go 允许您在类型签名中将通道限制为只发送或只接收。这既能说明用途,也能防止误用:

package main

// send-only: can only send, not receive or close (close is allowed by sender)
func sender(out chan<- int) {
    out <- 42
}

// receive-only: can only receive, not send or close
func receiver(in <-chan int) int {
    return <-in
}

// bidirectional: can do both
func relay(in <-chan int, out chan<- int) {
    out <- <-in
}

转换通道方向

双向通道可以隐式转换为任一单向类型,但单向通道不能反向转换为双向通道:

package main
import "fmt"

func produce(out chan<- int) { out <- 1; close(out) }
func consume(in <-chan int)  { fmt.Println(<-in) }

func main() {
    ch := make(chan int)   // bidirectional
    go produce(ch)         // auto-converts to chan<- int
    consume(ch)            // auto-converts to <-chan int
    // chan<- int cannot be converted back to <-chan int
}

流水线:第 1 阶段——生成器

流水线会将多个 goroutine 串联起来,每个阶段都从前一阶段的输出通道读取数据。生成器是第一个阶段:

package main
import "fmt"

func generate(nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        for _, n := range nums {
            out <- n
        }
        close(out)
    }()
    return out
}

func main() {
    ch := generate(2, 3, 4, 5)
    for v := range ch { fmt.Println(v) }
}

流水线:第 2 阶段——转换

中间阶段从一个通道接收数据,进行转换,然后发送到另一个通道:

package main
import "fmt"

func generate(nums ...int) <-chan int {
    out := make(chan int)
    go func() { for _, n := range nums { out <- n }; close(out) }()
    return out
}

func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        for v := range in { out <- v * v }
        close(out)
    }()
    return out
}

func main() {
    for v := range square(generate(2, 3, 4, 5)) {
        fmt.Println(v)  // 4, 9, 16, 25
    }
}

流水线:完整的三阶段示例

串联三个阶段——generate、转换和过滤:

package main
import "fmt"

func gen(nums ...int) <-chan int {
    c := make(chan int)
    go func() { for _, n := range nums { c <- n }; close(c) }()
    return c
}

func double(in <-chan int) <-chan int {
    c := make(chan int)
    go func() { for v := range in { c <- v * 2 }; close(c) }()
    return c
}

func evens(in <-chan int) <-chan int {
    c := make(chan int)
    go func() { for v := range in { if v%2 == 0 { c <- v } }; close(c) }()
    return c
}

func main() {
    for v := range evens(double(gen(1,2,3,4,5))) {
        fmt.Println(v)  // 2, 4, 6, 8, 10
    }
}

使用 done 通道取消流水线

传入一个 done 通道,使流水线阶段能够在取消时提前退出:

package main
import "fmt"

func gen(done <-chan struct{}, nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            select {
            case out <- n:
            case <-done:  // cancel signal
                return
            }
        }
    }()
    return out
}

func main() {
    done := make(chan struct{})
    ch := gen(done, 1, 2, 3, 4, 5)
    fmt.Println(<-ch)  // 1
    close(done)        // cancel remaining work
    fmt.Println("pipeline cancelled")
}

扇出:分发工作

扇出会将一个通道中的工作分发给多个 goroutine,以便并行处理:

package main
import ("fmt"; "sync")

func fanOut(in <-chan int, workers int) []<-chan int {
    outs := make([]<-chan int, workers)
    for i := 0; i < workers; i++ {
        out := make(chan int)
        outs[i] = out
        go func() {
            for v := range in { out <- v * v }  // each worker processes
            close(out)
        }()
    }
    return outs
}

func main() {
    fmt.Println("fan-out distributes work to multiple goroutines")
    _ = sync.WaitGroup{}
}

扇入:合并通道

扇入会将多个通道合并为一个:

package main
import ("fmt"; "sync")

func merge(cs ...<-chan int) <-chan int {
    out := make(chan int)
    var wg sync.WaitGroup
    for _, c := range cs {
        wg.Add(1)
        c := c
        go func() { defer wg.Done(); for v := range c { out <- v } }()
    }
    go func() { wg.Wait(); close(out) }()
    return out
}

func gen(nums ...int) <-chan int {
    c := make(chan int)
    go func() { for _, n := range nums { c <- n }; close(c) }()
    return c
}

func main() {
    for v := range merge(gen(1,2), gen(3,4), gen(5,6)) {
        fmt.Println(v)
    }
}

流水线最佳实践

流水线设计原则:

  • 每个阶段:协程从 <-chan 读取数据,向 chan<- 写入数据,并在完成后关闭输出通道
  • 只有生产者可以关闭通道——消费者绝不能关闭
  • 传入一个 done 通道(或上下文)来执行取消操作
  • 在所有函数签名中使用定向通道类型
  • 取消时排空通道,以防止协程泄漏

真实场景中的流水线:CSV 处理

流水线非常适合数据处理:读取行、解析、筛选和转换:

package main
import ("fmt"; "strings")

func splitLines(data string) <-chan string {
    out := make(chan string)
    go func() {
        for _, line := range strings.Split(data, "\n") {
            if line != "" { out <- line }
        }
        close(out)
    }()
    return out
}

func main() {
    csv := "alice,30\nbob,25\ncarol,35"
    for line := range splitLines(csv) {
        parts := strings.Split(line, ",")
        fmt.Printf("name=%s age=%s\n", parts[0], parts[1])
    }
}

快速检查

函数签名中使用定向通道类型的目的是什么?

回顾:通道方向与流水线

总结:

  • chan<- T 仅发送,<-chan T 仅接收——既记录了设计意图,又提供编译时安全保障
  • 流水线:将协程串联起来,每个协程从一个通道读取数据,再写入另一个通道
  • 扇出:一个输入通道 → 多个协程
  • 扇入:多个通道 → 合并为一个通道
  • 始终由发送方关闭通道;使用 done 通道执行取消操作
免费开始

用 AI 导师学习 Go — 免费

在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。

课程
51
课程
203

常见问题解答

「通道方向与流水线模式」课时是免费的吗?

是的 — 「通道方向与流水线模式」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Go Academy 课程的其余内容,请升级到 CoddyKit PRO。 Go Academy 课程共包含 4 节课。

「通道方向与流水线模式」这节课中我会学到什么?

只读通道、只写通道与流水线 你通过在浏览器中直接运行的动手代码来练习 Go Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Go Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Go Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。

「通道方向与流水线模式」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Go Academy 课中编写并运行代码吗?

能。每节 Go Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 启动 Goroutine
  2. 无缓冲通道
  3. 有缓冲通道
  4. 通道方向与流水线模式
← 返回 Go Academy