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

Bài 5: Goroutines & Channels

Goroutines, WaitGroup, buffered/unbuffered channels. Select statement, channel directions, fan-in/fan-out patterns. Concurrency vs parallelism.

💻 Lập trình — Bài 5 Bài 5: Goroutines & Channels

Golang: Từ Cơ bản đến Nâng cao

Phần 2: Concurrency & Networking

xdev.asia

1. Concurrency vs Parallelism

Trước khi bắt đầu, cần phân biệt 2 khái niệm thường bị nhầm lẫn:

  • Concurrency: Xử lý nhiều tasks cùng lúc bằng cách luân phiên (1 CPU core chạy xen kẽ). Giống 1 người nấu ăn chuẩn bị nhiều món.
  • Parallelism: Thực sự chạy đồng thời trên nhiều CPU cores. Giống nhiều người nấu ăn mỗi người 1 món.

Go được thiết kế cho concurrency — giúp viết chương trình xử lý nhiều tasks hiệu quả. Go runtime tự động schedule goroutines lên OS threads, tận dụng parallelism khi có nhiều CPU cores.

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

Goroutine là lightweight thread được quản lý bởi Go runtime. Khởi tạo chỉ với ~2KB stack (so với OS thread ~1-8MB), có thể tạo hàng triệu goroutines.

2.1. Tạo Goroutine

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. sync.WaitGroup

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 Leaks

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

Channels là cơ chế communication giữa goroutines. Go proverb: "Don't communicate by sharing memory, share memory by communicating."

3.1. Unbuffered Channel

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. Buffered Channel

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. Channel Directions

// 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. Close & Range

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. Select Statement

select cho phép chờ trên nhiều channel operations — giống switch nhưng cho channels.

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. Concurrency Patterns

5.1. Fan-Out / Fan-In

// 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 Pattern

// 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. Worker Pool

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. Done Channel (Cancellation)

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. Tổng kết

  • Goroutines: Lightweight threads (~2KB), tạo với go keyword
  • WaitGroup: Chờ group goroutines hoàn thành
  • Channels: Communication between goroutines (unbuffered = sync, buffered = async)
  • Select: Multiplexing trên nhiều channels
  • Patterns: Fan-Out/Fan-In, Pipeline, Worker Pool, Done Channel
  • Quy tắc: Sender close channel, dùng context/done cho cancellation, tránh goroutine leaks

Bài tiếp theo: Context, Sync & Concurrency Patterns — context package, mutex, và các patterns nâng cao.