Hình dung một dây chuyền lắp ráp trong nhà máy. Mỗi trạm làm một việc rồi đặt sản phẩm lên băng chuyền cho trạm sau — đó là pipeline, và băng chuyền chính là channel. Cần làm nhanh hơn? Đặt bốn người thợ giống nhau cùng bốc việc từ một băng chuyền, ai rảnh thì nhặt — đó là fan-out. Nhiều băng chuyền đổ về một băng gộp lại — đó là fan-in. Ba mẫu này chiếm phần lớn mã đồng thời trong Go thật, và 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 — một trạm có băng chuyền vào và băng chuyền 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 — đúng như mọi trạm trên dây chuyền làm cùng lúc trên những sản phẩm khác nhau.

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. Bốn người thợ cùng một băng chuyền, người nào xong tay thì nhặt món kế, không ai ngồi chờ phần đã chia sẵn cho mình.

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. Băng chuyền đầy ứ, mọi trạm phía trên đứng chôn chân với món hàng trên tay. Đó là rò rỉ goroutine, bài mai sẽ đo.

context là cái dây dừng khẩn chạy dọc cả dây chuyền. 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.

Muốn thấy cái rò rỉ tận mắt trong chục giây, lấy một giá trị rồi đếm goroutine còn sống:

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. Con số 2-đáng-ra-là-1 đó là toàn bộ lý do mọi giai đoạn pipeline cần ctx.

Mẫu số chung

Ba mẫu này không phải phát minh của Go — chúng là hình dạng muôn thuở của xử lý luồng, và tổ tiên trực tiếp là ống dẫn Unix: tao | loc | dem, mỗi lệnh một tiến trình, dấu | là một bộ đệm có giới hạn. Channel của Go chỉ là cái ống đó, có kiểu và nằm trong một tiến trình.

  • Unix pipe cho bạn pipeline + áp lực ngược miễn phí: người đọc chậm thì người ghi bị chặn.
  • Java có Stream.parallel(), ForkJoinPool (bài 76), và Reactive Streams (RxJava, Reactor) — vốn chính là fan-out/fan-in với áp lực ngược nâng thành một giao thức hẳn hoi.
  • Kafka: một consumer group chính là fan-out — mỗi partition tới đúng một consumer, tự cân bằng theo tốc độ, không chia trước; giống hệt "worker nào rảnh thì nhặt" ở trên.
  • Node có stream.pipe() với highWaterMark; Rust có mpsc cộng rayon; Erlang/Akka có actor với mailbox.

Hai hằng số đáng mang theo. Một: vấn đề thật không phải "bao nhiêu worker" mà là áp lực ngược — một nguồn nhanh gặp một đích chậm thì phải có bộ đệm giới hạn ghép chúng lại, nếu không nguồn phải chặn hoặc phải vứt dữ liệu; channel không đệm của Go cho bạn điều đó không mất công (lệnh gửi tự chặn). Hai: một luồng đồng thời bắt buộc phải có đường tháo dỡ sạch — Unix bắn SIGPIPE khi người đọc đóng, Go dùng context, Reactive Streams có cancel — thiếu nó là rò rỉ. Sợi chỉ chung: khi bạn nối các trạm bằng hàng đợi, câu hỏi thiết kế không phải "chạy mấy luồng" mà là "đầu chậm thì sao, và cả cỗ máy tắt thế nào" — thông lượng tự lo được, còn áp lực ngược và tháo dỡ thì không.

Ngày mai: rò rỉ goroutine — cách phát hiện và cách phòng.

Bài tập làm thử

Bài 1 (đọc hiểu). Đoạn mã sau lấy đúng một giá trị từ pipeline rồi bỏ đi, không huỷ context. Sau khi chạy, runtime.NumGoroutine() in ra bao nhiêu, và tại sao?

ch := sinh(ctx, 1000)
v := <-ch
fmt.Println(runtime.NumGoroutine())
Đáp án

