Pipeline
Pipeline — ma’lumotni bir nechta bosqichdan (stage) ketma-ket o’tkazish naqshi; bosqichlar bir-biriga channel orqali ulanadi. Har bir bosqich — alohida goroutine: kirish channel’idan o’qiydi, ishlaydi, natijani chiqish channel’iga yozadi. Xuddi konveyer kabi — har bosqich o’z ishini bajarib, keyingisiga uzatadi.
Quyida uch bosqich bor: gen sonlarni chiqaradi, square ularni kvadratga oshiradi, main esa natijani chop etadi. Har bosqich yozib bo’lgach o’z channel’ini close qiladi — shunda keyingi bosqichning range i o’zi tabiiy to’xtaydi.
package main
import "fmt"
func gen(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 n := range in {
out <- n * n
}
close(out)
}()
return out
}
func main() {
// gen -> square -> chop etish
for n := range square(gen(2, 3, 4)) {
fmt.Println(n)
}
}$ go run pipeline.go
4
9
16Har bosqich mustaqil goroutine bo’lgani uchun ular bir vaqtda ishlaydi: square birinchi qiymatni kvadratlayotganda gen allaqachon keyingisini tayyorlaydi.
Fan-out / fan-in
Bir bosqich sekin bo’lsa, uni fan-out qilamiz: bitta kirish channel’ini bir nechta bir xil goroutine parallel o’qiydi (yuk taqsimlanadi). So’ng ularning natijalarini bitta channel ga fan-in bilan yig’amiz:
package main
import (
"fmt"
"sync"
)
func gen(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 n := range in {
out <- n * n
}
close(out)
}()
return out
}
// merge - bir nechta channel ni bittaga yig'adi (fan-in)
func merge(cs ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, c := range cs {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for n := range c {
out <- n
}
}(c)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
in := gen(2, 3, 4, 5)
// fan-out: bitta kirishni ikki square parallel o'qiydi
c1 := square(in)
c2 := square(in)
// fan-in: natijalarni yig'amiz
sum := 0
for n := range merge(c1, c2) {
sum += n
}
fmt.Println("yig'indi:", sum)
}$ go run fanin.go
yig'indi: 544+9+16+25 = 54. Ikki square bir in channel’ini baham ko’rdi (fan-out - vazifa ular orasida taqsimlandi), merge esa ikkovining natijasini bitta oqimga yig’di (fan-in). Chiqish tartibi aralash bo’lishi mumkin, lekin yig’indi doim bir xil.
Qisqasi: pipeline bosqichlarni channel bilan ulaydi; fan-out yuk taqsimlaydi, fan-in natijalarni bitta channel’ga yig’adi.
Manba / batafsil: go.dev/blog/pipelines