プログラミング
Goの並行処理を凝縮
Go Concurrency Distilled (antonz.org)
要約
このミニブックは、Goにおける多くの並行処理トピックの概要を簡潔に提供します。各トピックにはインタラクティブな例が含まれており、コードを変更して実行することで自由に試すことができます。PDF版も用意されています。これはGoの並行処理のクイックリフレッシャーであり、初心者向けガイドではありません。実践的な演習で並行処理をゼロから学びたい場合は、別の書籍「Gist of Go: Concurrency」を参照してください。
全文翻訳
このミニブックは、Goにおける多くの並行処理トピックの概要を簡潔に提供します。各トピックにはインタラクティブな例が含まれており、コードを変更してクリックして実行することで自由に試すことができます。静的な例を含むPDF版もあります。
これはGoの並行処理のクイックリフレッシャーであり、初心者向けガイドではありません。実践的な演習で並行処理をゼロから学びたい場合は、私の別の書籍「Gist of Go: Concurrency」をチェックしてください。
この書籍はAIフリーです。
Goroutines • Channels • Select • Pipelines • Time • Context • Wait groups • Data races • Race conditions • Mutexes • Semaphores • Signaling • Run once • Object pool • Atomics • Testing • Scheduling • Diagnostics • Final thoughts
# Goroutines
Goの並行処理の基盤はgoroutineです。これはgoキーワードで開始される関数です。
func main() {
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
fmt.Println("worker 1")
}()
go func() {
defer wg.Done()
fmt.Println("worker 2")
}()
wg.Wait()
}
worker 2
worker 1
Goランタイムはこれらのgoroutineを管理し、CPUコアで実行されているオペレーティングシステムのスレッドに分散させます。OSスレッドと比較して、goroutineは軽量であるため、数百または数千ものgoroutineを作成できます。
Goroutineは完全に独立しています。main関数もgoroutineですが、プログラムが開始されるときに暗黙的に開始されます。mainが終了すると、他のgoroutineもシャットダウンします。
上記の例では、goroutineが終了するのを待つためにwait group (sync.WaitGroup) を使用します。wait groupには内部にカウンターがあります。Add(n) を呼び出すとnだけインクリメントされ、Done() を呼び出すと1だけデクリメントされます。Wait() は、カウンターがゼロになるまで呼び出し元のgoroutine(この場合はmain)をブロックします。このように、mainは終了する前に両方のワーカーが終了するのを待ちます。
WaitGroup.Goは、wait groupカウンターを自動的にインクリメントし、goroutineで関数を実行し、完了時にカウンターをデクリメントします。
func main() {
var wg sync.WaitGroup
wg.Go(func() {
fmt.Println("worker 1")
})
wg.Go(func() {
fmt.Println("worker 2")
})
wg.Wait()
}
worker 2
worker 1
# Channels
Goroutineはチャネルを介して互いに値を渡すことができます。チャネルは、一方のgoroutineが何かを投げ込み、もう一方がそれをキャッチできるウィンドウのようなものです。
func main() {
messages := make(chan string)
go func() {
messages <- "ping"
}()
msg := <-messages
fmt.Println(msg)
}
ping
チャネルへの値の送信は同期操作です。送信側のgoroutineがチャネルに値を書き込むと(ch <- val)、受信者がその値を受け取るまで(<-ch)ブロックして待機します。その後、処理が続行されます。
出力チャネル
関数から出力チャネルを返し、それを内部のgoroutineで埋めるのは、Goで一般的なパターンです。これにより、呼び出し元はチャネルを通じて値を受け取ることができますが、所有する関数はそのチャネルの制御を維持します。
func generate(start, stop int) chan int {
out := make(chan int)
go func() {
for i := start; i < stop; i++ {
out <- i
}
}()
return out
}
チャネルのクローズ
すべてのデータが送信されたことをリーダーに通知するために、ライターgoroutineはclose()でチャネルをクローズします。
func generate(start, stop int) chan int {
out := make(chan int)
go func() {
defer close(out)
for i := start; i < stop; i++ {
out <- i
}
}()
return out
}
リーダーは、読み取る際に2番目の値(「コンマOK」)でチャネルの状態を確認します。
func main() {
in := generate(5, 10)
for {
num, ok := <-in
if !ok {
break
}
fmt.Print(num, " ")
}
}
5 6 7 8 9
チャネルが開いている間、リーダーは次の値とtrueの状態を受け取ります。チャネルが閉じられると、リーダーはゼロ値とfalseの状態を受け取ります。
チャネルは一度しかクローズできません。再度クローズしたり、閉じられたチャネルに書き込んだりするとパニックが発生します。
チャネルをクローズする唯一の理由は、すべてのデータが送信されたことをリーダーに通知することです。これがリーダーにとって重要でない場合は、クローズする必要はありません。チャネルが使用されなくなると、Goのガベージコレクターは、クローズされているかどうかにかかわらず、そのリソースを解放します。
チャネルのイテレーション
rangeは自動的にチャネルから次の値を取得し、閉じられているかどうかを確認します。チャネルが閉じられている場合、ループを終了します。
func main() {
nums := generate(5, 10)
for n := range nums {
fmt.Print(n, " ")
}
}
5 6 7 8 9
スライスに対するrangeとは異なり、チャネルに対するRangeは単一の値(ペアではない)を返します。
方向性チャネル
チャネルの方向を設定することで、誤った書き込み/クローズエラーから自身を保護できます。チャネルには次のものがあります。
chan(双方向):読み書き用(デフォルト)。
chan<-(送信専用):書き込み専用。
<-chan(受信専用):読み取り専用。
送信専用チャネルから読み取ったり、受信専用チャネルに書き込んだりすることはできません(また、クローズすることもできません)。
チャネルは通常、読み書き両方で初期化され、関数パラメータで方向性として指定されます。Goは通常のチャネルを方向性チャネルに自動的に変換します。
stream := make(chan int)
go func(in chan<- int) {
in <- 42
}(stream)
func(out <-chan int) {
fmt.Println(<-out)
}(stream)
42
バッファ付きチャネル
バッファ付きチャネルは、固定サイズのバッファを持つFIFOキューのように機能し、値を格納します。
バッファに空きスペースがある限り、チャネルへの書き込みはgoroutineをブロックしません。同様に、バッファに値が含まれている限り、チャネルからの読み取りはgoroutineをブロックしません。
stream := make(chan int, 3)
stream <- 11
stream <- 12
stream <- 13
fmt.Println(<-stream)
fmt.Println(<-stream)
11
13
デフォルトでは、バッファサイズを指定しない場合、チャネルはバッファなし(バッファサイズはゼロ)です。
バッファ付きチャネルは、組み込みのlen()およびcap()関数で機能します。
stream := make(chan int, 3)
stream <- 11
fmt.Println(cap(stream), len(stream))
3 1
閉じられたバッファ付きチャネルからの読み取りは、バッファからの値とtrueの状態を返します。すべての値が取得されると、通常のチャネルと同様に、ゼロ値とfalseの状態を返します。
stream := make(chan int, 1)
stream <- 11
close(stream)
val, ok := <-stream
fmt.Println(val, ok) // 11 true
val, ok = <-stream
fmt.Println(val, ok) // 0 false
11 true
0 false
nilチャネル
Goの他の型と同様に、チャネルにはゼロ値があり、それはnilです。
nilチャネルへの書き込みまたはnilチャネルからの読み取りは、goroutineを無期限にブロックします。
var stream chan int
go func() {
// 無限にブロックします
stream <- 1
}()
// 無限にブロックします
<-stream
nilチャネルのクローズはパニックを引き起こします。
var stream chan int
close(stream) // panic: close of nil channel
# Select
selectステートメントはswitchに似ていますが、チャネル専用に設計されています。その機能は次のとおりです。
ブロックされていないケースをチェックします。
複数のケースが準備できている場合、実行するケースをランダムに選択します。
すべてのケースがブロックされており、defaultケースがある場合、それを実行します。
すべてのケースがブロックされており、defaultケースがない場合、いずれかが準備できるまで待機します。
Selectはパイプラインのデータフローを管理するために使用されます。
// mergeはin1とin2からの値をoutputチャネルに送信します。
func merge(in1, in2 <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for in1 != nil || in2 != nil {
select {
case val1, ok := <-in1:
if ok {
out <- val1
} else {
in1 = nil
}
case val2, ok := <-in2:
if ok {
out <- val2
} else {
in2 = nil
}
}
}
}()
return out
}
// in1に10..12を、in2に20..22を送信し、
// merge(in1, in2)を呼び出すと仮定します。
10 11 20 12 21 22
goroutineをキャンセルするには:
// processはinからの値を変更し、inが枯渇するかcancelが閉じられるまでoutに送信します。
func process(cancel chan struct{}, in <-chan int) <-chan int {
out := make(chan int)
go func() {
for val := range in {
select {
case out <- val*10:
case <-cancel:
fmt.Println("canceled")
return
}
}
}()
return out
}
// inに値11と12を送信し、
// その後close(cancel)を呼び出すと仮定します。
110
120
canceled
非ブロッキング操作の場合:
// multiplierは、入力を10倍してチャネルに送信するか、
// チャネルがビジーの場合はエラーを返す関数を返します。
func multiplier(ch chan<- int) func(n int) error {
return func(n int) error {
select {
case ch <- n*10:
return nil
default:
return errors.New("busy")
}
}
}
func main() {
nums := make(chan int, 1)
multiply := multiplier(nums)
err := multiply(11)
fmt.Println(<-nums, err) // 110 <nil>
err = multiply(12)
fmt.Println(<-nums, err) // 120 <n