go
Go dersleri
Go/Eşzamanlılık

Context ve Eşzamanlılık Desenleri

İptal, zaman aşımı, worker pool, pipeline ve fan-out/fan-in.

Ders 16 / 2130 dkİleri
Bu derste öğreneceklerin
  • context.Context: WithCancel, WithTimeout, WithValue
  • İptal sinyallerini yayma
  • Worker pool deseni
  • Pipeline deseni
  • Fan-out / fan-in
  • Semafor ile eşzamanlılığı sınırlama
  • errgroup ile hata yönetimi

Bir web sunucusu düşün: Kullanıcı bir arama isteği gönderiyor, sunucu üç farklı servise sorgu atıyor, sonuçları birleştirip döndürüyor. Kullanıcı sekmeyi kapatırsa ne olmalı? Yanıt açıktır: Devam eden işlerin hepsi durmalı. Aksi hâlde kimsenin okumayacağı sonuçlar için sunucu kaynak harcamaya devam eder ve yeterince istek iptal edildiğinde sistem çöker.

Go'nun bu soruna cevabı context.Context'tir. Context, bir iş birimi boyunca taşınan ve "bu iş hâlâ gerekli mi?" sorusunun cevabını taşıyan bir değerdir. İptal edildiğinde ya da süresi dolduğunda, onu dinleyen tüm goroutine'ler bunu öğrenir ve düzenli biçimde durur. Standart kütüphanenin ağ, veritabanı ve HTTP API'lerinin neredeyse tamamı ilk parametre olarak bir context alır; bu, dilin en yaygın sözleşmelerinden biridir.

Context tek başına bir eşzamanlılık aracı değildir; onu desenlerle birlikte kullanırsın. Bu derste iptal ve zaman aşımını, iptal sinyalinin nasıl yayıldığını, worker pool, pipeline ve fan-out/fan-in desenlerini, semaforla eşzamanlılık sınırlamayı ve birden çok goroutine'in hatalarını toplamayı öğreneceksin. Bu desenler, Go'da yazacağın hemen her eşzamanlı programın iskeletini oluşturur.

context.Context nedir?

Context, dört şeyi taşıyan küçük bir arayüzdür: bir iptal sinyali (Done kanalı), iptal nedeni (Err), bir son tarih (Deadline) ve isteğe bağlı değerler (Value). Bir context'i değiştiremezsin; ondan türetirsin. Türetilen context, atasının iptalini miras alır.

context.Background()

        ├── WithCancel      → elle iptal edilebilir
        │      │
        │      └── WithTimeout  → süre dolunca kendiliğinden iptal olur
        │             │
        │             └── WithValue → istek kimliği gibi değer taşır

        └── (ata iptal edilirse tüm alt dallar da iptal olur)
main.go
package main

import (
	"context"
	"errors"
	"fmt"
	"time"
)

func work(ctx context.Context, name string, steps int) error {
	for range steps {
		select {
		case <-ctx.Done(): // iptal veya zaman aşımı
			return fmt.Errorf("%s durduruldu: %w", name, ctx.Err())
		default:
		}
		time.Sleep(10 * time.Millisecond) // işi simüle et
	}
	return nil
}

func main() {
	// 1) Elle iptal
	ctx, cancel := context.WithCancel(context.Background())
	go func() {
		time.Sleep(30 * time.Millisecond)
		cancel() // işi durdur
	}()
	err1 := work(ctx, "elle-iptal", 1000)
	fmt.Println(err1)
	fmt.Println("iptal mi:", errors.Is(err1, context.Canceled))

	// 2) Zaman aşımı
	ctx2, cancel2 := context.WithTimeout(context.Background(), 30*time.Millisecond)
	defer cancel2() // sızıntıyı önlemek için her zaman çağrılır
	err2 := work(ctx2, "zaman-aşımı", 1000)
	fmt.Println(err2)
	fmt.Println("süre doldu mu:", errors.Is(err2, context.DeadlineExceeded))

	// 3) Zamanında biten iş
	ctx3, cancel3 := context.WithTimeout(context.Background(), time.Second)
	defer cancel3()
	fmt.Println("kısa iş hatası:", work(ctx3, "kısa", 3))
}
Çıktı
elle-iptal durduruldu: context canceled
iptal mi: true
zaman-aşımı durduruldu: context deadline exceeded
süre doldu mu: true
kısa iş hatası: <nil>

