Fan in Pattern
Şimdi biraz concurrent design pattern yani tasarım örüntülerine, kalıplarına bakalım.
Fan in pattern den başlayalım. Bu kalıpta, birden fazla kanalı tek bir kanala birleştiriyoruz. Onun için adı fan in, elektronikçiler bu terimi bilecektir. Bu örüntüye funnel yani huni tasarımı da denebiliyor. Burada multiple producer, single consumer bir yapıdan söz edebiliriz.
Aşağıdaki diagram durumu görselleştirmektedir.
graph LR
subgraph Producers [Üreticiler / Producers]
P1[Goroutine 1]
P2[Goroutine 2]
P3[Goroutine N]
end
P1 -- "Veri" --> CH((Ortak Kanal<br>Huni / Funnel))
P2 -- "Veri" --> CH
P3 -- "Veri" --> CH
CH -- "Toplanmış Veri" --> C[Tüketici / Consumer]
package main
import (
"fmt"
"sync"
)
func main() {
cift := make(chan int)
tek := make(chan int)
fanin := make(chan int)
go tx(cift, tek)
go rx(cift, tek, fanin)
for v := range fanin {
fmt.Println(v)
}
fmt.Println("Çıkış yapıyorum...")
}
func tx(cift, tek chan<- int) {
for i := 0; i < 10; i++ {
if i%2 == 0 {
cift <- i
} else {
tek <- i
}
}
close(cift)
close(tek)
}
func rx(cift, tek <-chan int, fanin chan<- int) {
var wg sync.WaitGroup
wg.Add(2)
go func() {
for v := range cift {
fanin <- v
}
wg.Done()
}()
go func() {
for v := range tek {
fanin <- v
}
wg.Done()
}()
wg.Wait()
close(fanin)
}
Çıktı:
0
2
1
3
4
5
7
6
8
9
Çıkış yapıyorum...
Yukarıdaki kodu anlamaya çalışalım. Önceki bölümlerde olduğu gibi tek ve çift
sayıları koymak için ayrı ayrı kanallar yaratıyoruz. Bu sefer bir de fanin
isminde bir kanalımız daha var. tx() fonksiyonu ayrı bir goroutine olarak bu
kanallara sayıları yerleştiriyor ve daha sonra kanalları kapatıyor. Yine rx()
fonksiyonunu da ayrı bir goroutine olarak çalıştırıyoruz. Bu fonksiyon kendi
içerisinde iki goroutine daha yaratıyor. Her bir goroutine tek veya çift
kanalını dinleyip huni yani funnel görevi görerek fanin kanalına
koyuyor. for v := range ... kullanımında kanal kapandığı zaman loop’un
sonlandığını önceden konuşmuştuk. Kanallar kapandıktan sonra sync.Waitgroup
ile bunun “sinyalini” veriyoruz. rx() fonksiyonu da yarattığı her iki
goroutine sonlandıktan sonra, wg.Wait() noktasını geçince, o da fanin
kanalını kapatıyor.
Ben burada rx() goroutine olmayabilir diye bir düşünmüştüm. Eğer go rx(...)
demeyip rx(...) dersek deadlock durumu oluşuyor. Onun sebebi de önceki
yazılardakine benzer aslında. Bu durumda fanin kanalını dinleyen kimse olmuyor
ve fanine gönderen iki goroutine bloke oluyor, günün sonunda rx() bloke
oluyor ve range fanin bloke oluyor. Tabii bu problem için fanin := make(chan int, 20) gibi yeteri kadar buffer kapasitesi olan bir kanal yaratırsak bu
problemi görmüyoruz. Buffer’ı olan kanalların önerilmemesinin bir sebebi de bu
muhtemelen. Aslında hatalı olan ya da daha başka/büyük veri boyutları için
çalışmayacak ama bu değerler için çalışan bir tasarımdaki bir hatayı saklamış
oluyor. Sıfır buffer’lı kanalla çalışmak birazcık er meydanı gibi.
Şimdi bir de Rob Pike’ın 2012 yılında yaptığı Go Concurrency Patterns isimli sunumdan esinlenerek hazırlanmış bir koda bakalım, bir örnek daha yapalım.
Önce kod ve çıktısı:
package main
import (
"fmt"
"math/rand"
"time"
)
func fanIn(input1, input2 <-chan string) <-chan string {
c := make(chan string)
go func() {
for {
c <- <-input1
}
}()
go func() {
for {
c <- <-input2
}
}()
return c
}
func boring(msg string) <-chan string {
c := make(chan string)
go func() {
for i := 0; ; i++ {
c <- fmt.Sprintf("%s %d", msg, i)
time.Sleep(time.Duration(rand.Intn(1e3)) * time.Millisecond)
}
}()
return c
}
func main() {
c := fanIn(boring("Hebele"), boring("Hübele"))
for i := 0; i < 10; i++ {
fmt.Println(<-c)
}
fmt.Println("Sıkıldım, çıkış yapıyorum...")
}
Çıktı:
Hübele 0
Hebele 0
Hübele 1
Hübele 2
Hebele 1
Hübele 3
Hübele 4
Hebele 2
Hübele 5
Hebele 3
Sıkıldım, çıkış yapıyorum...
Görselleştirmek adına konulan time.Sleep() ten dolayı kod “yavaş yavaş”
akmaktadır. Şimdi burada ne yapıyoruz?
Yine benzer bir örüntü var fakat biraz daha idiomatic bir kod, Rob Pike sonuçta,
dili yaratan ekipten biri. fanIn() ve boring() fonksiyonları doğrudan
string taşıyan kanallar dönüyorlar. boring() fonksiyonu iki farklı argüman
için çağrılıyor, "Hebele" ve "Hübele". boring() fonksiyonu kendi
içerisinde bir kanal yaratıyor ve arkasından goroutine ile bu kanallara acele
etmeden veri koyuyor, sonsuza kadar. Goroutine yaratılır yaratılmaz yarattığı
kanalı geri dönüyor. fanIn() fonksiyonuna da bu iki kanalı argüman olarak
geçiyoruz. Burada da bir kanal yaratıyoruz, funnel yani huninin ağzı olan kanal.
Parametre olarak gelen iki kanalın verilerini huni gibi bu kanala akıtmak için
birer goroutine oluşturuyor ve bu iki goroutine de sonsuza kadar, kanalda veri
oldukça verileri çekip huniden çıkarıyor. Bu yaratılan kanalı da dönüş olarak
dönüyor. main() içerisinde de belli bir sayıda okuma yapıp sonlandırıyoruz
programımızı, yoksa sonsuza kadar çalışabilir.