•14 min read

Goの高度な並行処理パターン

Goの高度な並行処理パターン

Goの並行処理モデルは、ゴルーチンとチャネルを中心に構築されており、高度に並行なシステムを構築するための強力かつ洗練された方法を提供します。ゴルーチン(go func())の起動やチャネルを介した通信の基本はほとんどのGo開発者に理解されていますが、堅牢でスケーラブル、かつ高性能なアプリケーションを構築するには、高度な並行処理パターンを習得することが不可欠です。この詳細な解説では、いくつかの高度なGo並行処理パターンを探求し、その実装を分析し、いつ、なぜそれらを使用すべきかを議論します。

Audio Briefing
0:00 / 0:00

1. パイプラインパターン

パイプラインパターンは、Goにおける並行データ処理の要です。これは、チャネルで接続された一連のステージで構成され、各ステージは同じ関数を実行するゴルーチンのグループです。各ステージで、ゴルーチンは次のことを行います。

  • インバウンドチャネルを介してアップストリームから値を受け取ります。
  • そのデータに対して何らかの関数を実行し、通常は新しい値を生成します。
  • アウトバウンドチャネルを介してダウンストリームに値を送信します。

パイプラインの実装

堅牢なパイプラインには、特にキャンセルとエラー処理に関して、ゴルーチンの慎重な管理が必要です。数値を生成し、それらを二乗し、偶数の結果をフィルタリングする多段階パイプラインを構築してみましょう。

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)
	}
}

このパイプラインは、コンテキストがキャンセルされた場合、すべてのゴルーチンが正常に終了し、リソースリークを防ぐことを保証します。context.Contextを明示的に使用することは、高度なGo並行処理におけるベストプラクティスです。

Advertisement

2. ファンアウトとファンイン

パイプラインの単一ステージが計算コストの高い場合、ボトルネックになる可能性があります。ファンアウト/ファンインパターンは、作業を複数のゴルーチンに分散させ(ファンアウト)、その結果を単一のチャネルに多重化し直す(ファンイン)ことで、この問題を解決します。

ファンアウト

ファンアウトは、同じチャネルから読み取るために複数のゴルーチンを起動するだけです。チャネルはワークキューとして機能し、ワーカー間でアイテムを分散します。

ファンイン

ファンインには、複数のチャネルを受け取り、単一のチャネルを返すマルチプレクサ関数が必要です。すべての入力チャネルが閉じられたときに追跡するためにsync.WaitGroupを使用します。

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
}

パイプラインステージとファンアウト/ファンインを組み合わせることで、CPUバウンドなタスクのスループットを劇的に向上させることができます。

3. ワーカープールパターン

ワーカープールパターンは、並行処理を制限するために不可欠です。これは、厳格なレート制限を持つリソースとやり取りする場合や、無制限のメモリ消費を避けようとする場合に重要です。

無制限のゴルーチン作成とは異なり、ワーカープールは、共有チャネルからタスクを消費する固定数のゴルーチンを維持します。

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
	}
}

このパターンは、同時操作の最大数を制御し、負荷の下での安定性と予測可能なパフォーマンスを提供します。

4. キャンセルとタイムアウトのためのコンテキスト

洗練された並行システムでは、実行中の操作をキャンセルする機能が不可欠です。contextパッケージは、API境界を越えてゴルーチン間でキャンセルシグナルとデッドラインを伝播するための標準的な方法です。

操作が指定されたタイムアウトを超えた場合、またはユーザーがリクエストを中止した場合、コンテキストはキャンセルされるべきです。操作に関与するすべてのゴルーチンは、定期的にctx.Done()をチェックするか、<-ctx.Done()を含むselectステートメントを使用する必要があります。

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
}

これにより、ネットワークリクエストが無期限にハングアップせず、不要になったリソースが速やかに解放されます。

Advertisement

5. セマフォパターン

アクティブなワーカーの数に関係なく、特定のリソースへのアクセスを制御し、同時にアクセスできるゴルーチンの数を制限する必要がある場合があります。sync.Mutexは一度に1つのゴルーチンのみがリソースにアクセスできるようにしますが、セマフォは最大N個のゴルーチンを許可します。

Goでは、バッファ付きチャネルがセマフォとして一般的に使用されます。

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()
}

このパターンは、サードパーティAPIへのリクエストをスロットリングしたり、高並行処理を処理できないレガシーデータベースへの接続を管理したりするのに非常に効果的です。

6. Or-Doneチャネルパターン

複数のチャネルを扱ったり、制御できない関数をラップしたりする場合、チャネル操作が完了する前にキャンセルが発生した場合にゴルーチンをリークさせないようにする方法が必要になることがよくあります。or-doneパターンは、読み取り操作をラップすることでこれをエレガントに解決します。

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
}

これにより、チャネルを消費する際に、ctx.Done()を毎回チェックする定型コードを抽象化し、よりクリーンなループを記述できます。

7. Teeチャネルパターン

単一のチャネルから発信されるデータストリームを2つの別々のストリームに分割し、独立した処理パスを可能にする必要がある場合があります。これはティーチャネル(Unixのteeコマンドに類似)として知られています。

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
}

ティーチャネルは、inチャネルから受信した各値が、次の値が処理される前にout1とout2の両方に送信されることを保証します。成功した送信後にローカルチャネル参照をnilに設定することは、Goにおける巧妙なトリックです。nilチャネルへの送信は永久にブロックされるため、selectステートメントはループの残りのイテレーションで残りの非nilケースのみを評価します。このパターンは、テレメトリーデータをメトリクスサービスに送信しながら同時にファイルにログを記録するなど、イベントストリームを複製するのに非常に役立ちます。

8. 協調タスクのためのErrGroup

sync.WaitGroupはゴルーチンのコレクションを待機するのに優れていますが、エラーを処理したり、1つのタスクが失敗した場合に残りのタスクをキャンセルしたりする必要がある場合には不十分です。golang.org/x/sync/errgroupパッケージは、共通のタスクのサブタスクに取り組むゴルーチンのグループに対して、同期、エラー伝播、およびContextキャンセルを提供します。

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.")
	}
}

errgroupを使用すると、個別のエラーチャネルとキャンセルロジックを手動で管理するよりも大幅にクリーンになります。これは、全体としての成功がその部分の成功に依存する、より大きな論理操作の一部である複数の並行タスクがある場合に推奨されるアプローチです。

結論

高度なGo並行処理を習得するには、基本的なgoキーワードを超えて、データフロー、リソース管理、およびライフサイクル制御について深く考える必要があります。パイプライン、ファンアウト/ファンイン、ワーカープール、堅牢なコンテキストキャンセル、セマフォなどのパターンを採用することで、高度に並行であるだけでなく、回復力があり、予測可能で、保守可能なシステムを設計できます。これらのパターンは、本番環境における高性能Goアプリケーションのアーキテクチャの基盤を形成します。

こちらもおすすめ

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement
開発者のための量子コンピューティング
tech

開発者のための量子コンピューティング

Qiskitで量子アルゴリズムを記述し、量子ゲートを理解し、古典的なハードウェアで回路をシミュレートする方法を解説する、開発者向けの量子コンピューティングガイド。

Read more
OpenTelemetryによる分散トレーシング
tech

OpenTelemetryによる分散トレーシング

OpenTelemetry分散トレーシングでマイクロサービスを計測し、サービス間のコンテキスト伝播、レイテンシーのボトルネックをトレースし、Jaegerにエクスポートします。

Read more