Chuyển đến nội dung chính

レッスン 5: ゴルーチンとチャネル

Goroutines、WaitGroup、バッファリングされたチャネル/バッファリングされていないチャネル。ステートメント、チャネル方向、ファンイン/ファンアウト パターンを選択します。同時実行と並列処理。

💻 プログラミング — レッスン 5 レッスン 5: ゴルーチンとチャネル

Golang: 基本から高度まで

パート 2: 同時実行性とネットワーキング

xdev.asia

1. 同時実行性と並列性

始める前に、よく混同される 2 つの概念を区別する必要があります。

  • 同時実行性: 複数のタスクを同時に処理します。 交互に (1つのCPUコアが交互に動作します)。まるでたくさんの料理を作る料理人のように。
  • 平行度:本当に 同時に実行する 複数の CPU コア上で。多くの人と同じように、各人が 1 つの料理を作ります。

Go は以下のために設計されています 同時実行性 — 複数のタスクを効率的に処理するプログラムを作成するのに役立ちます。 Go ランタイムは、複数の CPU コアがある場合の並列処理を利用して、OS スレッド上に goroutine を自動的にスケジュールします。

┌─────────── Go Runtime Scheduler ───────────┐
│                                             │
│  Goroutine 1  Goroutine 2  Goroutine 3     │
│      │             │             │          │
│      ▼             ▼             ▼          │
│  ┌────────┐   ┌────────┐   ┌────────┐     │
│  │ OS     │   │ OS     │   │ OS     │     │
│  │Thread 1│   │Thread 2│   │Thread 3│     │
│  └────────┘   └────────┘   └────────┘     │
│      │             │             │          │
│      ▼             ▼             ▼          │
│   CPU Core 1   CPU Core 2   CPU Core 3     │
└─────────────────────────────────────────────┘

2. ゴルーチン

ゴルーチンは、Go ランタイムによって管理される軽量のスレッドです。わずか約 2KB のスタック (OS スレッド約 1 ~ 8MB と比較) で初期化され、数百万のゴルーチンを作成できます。

2.1.ゴルーチンの作成

package main

import (
    "fmt"
    "time"
)

func sayHello(name string) {
    for i := 0; i < 3; i++ {
        fmt.Printf("Hello from %s (%d)\n", name, i)
        time.Sleep(100 * time.Millisecond)
    }
}

func main() {
    // go keyword tạo goroutine
    go sayHello("Goroutine 1")
    go sayHello("Goroutine 2")
    
    // Anonymous goroutine
    go func() {
        fmt.Println("Hello from anonymous goroutine")
    }()
    
    // ⚠️ main() là goroutine chính
    // Khi main() kết thúc, TẤT CẢ goroutines dừng lại
    time.Sleep(500 * time.Millisecond) // Chờ goroutines hoàn thành
    fmt.Println("Main done")
}

2.2.同期待機グループ

import (
    "fmt"
    "sync"
)

func worker(id int, wg *sync.WaitGroup) {
    defer wg.Done() // Giảm counter khi xong
    
    fmt.Printf("Worker %d starting\n", id)
    time.Sleep(time.Duration(id) * 100 * time.Millisecond)
    fmt.Printf("Worker %d done\n", id)
}

func main() {
    var wg sync.WaitGroup
    
    for i := 1; i <= 5; i++ {
        wg.Add(1) // Tăng counter
        go worker(i, &wg)
    }
    
    wg.Wait() // Block cho đến khi counter = 0
    fmt.Println("All workers completed")
}

2.3.ゴルーチンのリーク

// ⚠️ TRÁNH goroutine leaks - goroutine không bao giờ kết thúc

// ❌ Bad: goroutine leak
func badExample() {
    ch := make(chan int)
    go func() {
        val := <-ch // Block forever nếu không ai gửi
        fmt.Println(val)
    }()
    // Function return, goroutine vẫn chạy (leaked!)
}

// ✅ Good: dùng context để cancel
func goodExample(ctx context.Context) {
    ch := make(chan int)
    go func() {
        select {
        case val := <-ch:
            fmt.Println(val)
        case <-ctx.Done():
            fmt.Println("Cancelled")
            return
        }
    }()
}

// Kiểm tra goroutine leaks trong test:
// import "runtime"
// before := runtime.NumGoroutine()
// ... run test ...
// after := runtime.NumGoroutine()
// assert after == before

3. チャンネル