Üç kural context kullanımını güvenli kılar:

cancel'ı her zaman çağır. WithCancel, WithTimeout ve WithDeadline bir cancel fonksiyonu döndürür. Çağırmazsan ata context'te tutulan kaynak serbest bırakılmaz. defer cancel() yazmak standart alışkanlıktır ve iş zaten bitmişse zararsızdır.

Context'i ilk parametre olarak geçir. Gelenek func İş(ctx context.Context, ...) biçimindedir. Context'i struct alanında saklamak, onun yaşam döngüsünü belirsizleştirdiği için önerilmez.

Asla nil context geçirme. Hangi context'i kullanacağını bilmiyorsan context.Background() (kök) veya context.TODO() (henüz karar verilmedi) kullan.

İptal sinyalini yaymak

Context'in gücü, ağaç yapısından gelir. Bir context iptal edildiğinde ondan türetilen tüm context'ler de iptal olur. Böylece derin bir çağrı zincirinin en tepesindeki iptal, en alttaki goroutine'e kadar iner:

main.go
package main

import (
	"context"
	"fmt"
	"sort"
	"sync"
	"time"
)

// producer, iptal edilene kadar sayı üretir.
func producer(ctx context.Context, out chan<- int) {
	defer close(out)
	for i := 1; ; i++ {
		select {
		case out <- i:
			time.Sleep(5 * time.Millisecond)
		case <-ctx.Done():
			return // iptal: kanalı kapat ve çık
		}
	}
}

func main() {
	ctx, cancel := context.WithCancel(context.Background())
	numbers := make(chan int)
	go producer(ctx, numbers)

	var mu sync.Mutex
	var received []int
	var wg sync.WaitGroup

	// Üç tüketici aynı kanaldan okur
	for id := range 3 {
		wg.Go(func() {
			for n := range numbers {
				mu.Lock()
				received = append(received, n*10+id)
				mu.Unlock()
			}
		})
	}

	time.Sleep(60 * time.Millisecond)
	cancel() // tek bir iptal, tüm zinciri durdurur
	wg.Wait()

	mu.Lock()
	count := len(received)
	sort.Ints(received)
	mu.Unlock()

	fmt.Println("iptal sonrası tüm goroutine'ler durdu")
	fmt.Println("en az bir değer alındı mı:", count > 0)
	fmt.Println("context hatası:", ctx.Err())
}
Çıktı
iptal sonrası tüm goroutine'ler durdu
en az bir değer alındı mı: true
context hatası: context canceled

Bu programda üretici kanalı kapattığı için tüketicilerin range döngüleri kendiliğinden biter. Kanalı yalnızca yazan taraf kapatır; bu, Go'da kanal kullanımının değişmez kuralıdır.

select deyimi burada kilit rol oynar: Hem kanala yazmayı hem iptali aynı anda bekler, hangisi önce hazır olursa onu seçer. İptal kontrolünü default ile yapmak (önceki örnekteki gibi) yalnızca bloklanmayan bir yoklama isteniyorsa uygundur.

Worker pool deseni

Sınırsız goroutine başlatmak cazip görünür ama tehlikelidir: Her biri bellek tüketir, hepsi aynı veritabanına yüklenirse onu çökertir. Worker pool, sabit sayıda işçinin bir iş kanalından besleniyor olmasıdır.

main.go
package main

import (
	"fmt"
	"sort"
	"sync"
)

type Job struct {
	ID    int
	Input int
}

type Result struct {
	JobID  int
	Output int
}

func worker(jobs <-chan Job, results chan<- Result) {
	for job := range jobs {
		results <- Result{JobID: job.ID, Output: job.Input * job.Input}
	}
}

func main() {
	const (
		jobCount  = 12
		workerNum = 4
	)

	jobs := make(chan Job, jobCount)
	results := make(chan Result, jobCount)

	var wg sync.WaitGroup
	for range workerNum {
		wg.Go(func() { worker(jobs, results) })
	}

	for i := 1; i <= jobCount; i++ {
		jobs <- Job{ID: i, Input: i}
	}
	close(jobs) // işler bitti: işçiler range döngüsünden çıkar

	wg.Wait()
	close(results)

	var out []Result
	for r := range results {
		out = append(out, r)
	}
	sort.Slice(out, func(i, j int) bool { return out[i].JobID < out[j].JobID })

	for _, r := range out {
		fmt.Printf("%d%d ", r.JobID, r.Output)
	}
	fmt.Println()
	fmt.Println("işçi sayısı:", workerNum, "| sonuç sayısı:", len(out))
}
Çıktı
1→1 2→4 3→9 4→16 5→25 6→36 7→49 8→64 9→81 10→100 11→121 12→144
işçi sayısı: 4 | sonuç sayısı: 12

