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.

../_images/fan-in-gopher.png

https://go.dev/talks/2012/concurrency.slide#28

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.