Skip to Content
ConcurrencyPipeline

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.

pipeline.go
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 16

Har 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:

fanin.go
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: 54

4+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 

Last updated on