Deseni doğru kurmanın anahtarı kapatma sırasıdır: Önce iş kanalı kapatılır, işçiler doğal olarak biter, Wait döner, ancak ondan sonra sonuç kanalı kapatılabilir. Sonuç kanalını erken kapatırsan hâlâ yazan bir işçi panik oluşturur.

Kaç işçi olmalı? Cevap işin türüne bağlıdır. İşlemci yoğun işlerde çekirdek sayısı (runtime.NumCPU()) iyi bir başlangıçtır. Ağ veya disk beklemesi ağırlıklı işlerde çok daha fazlası mantıklıdır, çünkü işçilerin çoğu beklemede olur. Tek doğru cevap yoktur; ölçerek karar verilir.

Pipeline deseni

Pipeline, her aşaması bir kanaldan okuyup diğerine yazan goroutine'ler zinciridir. Unix boru hattının Go'daki karşılığıdır ve büyük veri akışlarını bellekte biriktirmeden işlemeni sağlar.

main.go
package main

import "fmt"

// Aşama 1: kaynak
func generate(nums ...int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for _, n := range nums {
			out <- n
		}
	}()
	return out
}

// Aşama 2: dönüştür
func square(in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for n := range in {
			out <- n * n
		}
	}()
	return out
}

// Aşama 3: filtrele
func odds(in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for n := range in {
			if n%2 == 1 {
				out <- n
			}
		}
	}()
	return out
}

func main() {
	pipeline := odds(square(generate(1, 2, 3, 4, 5, 6, 7)))

	total := 0
	for n := range pipeline {
		fmt.Print(n, " ")
		total += n
	}
	fmt.Println()
	fmt.Println("toplam:", total)
}
Çıktı
1 9 25 49
toplam: 84

Her aşamanın kendi kanalını kapatması ve yalnızca kendi çıkışını kapatması, zincirin sonuna kadar temiz biçimde kapanmasını sağlar. Aşamalar eşzamanlı çalıştığı için, ilk sayı son aşamaya ulaştığında kaynak hâlâ üretiyor olabilir — bu, belleği sabit tutar.

Gerçek bir pipeline'da her aşamaya context de geçirilir; böylece tüketici erken vazgeçtiğinde yukarıdaki aşamalar sonsuza kadar yazmaya çalışıp sızmaz.

Fan-out / fan-in

Bir aşama darboğaz oluşturuyorsa, o aşamayı çoğaltırsın: fan-out. Ardından birden çok kanalın çıktısını tek kanalda toplarsın: fan-in.

                ┌── işçi 1 ──┐
kaynak ──────►  ├── işçi 2 ──┤ ──► birleştir ──► tüketici
                └── işçi 3 ──┘
   (fan-out: iş dağıtılır)   (fan-in: sonuçlar toplanır)
main.go
package main

import (
	"fmt"
	"sort"
	"sync"
)

func generate(n int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for i := 1; i <= n; i++ {
			out <- i
		}
	}()
	return out
}

// Aynı girişi okuyan birden çok işçi: fan-out
func heavyWork(in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for n := range in {
			sum := 0
			for i := 1; i <= n; i++ { // biraz iş yapalım
				sum += i
			}
			out <- sum
		}
	}()
	return out
}

// Birden çok kanalı tek kanalda toplar: fan-in
func merge(channels ...<-chan int) <-chan int {
	out := make(chan int)
	var wg sync.WaitGroup

	for _, ch := range channels {
		wg.Go(func() {
			for v := range ch {
				out <- v
			}
		})
	}

	go func() {
		wg.Wait()
		close(out) // tüm kaynaklar bitince kapat
	}()
	return out
}