チャネルはゴルーチン間の通信メカニズムです。囲碁のことわざ: 「記憶を共有することでコミュニケーションするのではなく、コミュニケーションすることで記憶を共有する。」

3.1.バッファリングされていないチャネル

func main() {
    // Unbuffered channel - synchronous
    // Sender blocks cho đến khi receiver sẵn sàng (và ngược lại)
    ch := make(chan string)
    
    go func() {
        ch <- "Hello"  // Send - blocks cho đến khi có receiver
    }()
    
    msg := <-ch         // Receive - blocks cho đến khi có sender
    fmt.Println(msg)     // "Hello"
}

3.2.バッファリングされたチャネル

func main() {
    // Buffered channel - asynchronous (cho đến khi buffer đầy)
    ch := make(chan int, 3) // Buffer size = 3
    
    ch <- 1  // Không block (buffer chưa đầy)
    ch <- 2  // Không block
    ch <- 3  // Không block
    // ch <- 4 // ⚠️ Block! Buffer đầy, chờ receiver
    
    fmt.Println(<-ch)  // 1 (FIFO)
    fmt.Println(<-ch)  // 2
    fmt.Println(<-ch)  // 3
    
    fmt.Println(len(ch))  // 0 (số items trong buffer)
    fmt.Println(cap(ch))  // 3 (buffer capacity)
}

3.3.チャンネルの方向

// Send-only channel: chan<- T
// Receive-only channel: <-chan T

func producer(out chan<- int) {
    for i := 0; i < 5; i++ {
        out <- i
    }
    close(out) // Đóng channel khi xong gửi
}

func consumer(in <-chan int) {
    for val := range in { // range tự dừng khi channel closed
        fmt.Println("Received:", val)
    }
}

func main() {
    ch := make(chan int, 5)
    
    go producer(ch)  // chan tự convert thành chan<-
    consumer(ch)     // chan tự convert thành <-chan
}

3.4.近距離&遠距離

func fibonacci(n int, ch chan<- int) {
    a, b := 0, 1
    for i := 0; i < n; i++ {
        ch <- a
        a, b = b, a+b
    }
    close(ch) // Signal: không còn data
}

func main() {
    ch := make(chan int, 10)
    go fibonacci(10, ch)
    
    // Range over channel - dừng khi channel closed
    for val := range ch {
        fmt.Println(val)
    }
    
    // Check if channel is closed
    val, ok := <-ch
    if !ok {
        fmt.Println("Channel closed, val =", val) // zero value
    }
    
    // ⚠️ Rules:
    // - Chỉ SENDER close channel, không bao giờ receiver
    // - Send vào closed channel → panic
    // - Receive từ closed channel → zero value, false
    // - Close channel đã closed → panic
}

4. ステートメントの選択

選択する スイッチなどの複数のチャネル操作を待機できるようにしますが、チャネルを対象とします。

func main() {
    ch1 := make(chan string)
    ch2 := make(chan string)
    
    go func() {
        time.Sleep(100 * time.Millisecond)
        ch1 <- "from ch1"
    }()
    
    go func() {
        time.Sleep(200 * time.Millisecond)
        ch2 <- "from ch2"
    }()
    
    // select chờ channel nào sẵn sàng trước
    for i := 0; i < 2; i++ {
        select {
        case msg := <-ch1:
            fmt.Println(msg)
        case msg := <-ch2:
            fmt.Println(msg)
        }
    }
}

// Timeout pattern
func fetchWithTimeout(url string) (string, error) {
    ch := make(chan string, 1)
    
    go func() {
        // ... fetch data
        ch <- "data"
    }()
    
    select {
    case data := <-ch:
        return data, nil
    case <-time.After(5 * time.Second):
        return "", fmt.Errorf("timeout after 5s")
    }
}

// Non-blocking send/receive
func nonBlocking() {
    ch := make(chan int, 1)
    
    select {
    case val := <-ch:
        fmt.Println("Received:", val)
    default:
        fmt.Println("No data available") // Không block
    }
    
    select {
    case ch <- 42:
        fmt.Println("Sent")
    default:
        fmt.Println("Channel full") // Không block
    }
}

5. 同時実行パターン

5.1.ファンアウト / ファンイン

// Fan-Out: 1 producer → nhiều workers
// Fan-In: nhiều producers → 1 consumer