In ra 2 chứ không phải 1. Goroutine bên trong sinh đã gửi được giá trị đầu tiên qua channel, nhưng vì không ai đọc tiếp và không có ctx bị huỷ, nó kẹt mãi ở lệnh gửi giá trị thứ hai (out <- i) trong vòng lặp select. Đây là rò rỉ goroutine — con số "2 đáng ra là 1" đó chính là lý do mọi giai đoạn pipeline cần nhận và kiểm ctx.Done().

Bài 2 (sửa lỗi). Hàm gop sau đây dùng để fan-in nhiều channel, nhưng thiếu một bước quan trọng khiến channel đầu ra out không bao giờ được đóng, làm range phía người đọc bị treo mãi. Sửa lại.

func gop(chs ...<-chan int) <-chan int {
	out := make(chan int)
	for _, c := range chs {
		go func(c <-chan int) {
			for v := range c { out <- v }
		}(c)
	}
	return out
}
Đáp á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
}

Thiếu sync.WaitGroup: với nhiều goroutine cùng ghi vào out, không ai được phép tự ý đóng out (đóng hai lần sẽ panic), nên phải chờ tất cả các goroutine chuyển tiếp xong bằng wg.Wait() trong một goroutine riêng, rồi mới close(out).

Bài 3 (vận dụng thực tế). Bạn có 100 phần tử cần xử lý bằng 4 worker chạy song song, đọc từ cùng một channel nguồn. Không dùng errgroup, hãy viết mã fan-out theo đúng mẫu trong bài (nhiều goroutine đọc chung một channel).

Đáp án
ng := sinh(ctx, 100)
var w []<-chan int
for i := 0; i < 4; i++ {
	w = append(w, binhPhuong(ctx, ng))
}

Bốn worker (bốn lệnh gọi binhPhuong) cùng nhận một channel nguồn ng. Go tự chia việc: mỗi giá trị đi tới đúng một worker, và worker nào rảnh thì nhận — tải tự cân bằng theo tốc độ xử lý thật, không cần phân mảnh dữ liệu trước hay tự cân bằng tải thủ công.

Bài 4 (bẫy/đánh đổi). Một pipeline ba giai đoạn (mỗi giai đoạn một goroutine + một channel) được dùng để cộng dồn một mảng chỉ có 20 phần tử trong bộ nhớ. Đây có phải là thiết kế tốt không? Giải thích dựa trên nội dung bài.

Đáp án

Không nên. Pipeline 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 nhiều so với một lời gọi hàm thông thường. Với dữ liệu nhỏ trong bộ nhớ (20 phần tử), một vòng for tuần tự luôn nhanh hơn pipeline ba giai đoạn. Pipeline chỉ đáng dùng khi mỗi giai đoạn tốn đáng kể (I/O, tính toán nặng), khi dữ liệu là luồng vô hạn, hoặc khi cần áp lực ngược tự nhiên.

Bài 5 (đọc hiểu — bốn quy ước của một giai đoạn). Bài viết nêu bốn quy ước để một giai đoạn pipeline đúng: (1) trả về channel chỉ-nhận, (2) defer close(out), (3) select với ctx.Done() ở mọi lệnh gửi, (4) đọc bằng range. Giả sử một giai đoạn bỏ qua quy ước (3) — gửi trực tiếp out <- v không qua select. Hậu quả gì xảy ra nếu hạ nguồn ngừng đọc giữa chừng (ví dụ do lỗi và break sớm)?

Đáp án

Giai đoạn đó sẽ kẹt vĩnh viễn ở lệnh gửi out <- v, vì không còn ai đọc từ out nữa và channel không đệm sẽ chặn goroutine mãi mãi — đúng như mô tả "băng chuyền đầy ứ, mọi trạm phía trên đứng chôn chân với món hàng trên tay". Đây là rò rỉ goroutine. Có select với ctx.Done() mới cho phép giai đoạn thoát ra khi context bị huỷ, ngay cả khi không gửi được giá trị.