func main() {
	source := generate(9)

	// Üç işçi aynı kaynağı paylaşır
	merged := merge(heavyWork(source), heavyWork(source), heavyWork(source))

	var results []int
	for v := range merged {
		results = append(results, v)
	}
	sort.Ints(results) // varış sırası belirsizdir; sıralayarak deterministik yap

	fmt.Println(results)
	fmt.Println("sonuç sayısı:", len(results))
}
Çıktı
[1 3 6 10 15 21 28 36 45]
sonuç sayısı: 9

Fan-in fonksiyonundaki kapatma mantığına dikkat et: Çıkış kanalı, tüm kaynak kanallar bittikten sonra kapatılmalıdır. Bunu ayrı bir goroutine'de wg.Wait() bekleyip kapatarak yaparsın; doğrudan merge içinde beklersen fonksiyon geri dönmez ve tüketici hiç başlayamaz.

Sonuçların varış sırası işçilerin hızına bağlıdır ve öngörülemez. Sıralı çıktı gerekiyorsa ya sonuçları toplayıp sıralarsın ya da her sonuca bir sıra numarası eklersin.

Eşzamanlılığı sınırlama ve hata toplama

Bazen havuz kurmak fazla gelir; tek istediğin "aynı anda en fazla N tane" kuralıdır. Tamponlu kanal bunun için yeterlidir. Hataları toplamak için de her goroutine'in hatasını bir kanala ya da korumalı bir dilime yazması gerekir:

main.go
package main

import (
	"errors"
	"fmt"
	"sort"
	"sync"
)

func process(id int) (string, error) {
	if id%4 == 0 {
		return "", fmt.Errorf("iş %d başarısız", id)
	}
	return fmt.Sprintf("iş %d tamam", id), nil
}

func main() {
	const limit = 3
	sem := make(chan struct{}, limit) // semafor

	var (
		mu      sync.Mutex
		okList  []string
		errList []error
		wg      sync.WaitGroup
	)

	for id := 1; id <= 10; id++ {
		wg.Go(func() {
			sem <- struct{}{}        // yer kap
			defer func() { <-sem }() // yeri bırak

			out, err := process(id)
			mu.Lock()
			defer mu.Unlock()
			if err != nil {
				errList = append(errList, err)
				return
			}
			okList = append(okList, out)
		})
	}
	wg.Wait()

	sort.Strings(okList)
	// Hataların birikme sırası zamanlamaya bağlıdır; deterministik çıktı için sırala.
	sort.Slice(errList, func(i, j int) bool { return errList[i].Error() < errList[j].Error() })

	fmt.Println("başarılı:", len(okList))
	for _, s := range okList {
		fmt.Println(" ", s)
	}

	// Birden çok hatayı tek hatada birleştir (Go 1.20+)
	joined := errors.Join(errList...)
	fmt.Println("hatalı iş sayısı:", len(errList))
	fmt.Println("birleşik hata:\n" + joined.Error())
}
Çıktı
başarılı: 8
  iş 1 tamam
  iş 10 tamam
  iş 2 tamam
  iş 3 tamam
  iş 5 tamam
  iş 6 tamam
  iş 7 tamam
  iş 9 tamam
hatalı iş sayısı: 2
birleşik hata:
iş 4 başarısız
iş 8 başarısız

errors.Join birden çok hatayı tek bir hata değerinde toplar ve errors.Is ile her biri ayrı ayrı sorgulanabilir. Hata yönetimi araçlarının tamamını Hata Yönetimi dersinde bulabilirsin.

Bu deseni daha kısa yazmanın yaygın bir yolu, golang.org/x/sync/errgroup paketidir. Standart kütüphanenin parçası olmadığı için burada çalıştırılabilir bir örnek veremiyoruz, ama kullanımı şöyledir:

g, ctx := errgroup.WithContext(ctx)
g.SetLimit(3) // aynı anda en fazla 3 goroutine

for _, url := range urls {
	g.Go(func() error {
		return fetch(ctx, url) // ilk hata, ctx'i iptal eder
	})
}

if err := g.Wait(); err != nil { // ilk hatayı döndürür
	return err
}

errgroup'un en güçlü yanı şudur: Goroutine'lerden biri hata döndürdüğünde, gruba bağlı context otomatik olarak iptal edilir ve diğer işler de durur. Elle yazdığında bu davranışı kurmak birkaç düzine satır alır.

Desenleri seçerken

