Fan out Pattern (BİTMEDİ)

Bir işi birden fazla goroutine’e paylaştıracaksak bu kalıbı, pattern’i, kullanabiliriz. Fan in nasıl huni gibi çoklu girdiyi teke indiriyorsa bu da dağıtıyor gibi düşünülebilir.

Aşağıdaki koda ve çıktısına bakalım:

package main

import (
	"fmt"
	"math/rand"
	"sync"
	"time"
)

func main() {
	c1 := make(chan int)
	c2 := make(chan int)

	go doldur(c1)
	go fanOutIn(c1, c2)

	for v := range c2 {
		fmt.Println(v)
	}
	fmt.Println("Çıkış yapıyorum...")
}

func doldur(c chan int) {
	for i := 0; i < 10; i++ {
		c <- i
	}
	close(c)
}

func fanOutIn(c1, c2 chan int) {
	var wg sync.WaitGroup
	for v := range c1 {
		wg.Add(1)
		go func(v2 int) {
			c2 <- timeConsumingWork(v2)
			wg.Done()
		}(v)
	}
	wg.Wait()
	close(c2)
}

func timeConsumingWork(n int) int {
	time.Sleep(time.Microsecond * time.Duration(rand.Intn(500)))
	return n
}

Bir çıktı:

8
6
3
0
4
1
7
5
9
2
Çıkış yapıyorum...

Sizdeki çıktı farklı olabilir, sırası. Ama tüm rakamları “rastgele” sıra ile gördüğümüz bir çıktı olacaktır. Peki bu kod ne yapıyor?

timeConsumingWork() gerçekten iş yapan bir fonksiyon olsun. Aldığı parametreye göre bir şeyler hesaplıyor olsun. Biz burada rastgele bekleyip parametreyi dönüyoruz. doldur() fonksiyonun yaptığı şey ise sadece işlenecek verileri oluşturmak ve bunları bir kanala basmak. Gelelim esas fonksiyonumuza fanOutIn(). Bu arkadaşımız c1 kanalından aldığı verileri her bir veriyi işlemek için ayrı goroutine’ler oluşturarak aslında timeConsumingWork() işini yaptırtıyor, yani fan out işi yapıyor. Tek bir yerden veriler giriyor, her bir verinin işlenmesi için ayrı goroutine oluşturup onlara dağıtıyor ve gelen verileri, yani işlem çıktılarını, c2 kanalına koyuyor. Bunu bir dispatcher olarak düşünebiliriz. Her bir goroutine için wg.Add(1) çalıştırılıyor ve wg.Wait() ile hepsinin bitmesi bekleniyor.

Throttling, Limitleme, Worker Pool Yaklaşımı

Bu kalıpta veri işleyecek goroutine sayısını limitlemek isteyebiliriz. Mesela işlemci çekirdek sayımız kadar goroutine oluşturmak isteyebiliriz, ya da başka bir sayı. Bunu da şöyle yapabiliriz:

package main

import (
	"fmt"
	"math/rand"
	"sync"
	"time"
)

func main() {
	c1 := make(chan int)
	c2 := make(chan int)

	go doldur(c1)
	go fanOutIn(c1, c2)

	for v := range c2 {
		fmt.Println(v)
	}
	fmt.Println("Çıkış yapıyorum...")
}

func doldur(c chan int) {
	for i := 0; i < 10; i++ {
		c <- i
	}
	close(c)
}

func fanOutIn(c1, c2 chan int) {
	var wg sync.WaitGroup
	const numofgoroutine = 2
	wg.Add(numofgoroutine)
	for i := 0; i < numofgoroutine; i++ {
		go func() {
			for v := range c1 {
				func(v2, i int) {
					c2 <- timeConsumingWork(v2, i)
				}(v, i)
			}
			wg.Done()
		}()
	}
	wg.Wait()
	close(c2)
}

func timeConsumingWork(n, i int) int {
	time.Sleep(time.Millisecond * time.Duration(rand.Intn(500)))
	fmt.Println(i, n)
	return n
}

Bir çıktı:

0 1
1
1 0
0
0 2
2
1 3
3
1 5
5
0 4
4
1 6
6
0 7
7
0 9
9
1 8
8
Çıkış yapıyorum...

timeConsumingWork() fonksiyonuna yeni bir parametre geçtim, hangi goroutine’nin hangi veriyi işlediğini ekrana basıyor. Esas değişiklik fanOutIn() fonksiyonu içerisinde. numofgoroutine adet goroutine’i baştan başlatıyoruz. Tüm goroutine’ler c1 kanalından veri çekip işletip c2ye sonucu koyuyor. Eli boşa çıkan, yeni veriyi alıyor. wg.Add(numofgoroutine) ile de tüm goroutine’lerin, sayısı numofgoroutine kadar, bitmesini bekliyoruz. Burada ilk örnekten farklı olarak goroutine sayısını sabitlemiş olduk. İlk örnekte ise eğer c1 kalanında çok adet veri birikirse ve işleme de uzun sürüyorsa gereksiz fazla sayıda goroutine oluşabilir. Gereksiz diyorum çünkü çekirdek sayısından fazla goroutine alıştırmanın çok bir anlamı yok eğer CPU bound bir iş yapıyorsak, I/O bound ise anlamı olabilir (timer ile beklemek çok net I/O bound). Avantajı ise CPU’yu tam utilize etmeyi sağlayacaktır. İkinci çözümde seçtiğimiz goroutine sayısı 2 bu örnekte ama belki 40 çekirdekli makinadyız, oldukça az utilize etmiş olabiliriz.


v yi parametre olarak geçme AMA 1.22 ile gelen https://go.dev/blog/loopvar-preview seni kurtarsın, loopvar i’de gör. 1.21 ile podman’de v de de gör. Bunları ekle, biraz advanced konular.