Go: GoのSelectおよび並行処理パターン:多重化、タイムアウト制御、ファンイン/ファンアウト、パイプライン

select は、Go言語における並行プログラミングの「スイスアーミーナイフ」のような存在です。複数のチャネルを同時に監視し、最初に処理可能な状態になったチャネルを処理します。これは、ゴルーチンをつなぐ「スマートスイッチ」としての役割を果たします。

チャネルをゴルーチンをつなぐ電話回線だと考えると、selectは電話交換手のようなものです。すべての着信を同時に監視し、鳴った電話に応答します。このレッスンでは、selectの主要なパターンをすべて習得します。

1. 学習内容



2. クオンツ・トレーディング・エンジニアの実話

(1) 課題:3つのデータソースからデータを取得すると、CPUがフル稼働してしまう

ボブは、クオンツ・トレーディングチームのバックエンドエンジニアです。彼は、3つの株価データソースを同時に監視する必要があります:

「当社の戦略では、NYSE、NASDAQ、LSEの3つの取引所から同時にリアルタイムの株価情報を取得する必要があります。以前のJava実装では、1つのスレッドで3つのWebSocketをポーリングしており、そのたびに10ミリ秒のビジーウェイトが発生し、CPUの30%を消費していました。上司からは、『取引システム本体よりも多くのリソースを使っている』と言われたほどです。」

彼は現在のコードを開いた:

GO
// 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つのチャンネルを聞く

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

出力:

TEXT 📖 参照専用
[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文の構文

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

▶ サンプル:選択(ランダム選択)

GO
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)
        }
    }
}
▶ 試してみよう

出力(毎回異なる場合があります):

TEXT 📖 参照専用
from ch1
from ch2


4. タイムアウト制御

(1) タイムアウトモード

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

出力:

TEXT 📖 参照専用
Timeout! Operation exceeded 1 second

▶ サンプル:段階的なタイムアウト(段階的な待機)

GO 📖 参照専用
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!")
            }
        }
    }
}
論理コード 41 行(40 行制限超過、参照専用)
💡 ヒント: time.After(d)<-chan time.Time を返し、d 秒後に現在の時刻を送信します。selecttime.After を呼び出すたびに、新しいタイマーが作成されます。ループ内で頻繁に呼び出す場合は、リソースリークを防ぐために、必ず time.NewTimer を使用してタイマーを停止してください。



5. デフォルト:ノンブロッキング操作

(1) ノンブロッキング送受信

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

出力:

TEXT 📖 参照専用
No data (non-blocking)
Buffer full (non-blocking)

▶ サンプル:ノンブロッキングチャネル + ポーリング式サーキットブレーカー

GO
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)
    }
}
▶ 試してみよう

出力:

TEXT 📖 参照専用
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」終了モード

GO
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つの方法

GO
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) ファンアウト:タスクの割り当て

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

▶ サンプル:ファンイン:結果の集約

GO
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)
}
▶ 試してみよう

出力:

TEXT 📖 参照専用
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チャネル 結果の集約、ログの収集
100%
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つのデータソースからの株価情報を集約する

GO
// 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")
}

期待される出力:

TEXT 📖 参照専用
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 ループには、終了条件を含める必要があります。


❓ よくある質問