Eşzamanlı kod yazarken en sık yapılan hata, ihtiyaç duyulandan daha karmaşık bir yapı kurmaktır. Aşağıdaki sıralama, günlük hayatta işe yarayan bir karar rehberidir.

Önce eşzamanlılığa gerçekten ihtiyacın olduğundan emin ol. Sıralı bir döngü yeterince hızlıysa, goroutine eklemek yalnızca hata ayıklamayı zorlaştırır. Eşzamanlılığın kazandırdığı yer, işlerin birbirini beklediği durumlardır: ağ çağrıları, disk okumaları, bağımsız hesaplamalar.

Sabit sayıda bağımsız iş varsa en basit yol, hepsini başlatıp beklemektir. Sonuçları önceden boyutlandırılmış bir dilimin ayrı indekslerine yazarsan kilit bile gerekmez. Üç servise paralel istek atmak tam olarak böyle bir durumdur.

İş sayısı büyük ya da bilinmiyorsa havuz kur. Binlerce iş için binlerce goroutine başlatmak çalışır ama karşı tarafı boğar. Havuz, hem kaynak kullanımını hem de dış sistemlere binen yükü öngörülebilir kılar.

Veri aşamalardan geçiyorsa pipeline kur. Okuma, ayrıştırma, dönüştürme, yazma gibi adımlar birbirini izliyorsa her adımı ayrı bir aşama yapmak hem belleği sabit tutar hem de her aşamayı bağımsız olarak ölçeklendirmene izin verir.

Her durumda iptali baştan düşün. Sonradan eklenen iptal desteği, mevcut tüm imzaları değiştirmeyi gerektirir. Fonksiyonlarına en baştan bir context parametresi koymak, ileride çok daha ucuza gelir.

Son bir uyarı: Goroutine sızıntısı sessiz bir hatadır. Bir goroutine kimsenin okumadığı bir kanala yazmaya çalışırken sonsuza kadar bekleyebilir ve bunu hiçbir test fark etmez. Her goroutine için "bu nasıl ve ne zaman bitecek?" sorusunu yanıtlayabiliyor olmalısın.

Sık yapılan hatalar

  • cancel çağrısını atlamak. Kaynak sızar; defer cancel() yaz.
  • Context'i struct alanında saklamak. Yaşam döngüsü belirsizleşir; parametre olarak geçir.
  • Zorunlu parametreleri WithValue ile taşımak. İmza yalan söyler, derleyici yardım edemez.
  • Kanalı okuyan tarafın kapatması. Kapatmayı her zaman yazan taraf yapar.
  • Sonuç kanalını işçiler bitmeden kapatmak. Kapalı kanala yazmak panik oluşturur.
  • Fan-in'de wg.Wait()'i ayrı goroutine'e almamak. Fonksiyon dönmez, tüketici hiç başlamaz.
  • Eşzamanlı sonuçların sırasına güvenmek. Varış sırası öngörülemez; sırala veya sıra numarası ekle.
  • Sınırsız goroutine başlatmak. Karşı sistemi ve belleği tüketir; havuz veya semafor kullan.

Alıştırmalar

Alıştırma·Zaman aşımlı iş
Kolay

Verilen adım sayısı kadar çalışan, her adımda context'i kontrol eden bir fonksiyon yaz. Aynı fonksiyonu bir kez yeterli, bir kez yetersiz zaman aşımıyla çağır ve iki durumu da raporla. Hatanın zaman aşımından mı iptalden mi kaynaklandığını ayırt et.

İpucu

errors.Is(err, context.DeadlineExceeded) ve errors.Is(err, context.Canceled) ile ayrım yapabilirsin.

Çözümü göster
main.go
package main

import (
	"context"
	"errors"
	"fmt"
	"time"
)

func run(ctx context.Context, steps int) (int, error) {
	done := 0
	for range steps {
		select {
		case <-ctx.Done():
			return done, ctx.Err()
		default:
		}
		time.Sleep(5 * time.Millisecond)
		done++
	}
	return done, nil
}

