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

第 5 課:Goroutine 和 Channels

Goroutines、WaitGroup、緩衝/無緩衝通道。選擇語句、頻道方向、扇入/扇出模式。並發與並行。

💻 程式設計 — 第 5 課 第 5 課:Goroutine 和 Channels

Golang:從基礎到高級

第 2 部分:並發與網絡

亞洲開發網

1. 並發與並行

在開始之前,有必要區分兩個經常被混淆的概念:

  • 並發性:同時處理多個任務 交替地 (1個CPU核心交替運作)。就像一個廚師準備很多菜一樣。
  • 平行性: 真的 同時運行 在多個 CPU 核心上。和很多人一樣,每個人都會做一道菜。

Go 的設計目的是 並發性 - 協助編寫有效處理多個任務的程式。 Go 運行時自動將 goroutine 調度到作業系統執行緒上,在有多個 CPU 核心時利用並行性。

┌─────────── 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. Goroutine

Goroutines 是由 Go 運行時管理的輕量級執行緒。僅用 ~2KB 堆疊初始化(與作業系統線程 ~1-8MB 相比),可以創建數百萬個 goroutine。

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. Goroutine 洩漏

// ⚠️ 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. 頻道

通道是 goroutine 之間的通訊機制。走諺語: “不要透過共享記憶來交流,而是透過交流來共享記憶。”

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
}

六、總結

  • Goroutine:輕量級線程(~2KB),使用以下命令創建 走吧。去 關鍵字。關鍵字
  • 等待組:等待群組 goroutine 完成
  • 頻道:goroutine之間的通訊(無緩衝=同步,緩衝=非同步)
  • 選擇:多通道復用
  • 圖案:扇出/扇入、管道、工作池、完成通道
  • 規則:發送者關閉通道,使用context/done進行取消,避免goroutine洩漏

下一篇: 上下文、同步和並發模式 — 上下文套件、互斥體和進階模式。