Các mô hình đồng thời Go nâng cao

Table of Contents
Mô hình đồng thời của Go, được xây dựng xung quanh goroutine và channel, cung cấp một cách mạnh mẽ và tinh tế để xây dựng các hệ thống có tính đồng thời cao. Mặc dù những điều cơ bản về việc khởi tạo một goroutine (go func()) và giao tiếp qua channel đã được hầu hết các nhà phát triển Go hiểu rõ, việc nắm vững các mẫu đồng thời nâng cao là rất quan trọng để xây dựng các ứng dụng mạnh mẽ, có khả năng mở rộng và hiệu suất cao. Trong bài viết chuyên sâu này, chúng ta sẽ khám phá một số mẫu đồng thời Go nâng cao, phân tích cách triển khai và thảo luận về thời điểm cũng như lý do sử dụng chúng.
1. Mẫu Pipeline
Mẫu pipeline là nền tảng của xử lý dữ liệu đồng thời trong Go. Nó bao gồm một chuỗi các giai đoạn được kết nối bằng channel, trong đó mỗi giai đoạn là một nhóm goroutine chạy cùng một hàm. Trong mỗi giai đoạn, các goroutine:
- Nhận giá trị từ upstream qua các channel đầu vào.
- Thực hiện một số chức năng trên dữ liệu đó, thường tạo ra các giá trị mới.
- Gửi giá trị downstream qua các channel đầu ra.
Triển khai một Pipeline
Một pipeline mạnh mẽ đòi hỏi phải quản lý goroutine cẩn thận, đặc biệt là liên quan đến việc hủy bỏ và xử lý lỗi. Hãy cùng xây dựng một pipeline đa giai đoạn để tạo số, bình phương chúng, sau đó lọc ra các kết quả chẵn.
package main
import (
"context"
"fmt"
"sync"
)
// Stage 1: Generator
func generator(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}()
return out
}
// Stage 2: Squarer
func squarer(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
select {
case out <- n * n:
case <-ctx.Done():
return
}
}
}()
return out
}
// Stage 3: Filter
func filterEven(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
if n%2 == 0 {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}
}()
return out
}
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
nums := []int{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
genChan := generator(ctx, nums...)
sqChan := squarer(ctx, genChan)
filterChan := filterEven(ctx, sqChan)
for val := range filterChan {
fmt.Println(val)
}
}
Pipeline này đảm bảo rằng nếu context bị hủy, tất cả các goroutine sẽ thoát một cách duyên dáng, ngăn chặn rò rỉ tài nguyên. Việc sử dụng rõ ràng context.Context là một phương pháp hay nhất cho tính đồng thời Go nâng cao.
2. Fan-Out và Fan-In
Khi một giai đoạn duy nhất trong pipeline tốn nhiều tính toán, nó có thể trở thành nút thắt cổ chai. Mẫu fan-out/fan-in giải quyết vấn đề này bằng cách phân phối công việc trên nhiều goroutine (fan-out) và sau đó ghép nối kết quả của chúng trở lại một channel duy nhất (fan-in).
Fan-Out
Fan-out đơn giản là khởi tạo nhiều goroutine để đọc từ cùng một channel. Channel hoạt động như một hàng đợi công việc, phân phối các mục giữa các worker.
Fan-In
Fan-in yêu cầu một hàm multiplexer nhận nhiều channel và trả về một channel duy nhất. Nó sử dụng một sync.WaitGroup để theo dõi khi tất cả các channel đầu vào đã được đóng.
func merge(ctx context.Context, cs ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
output := func(c <-chan int) {
defer wg.Done()
for n := range c {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}
wg.Add(len(cs))
for _, c := range cs {
go output(c)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
Bằng cách kết hợp các giai đoạn pipeline với fan-out/fan-in, chúng ta có thể tăng đáng kể thông lượng cho các tác vụ bị giới hạn bởi CPU.
3. Mẫu Worker Pool
Mẫu worker pool rất cần thiết để giới hạn tính đồng thời, điều này rất quan trọng khi tương tác với các tài nguyên có giới hạn tốc độ nghiêm ngặt hoặc khi cố gắng tránh tiêu thụ bộ nhớ không giới hạn.
Không giống như việc tạo goroutine không giới hạn, một worker pool duy trì một số lượng goroutine cố định tiêu thụ các tác vụ từ một channel dùng chung.
type Task struct {
ID int
Value string
}
type Result struct {
TaskID int
Output string
}
func worker(id int, tasks <-chan Task, results chan<- Result) {
for t := range tasks {
// Simulate work
output := fmt.Sprintf("Worker %d processed Task %d: %s", id, t.ID, t.Value)
results <- Result{TaskID: t.ID, Output: output}
}
}
func main() {
numWorkers := 3
numTasks := 10
tasks := make(chan Task, numTasks)
results := make(chan Result, numTasks)
// Start workers
for w := 1; w <= numWorkers; w++ {
go worker(w, tasks, results)
}
// Send tasks
for j := 1; j <= numTasks; j++ {
tasks <- Task{ID: j, Value: fmt.Sprintf("Data %d", j)}
}
close(tasks)
// Collect results
for a := 1; a <= numTasks; a++ {
<-results
}
}
Mẫu này kiểm soát số lượng hoạt động đồng thời tối đa, cung cấp sự ổn định và hiệu suất dự đoán được dưới tải.
4. Context để Hủy bỏ và Hết thời gian
Trong bất kỳ hệ thống đồng thời phức tạp nào, khả năng hủy bỏ các hoạt động đang diễn ra là rất quan trọng. Gói context là cách tiêu chuẩn để truyền tín hiệu hủy bỏ và thời hạn qua các ranh giới API và giữa các goroutine.
Khi một hoạt động vượt quá thời gian chờ được chỉ định hoặc khi người dùng hủy một yêu cầu, context sẽ bị hủy. Mọi goroutine liên quan đến hoạt động phải định kỳ kiểm tra ctx.Done() hoặc sử dụng câu lệnh select bao gồm <-ctx.Done().
func fetchResource(ctx context.Context, url string) (string, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return "", err
}
client := http.DefaultClient
resp, err := client.Do(req)
if err != nil {
return "", err
}
defer resp.Body.Close()
// process response...
return "Success", nil
}
Điều này đảm bảo rằng các yêu cầu mạng không bị treo vô thời hạn và tài nguyên được giải phóng kịp thời khi không còn cần thiết.
5. Mẫu Semaphore
Đôi khi bạn cần kiểm soát quyền truy cập vào một tài nguyên cụ thể, giới hạn số lượng goroutine có thể truy cập đồng thời, độc lập với số lượng worker đang hoạt động. Trong khi một sync.Mutex chỉ cho phép một goroutine truy cập tài nguyên tại một thời điểm, một semaphore cho phép tối đa N goroutine.
Trong Go, một channel có bộ đệm thường được sử dụng làm semaphore.
func main() {
concurrencyLimit := 5
semaphore := make(chan struct{}, concurrencyLimit)
var wg sync.WaitGroup
for i := 0; i < 20; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
// Acquire semaphore
semaphore <- struct{}{}
// Critical section
fmt.Printf("Goroutine %d executing\n", id)
time.Sleep(time.Second) // Simulate work
// Release semaphore
<-semaphore
}(i)
}
wg.Wait()
}
Mẫu này rất hiệu quả để điều tiết các yêu cầu đến các API của bên thứ ba hoặc quản lý các kết nối đến các cơ sở dữ liệu cũ không thể xử lý tính đồng thời cao.
6. Mẫu Or-Done Channel
Khi xử lý nhiều channel hoặc khi gói các hàm mà bạn không kiểm soát, bạn thường cần một cách để đảm bảo rằng bạn không làm rò rỉ goroutine nếu việc hủy bỏ xảy ra trước khi một hoạt động channel hoàn tất. Mẫu or-done giải quyết vấn đề này một cách thanh lịch bằng cách gói hoạt động đọc.
func orDone(ctx context.Context, c <-chan interface{}) <-chan interface{} {
valStream := make(chan interface{})
go func() {
defer close(valStream)
for {
select {
case <-ctx.Done():
return
case v, ok := <-c:
if !ok {
return
}
select {
case valStream <- v:
case <-ctx.Done():
}
}
}
}()
return valStream
}
Điều này cho phép bạn viết các vòng lặp sạch hơn khi tiêu thụ channel, trừu tượng hóa phần mã lặp đi lặp lại của việc kiểm tra ctx.Done() trên mỗi lần đọc.
7. Mẫu Tee Channel
Đôi khi bạn cần chia một luồng dữ liệu bắt nguồn từ một channel duy nhất thành hai luồng riêng biệt, cho phép các đường xử lý độc lập. Đây được gọi là tee channel (tương tự như lệnh tee trong Unix).
func tee(ctx context.Context, in <-chan interface{}) (<-chan interface{}, <-chan interface{}) {
out1 := make(chan interface{})
out2 := make(chan interface{})
go func() {
defer close(out1)
defer close(out2)
for val := range orDone(ctx, in) {
// Create local copies of channels to track which ones we've sent to
var out1Local, out2Local = out1, out2
for i := 0; i < 2; i++ {
select {
case <-ctx.Done():
return
case out1Local <- val:
out1Local = nil // Prevent sending to this channel again
case out2Local <- val:
out2Local = nil // Prevent sending to this channel again
}
}
}
}()
return out1, out2
}
Tee channel đảm bảo rằng mỗi giá trị nhận được từ channel in được gửi đến cả out1 và out2 trước khi giá trị tiếp theo được xử lý. Đặt tham chiếu channel cục bộ thành nil sau khi gửi thành công là một thủ thuật thông minh trong Go; gửi đến một channel nil sẽ chặn mãi mãi, nghĩa là câu lệnh select sẽ chỉ đánh giá các trường hợp không nil còn lại cho phần còn lại của vòng lặp. Mẫu này đặc biệt hữu ích để nhân đôi các luồng sự kiện, chẳng hạn như gửi dữ liệu đo từ xa đến một dịch vụ đo lường trong khi đồng thời ghi nhật ký vào một tệp.
8. ErrGroup cho các tác vụ phối hợp
Trong khi sync.WaitGroup rất tuyệt vời để chờ một tập hợp các goroutine, nó lại không hiệu quả khi bạn cần xử lý lỗi hoặc hủy bỏ các tác vụ còn lại nếu một tác vụ thất bại. Gói golang.org/x/sync/errgroup cung cấp đồng bộ hóa, truyền lỗi và hủy Context cho các nhóm goroutine làm việc trên các tác vụ con của một tác vụ chung.
package main
import (
"context"
"fmt"
"golang.org/x/sync/errgroup"
)
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
g, ctx := errgroup.WithContext(ctx)
urls := []string{
"http://example.com",
"http://example.net",
"http://example.org",
}
for _, url := range urls {
// Launch a goroutine to fetch the URL.
url := url // https://golang.org/doc/faq#closures_and_goroutines
g.Go(func() error {
// Use the context returned by errgroup.WithContext.
// If any goroutine returns an error, this context will be canceled.
// return fetchResource(ctx, url)
fmt.Printf("Fetching %s\n", url)
return nil
})
}
// Wait for all HTTP fetches to complete.
if err := g.Wait(); err != nil {
fmt.Printf("Failed to fetch all URLs: %v\n", err)
} else {
fmt.Println("Successfully fetched all URLs.")
}
}
Sử dụng errgroup sạch hơn đáng kể so với việc quản lý một channel lỗi riêng biệt và logic hủy bỏ thủ công. Đây là cách tiếp cận được khuyến nghị bất cứ khi nào bạn có nhiều tác vụ đồng thời là một phần của một hoạt động logic lớn hơn, trong đó sự thành công của toàn bộ phụ thuộc vào sự thành công của các phần của nó.
Kết luận
Nắm vững tính đồng thời Go nâng cao đòi hỏi phải vượt ra ngoài từ khóa go cơ bản và suy nghĩ sâu sắc về luồng dữ liệu, quản lý tài nguyên và kiểm soát vòng đời. Bằng cách sử dụng các mẫu như pipeline, fan-out/fan-in, worker pool, hủy context mạnh mẽ và semaphore, bạn có thể thiết kế các hệ thống không chỉ có tính đồng thời cao mà còn kiên cường, dễ dự đoán và dễ bảo trì. Những mẫu này tạo thành nền tảng kiến trúc của các ứng dụng Go hiệu suất cao trong sản xuất.
Bạn cũng có thể thích
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

gRPC vs ConnectRPC: Microservices hiện đại và Protobuf gốc trình duyệt
Đánh giá kiến trúc gRPC vs ConnectRPC trong TypeScript và Go, khám phá streaming HTTP/1.1 vs HTTP/2, client trình duyệt không cần proxy Envoy và độ trễ RPC p99.
Read more
Điện toán lượng tử cho nhà phát triển
Hướng dẫn dành cho nhà phát triển về điện toán lượng tử: viết thuật toán lượng tử với Qiskit, hiểu các cổng lượng tử và mô phỏng mạch trên phần cứng cổ điển.
Read more
Distributed Tracing với OpenTelemetry
Theo dõi vi dịch vụ bằng OpenTelemetry distributed tracing: theo dõi sự lan truyền ngữ cảnh giữa các dịch vụ, các điểm nghẽn độ trễ và xuất sang Jaeger.
Read more