// report, zamanlamaya bağlı adım sayısını değil yalnızca sonucun türünü yazar:
// böylece çıktı her çalıştırmada aynı olur.
func report(name string, err error) {
	switch {
	case err == nil:
		fmt.Printf("%s: iş tamamlandı\n", name)
	case errors.Is(err, context.DeadlineExceeded):
		fmt.Printf("%s: süre doldu\n", name)
	case errors.Is(err, context.Canceled):
		fmt.Printf("%s: iptal edildi\n", name)
	default:
		fmt.Printf("%s: beklenmeyen hata: %v\n", name, err)
	}
}

func main() {
	ctx1, cancel1 := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel1()
	done1, err1 := run(ctx1, 5)
	report("yeterli süre", err1)
	fmt.Println("  tamamlanan adım:", done1)

	ctx2, cancel2 := context.WithTimeout(context.Background(), 20*time.Millisecond)
	defer cancel2()
	done2, err2 := run(ctx2, 10_000)
	report("yetersiz süre", err2)
	fmt.Println("  tüm adımlar bitti mi:", done2 == 10_000)

	ctx3, cancel3 := context.WithCancel(context.Background())
	go func() {
		time.Sleep(20 * time.Millisecond)
		cancel3()
	}()
	done3, err3 := run(ctx3, 10_000)
	report("elle iptal", err3)
	fmt.Println("  tüm adımlar bitti mi:", done3 == 10_000)
}
Çıktı
yeterli süre: iş tamamlandı
  tamamlanan adım: 5
yetersiz süre: süre doldu
  tüm adımlar bitti mi: false
elle iptal: iptal edildi
  tüm adımlar bitti mi: false

errors.Is ile iki iptal nedenini ayırt edebilmek pratikte önemlidir: Süre dolduğunda genellikle isteği yeniden denemek mantıklıdır, ama kullanıcı vazgeçtiğinde denemenin bir anlamı yoktur. İkisini aynı kefeye koyan kod, gereksiz yeniden denemelerle sistemi meşgul eder.

Çıktıda tamamlanan adım sayısını doğrudan yazmadığımıza dikkat et: O sayı makinenin hızına ve planlayıcının davranışına bağlıdır. Zamanlamaya bağlı programları sınarken sayıları değil, değişmezleri doğrulamak gerekir: "iş tamamlanmadı", "hata türü zaman aşımıdır" gibi.

Alıştırma·Worker pool ile toplama
Orta

1'den 50'ye kadar sayıların karelerini 5 işçiyle hesaplayan bir worker pool yaz. Sonuçları topla, toplamın doğru olduğunu formülle karşılaştırarak doğrula. Sonuçları sıralı biçimde de yazdır.

İpucu

İlk n sayının kareleri toplamı n(n+1)(2n+1)/6 formülüyle bulunur; hesabının doğruluğunu böyle sınayabilirsin.

Çözümü göster
main.go
package main

import (
	"fmt"
	"sort"
	"sync"
)

func main() {
	const (
		n         = 50
		workerNum = 5
	)

	jobs := make(chan int, n)
	results := make(chan int, n)

	var wg sync.WaitGroup
	for range workerNum {
		wg.Go(func() {
			for j := range jobs {
				results <- j * j
			}
		})
	}

	for i := 1; i <= n; i++ {
		jobs <- i
	}
	close(jobs)
	wg.Wait()
	close(results)

	var squares []int
	total := 0
	for r := range results {
		squares = append(squares, r)
		total += r
	}
	sort.Ints(squares)

	expected := n * (n + 1) * (2*n + 1) / 6
	fmt.Println("ilk beş kare:", squares[:5])
	fmt.Println("son üç kare:", squares[len(squares)-3:])
	fmt.Println("toplam:", total, "| formül:", expected, "| eşit mi:", total == expected)
	fmt.Println("sonuç sayısı:", len(squares))
}
Çıktı
ilk beş kare: [1 4 9 16 25]
son üç kare: [2304 2401 2500]
toplam: 42925 | formül: 42925 | eşit mi: true
sonuç sayısı: 50

İş kanalını tamponlu yapmak, ana goroutine'in tüm işleri bloklanmadan kuyruğa koymasını sağlar. Sonuç kanalı da tamponlu olmasaydı, işçiler kimse okumadığı için kilitlenirdi — çünkü ana goroutine wg.Wait() üzerinde bekliyor olurdu. Bu, worker pool kurarken en sık düşülen tuzaklardan biridir.

Alıştırma·İptal edilebilir pipeline
Zor

