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
gokeyword - 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.