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.