Üç aşamalı bir pipeline kur: üret → karesini al → süz. Her aşama context'i dinlesin ve iptal edildiğinde temiz biçimde kapansın. Tüketici ilk beş sonucu aldıktan sonra iptal etsin; hiçbir goroutine'in sızmadığını runtime.NumGoroutine() ile göster.

İpucu

Her aşamada select ile hem çıkışa yazmayı hem ctx.Done()'ı bekle. İptalden sonra goroutine'lerin kapanması için kısa bir süre beklemen gerekebilir.

Çözümü göster
main.go
package main

import (
	"context"
	"fmt"
	"runtime"
	"time"
)

func generate(ctx context.Context) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for i := 1; ; i++ {
			select {
			case out <- i:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

func square(ctx context.Context, in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for n := range in {
			select {
			case out <- n * n:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

func onlyEven(ctx context.Context, in <-chan int) <-chan int {
	out := make(chan int)
	go func() {
		defer close(out)
		for n := range in {
			if n%2 != 0 {
				continue
			}
			select {
			case out <- n:
			case <-ctx.Done():
				return
			}
		}
	}()
	return out
}

func main() {
	before := runtime.NumGoroutine()

	ctx, cancel := context.WithCancel(context.Background())
	pipeline := onlyEven(ctx, square(ctx, generate(ctx)))

	var got []int
	for v := range pipeline {
		got = append(got, v)
		if len(got) == 5 {
			break // tüketici erken vazgeçti
		}
	}
	cancel() // yukarıdaki tüm aşamaları durdur

	// Goroutine'lerin kapanmasına zaman tanı
	for range 100 {
		if runtime.NumGoroutine() <= before {
			break
		}
		time.Sleep(time.Millisecond)
	}

	fmt.Println("alınan değerler:", got)
	fmt.Println("sızıntı var mı:", runtime.NumGoroutine() > before)
}
Çıktı
alınan değerler: [4 16 36 64 100]
sızıntı var mı: false

Buradaki kritik ayrıntı, her aşamanın çıkışa yazarken select kullanmasıdır. out <- v satırını doğrudan yazsaydın, tüketici vazgeçtikten sonra aşama sonsuza kadar yazmayı beklerdi: klasik goroutine sızıntısı. ctx.Done()'ı aynı select içinde beklemek, iptal geldiğinde yazma denemesinden vazgeçilmesini sağlar.

for range ile kısa aralıklarla beklemek, goroutine'lerin planlayıcı tarafından sonlandırılmasına fırsat tanır; gerçek kodda bunun yerine sync.WaitGroup ile beklemek daha kesin bir yoldur.

Kısa sınav

Kısa sınav

context.WithTimeout çağrısının döndürdüğü cancel fonksiyonu neden her zaman çağrılmalıdır?

Bir kanalı kim kapatmalıdır?

Worker pool'da sonuç kanalı ne zaman kapatılmalıdır?

Fan-in fonksiyonunda wg.Wait() neden ayrı bir goroutine içinde çağrılır?

context.WithValue hangi bilgi için uygundur?

Goroutine sızıntısının en yaygın nedeni nedir?

Özet

  • context.Context, bir iş biriminin hâlâ gerekli olup olmadığını taşır; iptal, son tarih ve istek kapsamlı değerler sağlar.
  • Context'ler ağaç oluşturur: Ata iptal edilince tüm alt dallar da iptal olur. cancel her zaman çağrılmalıdır.
  • Context ilk parametre olarak geçirilir, struct'ta saklanmaz; WithValue yalnızca yardımcı veriler içindir.
  • Worker pool, sabit sayıda işçiyle sınırsız goroutine başlatmayı önler; kapatma sırası kritiktir.
  • Pipeline, aşamaları kanallarla bağlar ve belleği sabit tutar; her aşama yalnızca kendi çıkışını kapatır.
  • Fan-out işi çoğaltır, fan-in sonuçları tek kanalda toplar; birleştirici kanalı tüm kaynaklar bitince kapatır.
  • Tamponlu kanal doğal bir semafordur; errors.Join birden çok hatayı tek değerde toplar.
  • Her goroutine için "bu nasıl bitecek?" sorusunu yanıtlayabilmelisin; yanıtlayamıyorsan sızıntı vardır.
Bu dersi bitirdin mi?
İlerlemen bu tarayıcıda saklanır.