Go: GoのSelectおよび並行処理パターン:多重化、タイムアウト制御、ファンイン/ファンアウト、パイプライン
selectは、Go言語における並行プログラミングの「スイスアーミーナイフ」のような存在です。複数のチャネルを同時に監視し、最初に処理可能な状態になったチャネルを処理します。これは、ゴルーチンをつなぐ「スマートスイッチ」としての役割を果たします。
チャネルをゴルーチンをつなぐ電話回線だと考えると、selectは電話交換手のようなものです。すべての着信を同時に監視し、鳴った電話に応答します。このレッスンでは、selectの主要なパターンをすべて習得します。
1. 学習内容
selectマルチプレクシングの構文と動作- 無作為抽出(複数の事例が同時に準備されている場合)
time.Afterタイムアウト制御default: ノンブロッキング操作for-selectループ終了(処理済みチャンネル)- ファンアウトモードを使用してタスクを割り当てる
- ファンインモードにおける集約結果
- パイプラインモデルを用いてデータ処理パイプラインを構築する
2. クオンツ・トレーディング・エンジニアの実話
(1) 課題:3つのデータソースからデータを取得すると、CPUがフル稼働してしまう
ボブは、クオンツ・トレーディングチームのバックエンドエンジニアです。彼は、3つの株価データソースを同時に監視する必要があります:
「当社の戦略では、NYSE、NASDAQ、LSEの3つの取引所から同時にリアルタイムの株価情報を取得する必要があります。以前のJava実装では、1つのスレッドで3つのWebSocketをポーリングしており、そのたびに10ミリ秒のビジーウェイトが発生し、CPUの30%を消費していました。上司からは、『取引システム本体よりも多くのリソースを使っている』と言われたほどです。」
彼は現在のコードを開いた:
// Bad code: busy polling
func pollDataSources() {
for {
// Poll every 10ms, wasting CPU
if data1 := pollNYSE(); data1 != nil {
process(data1)
}
if data2 := pollNASDAQ(); data2 != nil {
process(data2)
}
if data3 := pollLSE(); data3 != nil {
process(data3)
}
time.Sleep(10 * time.Millisecond)
}
}
3つの問題点がある:(1) 10ミリ秒のポーリング間隔はCPUリソースを浪費する;(2) データの到着と処理が同期していない;(3) 0~10ミリ秒のポーリング遅延は予測できない。
(2) Goの解決策:selectを使用して3つのチャンネルを聞く
// market_data.go
package main
import (
"fmt"
"math/rand"
"time"
)
type Quote struct {
Source string
Symbol string
Price float64
}
func simulateExchange(name string, out chan<- Quote) {
symbols := []string{"AAPL", "GOOGL", "MSFT", "AMZN"}
for {
time.Sleep(time.Duration(50+rand.Intn(200)) * time.Millisecond)
quote := Quote{
Source: name,
Symbol: symbols[rand.Intn(len(symbols))],
Price: 100 + rand.Float64()*200,
}
out <- quote
}
}
func main() {
nyse := make(chan Quote)
nasdaq := make(chan Quote)
lse := make(chan Quote)
go simulateExchange("NYSE", nyse)
go simulateExchange("NASDAQ", nasdaq)
go simulateExchange("LSE", lse)
// select listens to all three channels simultaneously
timeout := time.After(2 * time.Second)
for {
select {
case q := <-nyse:
fmt.Printf("[NYSE] %s: $%.2f\n", q.Symbol, q.Price)
case q := <-nasdaq:
fmt.Printf("[NASDAQ] %s: $%.2f\n", q.Symbol, q.Price)
case q := <-lse:
fmt.Printf("[LSE] %s: $%.2f\n", q.Symbol, q.Price)
case <-timeout:
fmt.Println("Demo ended")
return
}
}
}
出力:
[NASDAQ] AMZN: $198.32
[NYSE] AAPL: $150.45
[LSE] MSFT: $287.10
...
(3) パフォーマンス:SELECT Listen 対 ポーリング
| 手法 | CPU 使用率 | 応答遅延 | コードの複雑度 |
|---|---|---|---|
| ポーリング 10 ms | 30% | 0~10 ms | 低 |
| ポーリング間隔 100 ms | 3% | 0~100 ms | 低 |
| 「Listen」を選択 | 0% | 0 ms (リアルタイム) | 中程度 |
select + channel は イベント駆動型 です。データがないときはゴルーチンが(CPU を消費せずに)スリープ状態になり、データが到着するとランタイムによって呼び起こされます。これこそが、Go の並行処理モデルを非常に効率的なものにしている理由です。
3. SELECT の基礎
(1) select文の構文
select {
case v := <-ch1:
// ch1 is ready
case v := <-ch2:
// ch2 is ready
case ch3 <- value:
// Can send to ch3 (ch3 has space or has a receiver)
default:
// All channels are not ready
}
▶ サンプル:選択(ランダム選択)
package main
import (
"fmt"
)
func main() {
ch1 := make(chan string, 1)
ch2 := make(chan string, 1)
ch1 <- "from ch1"
ch2 <- "from ch2"
// Both channels are ready; select picks one at random
for i := 0; i < 2; i++ {
select {
case msg := <-ch1:
fmt.Println(msg)
case msg := <-ch2:
fmt.Println(msg)
}
}
}
出力(毎回異なる場合があります):
from ch1
from ch2
4. タイムアウト制御
(1) タイムアウトモード
package main
import (
"fmt"
"time"
)
func slowOperation() string {
time.Sleep(2 * time.Second)
return "result"
}
func main() {
ch := make(chan string)
go func() {
ch <- slowOperation()
}()
select {
case result := <-ch:
fmt.Println("Success:", result)
case <-time.After(1 * time.Second):
fmt.Println("Timeout! Operation exceeded 1 second")
}
}
出力:
Timeout! Operation exceeded 1 second
▶ サンプル:段階的なタイムアウト(段階的な待機)
package main
import (
"fmt"
"time"
)
func fetchFromCache() string {
time.Sleep(50 * time.Millisecond)
return "cache-data"
}
func fetchFromDB() string {
time.Sleep(200 * time.Millisecond)
return "db-data"
}
func fetchFromAPI() string {
time.Sleep(500 * time.Millisecond)
return "api-data"
}
func main() {
cache := make(chan string)
db := make(chan string)
api := make(chan string)
go func() { cache <- fetchFromCache() }()
go func() { db <- fetchFromDB() }()
go func() { api <- fetchFromAPI() }()
select {
case r := <-cache:
fmt.Println("Cache hit:", r)
case <-time.After(100 * time.Millisecond):
select {
case r := <-db:
fmt.Println("DB returned:", r)
case <-time.After(300 * time.Millisecond):
select {
case r := <-api:
fmt.Println("API returned:", r)
case <-time.After(600 * time.Millisecond):
fmt.Println("All data sources timed out!")
}
}
}
}
time.After(d) は <-chan time.Time を返し、d 秒後に現在の時刻を送信します。select が time.After を呼び出すたびに、新しいタイマーが作成されます。ループ内で頻繁に呼び出す場合は、リソースリークを防ぐために、必ず time.NewTimer を使用してタイマーを停止してください。
5. デフォルト:ノンブロッキング操作
(1) ノンブロッキング送受信
package main
import (
"fmt"
)
func main() {
ch := make(chan int, 1)
// Non-blocking receive
select {
case v := <-ch:
fmt.Println("Received:", v)
default:
fmt.Println("No data (non-blocking)")
}
// Non-blocking send
ch <- 1
select {
case ch <- 2:
fmt.Println("Send successful")
default:
fmt.Println("Buffer full (non-blocking)")
}
}
出力:
No data (non-blocking)
Buffer full (non-blocking)
▶ サンプル:ノンブロッキングチャネル + ポーリング式サーキットブレーカー
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int, 3)
go func() {
for i := 1; i <= 10; i++ {
select {
case ch <- i:
// Send successful
default:
fmt.Printf("Buffer full, dropped %d\n", i)
}
time.Sleep(10 * time.Millisecond)
}
close(ch)
}()
// Slow consumer
for v := range ch {
fmt.Printf("Processing: %d\n", v)
time.Sleep(50 * time.Millisecond)
}
}
出力:
Processing: 1
Processing: 2
Processing: 3
Buffer full, dropped 4
Buffer full, dropped 5
Processing: 6
...
default によって select が 即座に終了 してしまいます。どのチャネルも準備ができていない場合、default が実行されます。これはノンブロッキングなチャネル操作には理想的ですが、注意が必要です。forループ内でdefaultを使用すると、ビジーループが発生する可能性があります。
6. for-select ループと done チャネル
(1) 「done channel」終了モード
package main
import (
"fmt"
"time"
)
func worker(done <-chan struct{}) {
for {
select {
case <-done:
fmt.Println("worker exiting")
return
default:
fmt.Println("worker working...")
time.Sleep(200 * time.Millisecond)
}
}
}
func main() {
done := make(chan struct{})
go worker(done)
time.Sleep(1 * time.Second)
close(done)
time.Sleep(100 * time.Millisecond)
fmt.Println("main exiting")
}
▶ サンプル:for-selectループを終了する3つの方法
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int)
done := make(chan struct{})
// Producer
go func() {
for i := 1; i <= 5; i++ {
ch <- i
time.Sleep(100 * time.Millisecond)
}
close(ch)
}()
// Consumer (for-select loop)
go func() {
for {
select {
case v, ok := <-ch:
if !ok {
fmt.Println("Style 1: channel closed exit")
close(done)
return
}
fmt.Printf("Processing: %d\n", v)
case <-time.After(1 * time.Second):
fmt.Println("Style 2: timeout exit")
close(done)
return
}
}
}()
<-done
fmt.Println("main exiting")
}
(3) for-select の終了モードの比較
| モード | トリガー条件 | メリット | デメリット |
|---|---|---|---|
| close(ch) | 送信側がチャネルを閉じる | 自然な終了 | 受信側のみが使用可能 |
| 完了チャンネル | close(done) 信号 |
柔軟性があり、外部からトリガー可能 | 追加のチャンネルが必要 |
| タイムアウト | time.After タイムアウト |
デッドロックを防止する | ハードタイムアウトは柔軟性に欠ける |
| コンテキスト | ctx.Done() |
キャンセル信号を送信可能 | 第17課:詳細解説 |
7. ファンアウト/ファンインモード
(1) ファンアウト:タスクの割り当て
package main
import (
"fmt"
"sync"
)
func fanOut(jobs <-chan int, workers int) []<-chan int {
channels := make([]<-chan int, workers)
for w := 0; w < workers; w++ {
ch := make(chan int, 10)
channels[w] = ch
go func(id int, out chan<- int) {
defer close(out)
for job := range jobs {
result := job * job
fmt.Printf("Worker %d: %d^2 = %d\n", id, job, result)
out <- result
}
}(w, ch)
}
return channels
}
▶ サンプル:ファンイン:結果の集約
package main
import (
"fmt"
"sync"
)
func fanIn(channels ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, ch := range channels {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for v := range c {
out <- v
}
}(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
jobs := make(chan int, 10)
for i := 1; i <= 6; i++ {
jobs <- i
}
close(jobs)
// fan-out: distribute to 3 workers
workers := fanOut(jobs, 3)
// fan-in: aggregate all worker results
results := fanIn(workers...)
// Collect results
sum := 0
count := 0
for r := range results {
sum += r
count++
}
fmt.Printf("Total %d results, sum = %d\n", count, sum)
}
出力:
Worker 2: 1^2 = 1
Worker 0: 3^2 = 9
Worker 1: 2^2 = 4
Worker 1: 5^2 = 25
Worker 0: 4^2 = 16
Worker 2: 6^2 = 36
Total 6 results, sum = 91
(3) ファンアウトとファンイン
| モード | 方向 | 目的 |
|---|---|---|
| ファンアウト | 1チャネル → Nチャネル | タスク分散、並列計算 |
| ファンイン | Nチャネル → 1チャネル | 結果の集約、ログの収集 |
flowchart LR
subgraph FanOut [Fan-Out]
J[Jobs Channel] --> W1[Worker 1]
J --> W2[Worker 2]
J --> W3[Worker 3]
end
subgraph FanIn [Fan-In]
W1 --> R[Results Channel]
W2 --> R
W3 --> R
end
style J fill:#e1f5fe
style R fill:#fff3e0
8. 完全な例:3つのデータソースからの株価情報を集約する
// market_aggregator.go
package main
import (
"fmt"
"math/rand"
"sync"
"time"
)
// ---------- Data Types ----------
type Trade struct {
Source string
Symbol string
Price float64
Volume int
Time time.Time
}
type AggregatedQuote struct {
Symbol string
AvgPrice float64
TotalVolume int
Sources int
High float64
Low float64
Time time.Time
}
// ---------- Simulated Exchange Data Source ----------
var symbols = []string{"AAPL", "GOOGL", "MSFT", "AMZN", "TSLA", "META"}
func simulateExchange(name string, out chan<- Trade, done <-chan struct{}) {
for {
select {
case <-done:
return
default:
time.Sleep(time.Duration(100+rand.Intn(300)) * time.Millisecond)
trade := Trade{
Source: name,
Symbol: symbols[rand.Intn(len(symbols))],
Price: 100 + rand.Float64()*200,
Volume: rand.Intn(1000) + 100,
Time: time.Now(),
}
select {
case out <- trade:
case <-done:
return
}
}
}
}
// ---------- Aggregator ----------
type Aggregator struct {
trades chan Trade
done chan struct{}
wg sync.WaitGroup
}
func NewAggregator() *Aggregator {
return &Aggregator{
trades: make(chan Trade, 100),
done: make(chan struct{}),
}
}
func (a *Aggregator) AddSource(name string) {
a.wg.Add(1)
go func() {
defer a.wg.Done()
simulateExchange(name, a.trades, a.done)
}()
fmt.Printf("Added data source: %s\n", name)
}
func (a *Aggregator) Start(interval time.Duration, callback func(map[string]*AggregatedQuote)) {
// Aggregation window
window := make(map[string][]Trade)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case trade := <-a.trades:
window[trade.Symbol] = append(window[trade.Symbol], trade)
case <-ticker.C:
// Time window reached; compute aggregated quotes
quotes := make(map[string]*AggregatedQuote)
for symbol, trades := range window {
if len(trades) == 0 {
continue
}
var sumPrice float64
var totalVol int
high := trades[0].Price
low := trades[0].Price
for _, t := range trades {
sumPrice += t.Price * float64(t.Volume)
totalVol += t.Volume
if t.Price > high {
high = t.Price
}
if t.Price < low {
low = t.Price
}
}
sources := make(map[string]bool)
for _, t := range trades {
sources[t.Source] = true
}
quotes[symbol] = &AggregatedQuote{
Symbol: symbol,
AvgPrice: sumPrice / float64(totalVol),
TotalVolume: totalVol,
Sources: len(sources),
High: high,
Low: low,
Time: time.Now(),
}
}
callback(quotes)
window = make(map[string][]Trade) // Reset window
case <-a.done:
return
}
}
}
func (a *Aggregator) Stop() {
close(a.done)
a.wg.Wait()
}
// ---------- Main Function ----------
func main() {
aggregator := NewAggregator()
// Add 3 data sources
aggregator.AddSource("NYSE")
aggregator.AddSource("NASDAQ")
aggregator.AddSource("LSE")
fmt.Println("\nStarting aggregation (output every 2 seconds)...\n")
// Start aggregation with 2-second time window
done := make(chan struct{})
go func() {
aggregator.Start(2*time.Second, func(quotes map[string]*AggregatedQuote) {
fmt.Printf("=== Aggregation Report %s ===\n", time.Now().Format("15:04:05"))
for _, q := range quotes {
fmt.Printf("%-6s | $%.2f (avg) | vol=%d | src=%d | H=%.2f L=%.2f\n",
q.Symbol, q.AvgPrice, q.TotalVolume, q.Sources, q.High, q.Low)
}
fmt.Println()
})
close(done)
}()
// Run for 10 seconds then stop
time.Sleep(10 * time.Second)
aggregator.Stop()
<-done
fmt.Println("Quote aggregator stopped")
}
期待される出力:
Added data source: NYSE
Added data source: NASDAQ
Added data source: LSE
Starting aggregation (output every 2 seconds)...
=== Aggregation Report 10:00:02 ===
AAPL | $152.34 (avg) | vol=2341 | src=3 | H=165.20 L=142.10
GOOGL | $178.90 (avg) | vol=1567 | src=2 | H=185.00 L=172.30
MSFT | $295.40 (avg) | vol=890 | src=3 | H=301.20 L=288.50
=== Aggregation Report 10:00:04 ===
...
Quote aggregator stopped
select ブロックでは、case <-done: と case trade := <-a.trades: の両方にサブスクライブする必要があります。trades チャンネルのみにサブスクライブした場合、Stop シグナルは配信されません。すべての for-select ループには、終了条件を含める必要があります。
❓ よくある質問
select文をネストさせるか、default節と手動によるチェックを組み合わせて使用してください。time.After はタイムアウトをどのように処理しますか?time.After(d) は、d 単位の時間が経過した後に値を受け取るチャネルを返します。select ステートメントでは、タイムアウト分岐として case <- time.After(d): を使用してください。注:select による time.After の呼び出しごとに新しいタイマーが作成されます。ループ内では、time.NewTimer を使用してそれを再利用することをお勧めします。defaultはselectステートメント内でどのような役割を果たしますか?defaultは、selectステートメントをノンブロッキングにします。つまり、どのチャネルも準備が整っていない場合でも、defaultブロックを直ちに実行します。適した用途:ブロッキングせずに送受信を試みる場合、失敗率の算出、サーキットブレーキングやパフォーマンスの低下。注:forループ内でdefaultを使用するとビジーループが発生します。通常はtime.Sleepを追加するか、実行頻度を制限する必要があります。done := make(chan struct{})を作成し、selectブロック内でcase <-done:をリッスンし、終了したいときにclose(done)を呼び出します。doneをリッスンしているすべてのgoroutineは、同時に終了シグナルを受け取ります。closeは、値を送信するよりも好ましい方法です。これは複数回受信できるからです。select は送信に使用できますか?case ch <- value: も有効な select case です。チャネルに空き(バッファあり)がある場合、または受信側(バッファなし)がある場合に、送信準備が整います。この機能は、レート制限モードおよびセマフォモードで使用されます。select {} ではどうなりますか?select {} は無期限にブロックされます。これは、ケースもデフォルトも存在しないためです。これは「無期限に待機する」というシグナルとして機能します。ただし、本番環境のコードでは、必ず終了パスを用意する必要があります。そうしないと、goroutine がリークしてしまいます。for-select ループ内の break ステートメントで、そのループから抜け出すことはできますか?select ブロック内の break ステートメントは、select ブロックから抜け出すだけであり、for ループからは抜け出せません。for-selectループを終了させるには、ラベルの後にbreak label、return、またはdoneチャネルを指定してください。📖 まとめ
- selectは複数のチャネルを同時に監視し、最初に処理可能な状態になったチャネルを処理します
- 複数の案件が同時に処理可能になった場合は、その中からランダムに1件を選択する
time.Afterはタイムアウト制御に使用され、time.NewTimerはループ内で再利用されますdefaultは、select 操作をノンブロッキングにし、サーキットブレーカーやリトライ操作に適していますdone channelは、for-select ループを終了するための標準的な方法です- 「ファンアウト」はタスクを分散させ、「ファンイン」は結果を集約し、この2つが組み合わさって並列パイプラインを形成する
- すべての for-select には、終了条件を含める必要があります
📝 練習問題
-
基本問題(難易度 ⭐):3つのgoroutineを起動するプログラムを作成してください。各goroutineは、自身のチャンネルに整数を送信します。
selectを使用して、それらを同時に受信し、出力してください。3秒後に終了するタイムアウト分岐を追加してください。 -
上級問題(難易度 ⭐⭐):タイムアウトキャッシュを実装してください:
func FetchWithCache(key string, cache map[string]string) string。まず、キャッシュチャネルを確認します(タイムアウト 10 ms)。ヒットがない場合は、「シミュレートされたデータベース」にクエリを実行し(500 ms)、その結果をキャッシュに格納します。階層型タイムアウトを実装するには、selectとtime.Afterを組み合わせて使用する必要があります。 -
課題(難易度:⭐⭐⭐):ログ集約システムを実装する: 3つのログソース(goroutine)がそれぞれ異なるレベル(INFO/WARN/ERROR)のログを生成し、これらはファンイン操作を用いて単一のチャネルに集約された後、
selectによってレベルごとに処理されます。ERROR ログは即座にアラート(print)を発生させ、WARN ログはバッファに蓄積され(5件ずつバッチ処理でフラッシュ)、INFO ログはバッチ処理で書き込まれます (10件単位でフラッシュ)されます。要件:タイムアウトによるフラッシュおよびdoneチャネルを介した正常な終了。