Ba mẫu này chiếm phần lớn mã đồng thời trong Go thật. Chúng đơn giản tới mức không cần thư viện — chỉ goroutine và channel.
Pipeline: nối các giai đoạn
Mỗi giai đoạn là một hàm nhận channel vào, trả channel ra:
func sinh(ctx context.Context, n int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := 1; i <= n; i++ {
select {
case out <- i:
case <-ctx.Done():
return
}
}
}()
return out
}
Nối lại:
for v := range binhPhuong(ctx, sinh(ctx, 5)) {
tong += v
}
tổng bình phương 1..5 = 55
Bốn quy ước làm nên một giai đoạn đúng:
Trả về channel chỉ-nhận (<-chan int) — người gọi không đóng được, đúng quy tắc ở bài 32.
defer close(out) — giai đoạn tự đóng đầu ra của mình khi xong.
select với ctx.Done() ở mọi lệnh gửi — nếu không, giai đoạn kẹt khi hạ nguồn ngừng đọc.
Đọc bằng range — tự dừng khi thượng nguồn đóng.
Điểm mạnh: mỗi giai đoạn chạy song song với các giai đoạn khác. Trong khi giai đoạn 2 xử lý phần tử thứ nhất, giai đoạn 1 đã sinh phần tử thứ hai.
Fan-out: nhiều worker, một nguồn
ng := sinh(ctx, 100)
var w []<-chan int
for i := 0; i < 4; i++ {
w = append(w, binhPhuong(ctx, ng)) // BỐN worker, CÙNG một channel vào
}
4 worker xử lý 100 phần tử, tổng=338350
Đây là chỗ channel của Go tỏa sáng: nhiều goroutine đọc cùng một channel, và Go tự chia việc. Không cần phân mảnh dữ liệu, không cần hàng đợi riêng cho từng worker, không cần cân bằng tải.
Mỗi giá trị đi tới đúng một worker. Worker nào rảnh thì nhận — nên tải tự cân bằng theo tốc độ xử lý thật, không theo số lượng chia trước.
Fan-in: gộp nhiều nguồn
func gop(chs ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, c := range chs {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for v := range c { out <- v }
}(c)
}
go func() { wg.Wait(); close(out) }()
return out
}
Mỗi channel vào có một goroutine chuyển tiếp. Goroutine cuối cùng chỉ làm một việc: chờ tất cả xong rồi đóng đầu ra.
Đây chính là mẫu "nhiều bên gửi thì không ai đóng, dùng WaitGroup rồi đóng ở goroutine riêng" ở bài 32.
Nếu chỉ có hai nguồn cố định, select với mẹo gán nil ở bài 33 gọn hơn và không cần goroutine phụ.
Giới hạn số worker
Bài 31 đã nói: goroutine rẻ, nhưng thứ nó dùng thì không. Semaphore bằng channel có đệm:
sem := make(chan struct{}, 10) // tối đa 10 việc cùng lúc
for _, v := range ds {
sem <- struct{}{} // xin chỗ, chặn nếu đủ 10
go func(v int) {
defer func() { <-sem }() // trả chỗ
xuLy(v)
}(v)
}
chan struct{} vì ta chỉ cần đếm, không cần dữ liệu — và bài 11 đã đo struct{} chiếm 0 byte.
Bài 41 sẽ cho thấy errgroup.SetLimit làm việc này gọn hơn và có xử lý lỗi.
Toàn bộ pipeline phải dừng được
Đây là phần quan trọng nhất và cũng hay bị bỏ qua.
Nếu hạ nguồn ngừng đọc — vì lỗi, vì break, vì đủ dữ liệu — mọi giai đoạn thượng nguồn sẽ kẹt ở lệnh gửi và không bao giờ thoát. Đó là rò rỉ goroutine, bài mai sẽ đo.
context giải chuyện này. Với mẫu ở đầu bài:
c2, h2 := context.WithCancel(context.Background())
g := binhPhuong(c2, sinh(c2, 1_000_000))
<-g // chỉ lấy MỘT phần tử
h2() // huỷ -> mọi giai đoạn thoát
Không có h2(), hai goroutine sinh và bình phương sẽ kẹt mãi mãi dù bạn chỉ cần một giá trị.
Quy tắc: hàm nào tạo pipeline thì hàm đó chịu trách nhiệm huỷ nó, thường bằng defer cancel().
Khi nào KHÔNG dùng pipeline
Mẫu này đẹp nhưng không miễn phí. Mỗi giai đoạn là một goroutine cộng một channel, và mỗi lần chuyển giá trị qua channel tốn hơn một lời gọi hàm nhiều lần.
Đừng dựng pipeline cho dữ liệu nhỏ. Một vòng for xử lý một nghìn phần tử trong bộ nhớ luôn nhanh hơn pipeline ba giai đoạn.
Chỉ đáng khi mỗi giai đoạn tốn đáng kể (I/O, tính toán nặng), hoặc khi dữ liệu là luồng vô hạn, hoặc khi bạn cần áp lực ngược tự nhiên.
Đây là kết luận giống bài 76 sê-ri Java về ForkJoinPool: song song hoá có chi phí cố định, và với dữ liệu nhỏ thì tuần tự thắng.
Thử ba mươi giây
ch := sinh(ctx, 1000)
v := <-ch // lấy đúng một giá trị rồi bỏ đi
fmt.Println(runtime.NumGoroutine())
In ra 2 chứ không phải 1 — goroutine sinh vẫn còn, kẹt ở lệnh gửi thứ hai.
Thêm defer cancel() và chạy lại. Ba mươi giây đó là toàn bộ lý do mọi giai đoạn pipeline cần ctx.
Ngày mai: rò rỉ goroutine — cách phát hiện và cách phòng.