Q 複数のチャネルが同時に実行可能になった場合はどうすればよいですか?
A Goは、実行するケースを擬似ランダムに1つ選択します。これは、開発者がチャネルの優先順位に依存することを防ぐため、言語仕様の一部となっています。優先順位が必要な場合は、2つのselect文をネストさせるか、default節と手動によるチェックを組み合わせて使用してください。
Q time.After はタイムアウトをどのように処理しますか?
A time.After(d) は、d 単位の時間が経過した後に値を受け取るチャネルを返します。select ステートメントでは、タイムアウト分岐として case <- time.After(d): を使用してください。注:select による time.After の呼び出しごとに新しいタイマーが作成されます。ループ内では、time.NewTimer を使用してそれを再利用することをお勧めします。
Q defaultselectステートメント内でどのような役割を果たしますか?
A defaultは、selectステートメントをノンブロッキングにします。つまり、どのチャネルも準備が整っていない場合でも、defaultブロックを直ちに実行します。適した用途:ブロッキングせずに送受信を試みる場合、失敗率の算出、サーキットブレーキングやパフォーマンスの低下。注:forループ内でdefaultを使用するとビジーループが発生します。通常はtime.Sleepを追加するか、実行頻度を制限する必要があります。
Q doneチャネルを使ってfor-selectループを終了するにはどうすればよいですか?
A done := make(chan struct{})を作成し、selectブロック内でcase <-done:をリッスンし、終了したいときにclose(done)を呼び出します。doneをリッスンしているすべてのgoroutineは、同時に終了シグナルを受け取ります。closeは、値を送信するよりも好ましい方法です。これは複数回受信できるからです。
Q ファンインとファンアウトとは何ですか?
A ファンアウトとは、1つの入力チャネルをN個のワーカーgoroutineに分散させること(並列処理)です。ファンインとは、N個のチャネルを1つの出力チャネルに集約すること(結果の統合)です。この2つの組み合わせは、Go言語における並列計算の典型的なパターンです。
Q select は送信に使用できますか?
A はい。case ch <- value: も有効な select case です。チャネルに空き(バッファあり)がある場合、または受信側(バッファなし)がある場合に、送信準備が整います。この機能は、レート制限モードおよびセマフォモードで使用されます。
Q 空の select {} ではどうなりますか?
A select {} は無期限にブロックされます。これは、ケースもデフォルトも存在しないためです。これは「無期限に待機する」というシグナルとして機能します。ただし、本番環境のコードでは、必ず終了パスを用意する必要があります。そうしないと、goroutine がリークしてしまいます。
Q for-select ループ内の break ステートメントで、そのループから抜け出すことはできますか?
A いいえ。select ブロック内の break ステートメントは、select ブロックから抜け出すだけであり、for ループからは抜け出せません。for-selectループを終了させるには、ラベルの後にbreak labelreturn、またはdoneチャネルを指定してください。

📖 まとめ


📝 練習問題

  1. 基本問題(難易度 ⭐):3つのgoroutineを起動するプログラムを作成してください。各goroutineは、自身のチャンネルに整数を送信します。selectを使用して、それらを同時に受信し、出力してください。3秒後に終了するタイムアウト分岐を追加してください。

  2. 上級問題(難易度 ⭐⭐)タイムアウトキャッシュを実装してください:func FetchWithCache(key string, cache map[string]string) string。まず、キャッシュチャネルを確認します(タイムアウト 10 ms)。ヒットがない場合は、「シミュレートされたデータベース」にクエリを実行し(500 ms)、その結果をキャッシュに格納します。階層型タイムアウトを実装するには、selecttime.Afterを組み合わせて使用する必要があります。

  3. 課題(難易度:⭐⭐⭐)ログ集約システムを実装する: 3つのログソース(goroutine)がそれぞれ異なるレベル(INFO/WARN/ERROR)のログを生成し、これらはファンイン操作を用いて単一のチャネルに集約された後、select によってレベルごとに処理されます。ERROR ログは即座にアラート(print)を発生させ、WARN ログはバッファに蓄積され(5件ずつバッチ処理でフラッシュ)、INFO ログはバッチ処理で書き込まれます (10件単位でフラッシュ)されます。要件:タイムアウトによるフラッシュおよびdoneチャネルを介した正常な終了。

Web-Tutorial.com

Web-Tutorial 技術チーム

複数の開発者によって共同維持されているプログラミングチュートリアルプラットフォーム。各チュートリアルは専門分野の開発者が執筆・レビューしています。正確で信頼性の高いコンテンツを目指しています — 問題を見つけた場合はお知らせください。

100%