// Fan-Out
func fanOut(jobs <-chan int, numWorkers int) []<-chan int {
    workers := make([]<-chan int, numWorkers)
    for i := 0; i < numWorkers; i++ {
        workers[i] = worker(i, jobs)
    }
    return workers
}

func worker(id int, jobs <-chan int) <-chan int {
    results := make(chan int)
    go func() {
        defer close(results)
        for job := range jobs {
            // Process job
            results <- job * job
            fmt.Printf("Worker %d processed job %d\n", id, job)
        }
    }()
    return results
}

// Fan-In: merge nhiều channels thành 1
func fanIn(channels ...<-chan int) <-chan int {
    merged := make(chan int)
    var wg sync.WaitGroup
    
    for _, ch := range channels {
        wg.Add(1)
        go func(c <-chan int) {
            defer wg.Done()
            for val := range c {
                merged <- val
            }
        }(ch)
    }
    
    go func() {
        wg.Wait()
        close(merged)
    }()
    
    return merged
}

func main() {
    // Tạo jobs
    jobs := make(chan int, 10)
    go func() {
        for i := 1; i <= 10; i++ {
            jobs <- i
        }
        close(jobs)
    }()
    
    // Fan-Out: 3 workers
    workers := fanOut(jobs, 3)
    
    // Fan-In: merge results
    results := fanIn(workers...)
    
    for result := range results {
        fmt.Println("Result:", result)
    }
}

5.2.パイプラインパターン

// Pipeline: chuỗi stages, mỗi stage nhận input channel và trả output channel

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

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
}

func filter(in <-chan int, predicate func(int) bool) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            if predicate(n) {
                out <- n
            }
        }
    }()
    return out
}

func main() {
    // Pipeline: generate → square → filter (> 10)
    numbers := generate(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
    squared := square(numbers)
    filtered := filter(squared, func(n int) bool {
        return n > 10
    })
    
    for val := range filtered {
        fmt.Println(val) // 16, 25, 36, 49, 64, 81, 100
    }
}

5.3.労働者プール

type Job struct {
    ID   int
    Data string
}

type Result struct {
    JobID  int
    Output string
}

func workerPool(numWorkers int, jobs <-chan Job) <-chan Result {
    results := make(chan Result, len(jobs))
    var wg sync.WaitGroup
    
    for i := 0; i < numWorkers; i++ {
        wg.Add(1)
        go func(workerID int) {
            defer wg.Done()
            for job := range jobs {
                // Process job
                output := fmt.Sprintf("Worker %d processed: %s", workerID, job.Data)
                results <- Result{JobID: job.ID, Output: output}
            }
        }(i)
    }
    
    go func() {
        wg.Wait()
        close(results)
    }()
    
    return results
}

func main() {
    jobs := make(chan Job, 20)
    
    // Enqueue jobs
    for i := 1; i <= 20; i++ {
        jobs <- Job{ID: i, Data: fmt.Sprintf("task-%d", i)}
    }
    close(jobs)
    
    // Start worker pool (5 workers)
    results := workerPool(5, jobs)
    
    for result := range results {
        fmt.Printf("Job %d: %s\n", result.JobID, result.Output)
    }
}

5.4.完了チャンネル(キャンセル)

func generator(done <-chan struct{}) <-chan int {
    ch := make(chan int)
    go func() {
        defer close(ch)
        i := 0
        for {
            select {
            case <-done:
                return // Cleanup khi cancelled
            case ch <- i:
                i++
            }
        }
    }()
    return ch
}

func main() {
    done := make(chan struct{})
    nums := generator(done)
    
    // Lấy 5 số rồi cancel
    for i := 0; i < 5; i++ {
        fmt.Println(<-nums)
    }
    
    close(done) // Signal all goroutines to stop
}

6. まとめ

  • ゴルーチン: 軽量スレッド (~2KB)、で作成されました。 行きなさい。行きます キーワード。キーワード
  • 待機グループ: グループ goroutine が完了するまで待ちます
  • チャンネル: ゴルーチン間の通信 (バッファなし = 同期、バッファあり = 非同期)
  • 選択: 複数チャンネルでの多重化
  • パターン: ファンアウト/ファンイン、パイプライン、ワーカー プール、完了チャネル
  • ルール: 送信者はチャネルを閉じ、キャンセルには context/done を使用し、ゴルーチン リークを回避します

次の記事: コンテキスト、同期、同時実行パターン — コンテキスト パッケージ、ミューテックス、および高度なパターン。