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洩漏
下一篇: 上下文、同步和並發模式 — 上下文套件、互斥體和進階模式。