Go: Goの並行処理パターン:ファンイン/ファンアウト、パイプライン、ワーカープール、Or-Done、errgroup

最終更新:2026-08-26

Goにおける並行処理は、ライブラリではなく、「ファンイン」、「ファンアウト」、「パイプライン」といったパターンに依存しています。これらのパターンを組み合わせることで、並行処理の問題の90%を解決することができます。

大量のデータを処理する必要がある場合――ジェネレータからデータを読み込み、複数のワーカーで並列に処理を行い、最後に結果を統合する――どのようにすれば洗練されたソリューションを設計できるでしょうか?このレッスンでは、Goコミュニティで最も定番の並行処理パターンを習得します。

1. 学習内容



2. あるデータエンジニアの実話

(1) 課題:1,000万件のログエントリを単一スレッドで処理するのに一晩かかってしまった

ファティマはデータプラットフォームエンジニアです。同社では毎日1,000万件のアクセスログが生成されており、これらをリアルタイムで処理する必要があります:

「最初のバージョンでは、forループを使用していました。レコードを読み込み→解析→データベースに書き込むという流れです。本番運用開始初日、処理速度が生成速度に追いつかないことが判明しました。ログキューが1時間あたり200万件のペースで増加していたのです。上司は『明日のレポートはこのデータにかかっている』と言いました。私はサーバーを追加しましたが、コードは依然としてシングルスレッドのままで、CPUの稼働率はわずか5%にとどまっていました。」

当時の彼女のコード:

GO
// Bad code: single-threaded, not utilizing CPU
func processLogs(logs []string) {
    for _, log := range logs {
        parsed := parse(log)       // CPU-intensive
        enriched := enrich(parsed) // CPU-intensive
        save(enriched)              // I/O-intensive
    }
    // 10 million log entries: 3 hours!
}

(2) Goにおける並行処理のアプローチ:パイプライン+ファンアウト

GO
// pipeline.go
package main

import (
    "fmt"
    "sync"
    "time"
)

type LogEntry struct {
    Raw     string
    Level   string
    Message string
    Time    time.Time
}

func main() {
    logs := make([]string, 100)
    for i := range logs {
        logs[i] = fmt.Sprintf("log-%d", i)
    }
    start := time.Now()

    // Pipeline: generator → parse → enrich → save
    rawCh := generator(logs)
    parsedCh := parseStage(rawCh, 4)       // 4 parse workers
    enrichedCh := enrichStage(parsedCh, 4)  // 4 enrich workers
    saveStage(enrichedCh, 2)                // 2 save workers

    fmt.Printf("Processing complete, elapsed: %v\n", time.Since(start))
}

// Stage 1: Generator (data source)
func generator(logs []string) <-chan string {
    out := make(chan string, 100)
    go func() {
        defer close(out)
        for _, log := range logs {
            out <- log
        }
    }()
    return out
}

// Stage 2: Parse (Fan-out)
func parseStage(in <-chan string, workers int) <-chan LogEntry {
    out := make(chan LogEntry, 100)
    for i := 0; i < workers; i++ {
        go func() {
            for raw := range in {
                time.Sleep(1 * time.Millisecond) // Simulate parsing
                out <- LogEntry{Raw: raw, Level: "INFO", Message: raw, Time: time.Now()}
            }
        }()
    }
    return out
}

// Stage 3: Enrich (Fan-out)
func enrichStage(in <-chan LogEntry, workers int) <-chan LogEntry {
    out := make(chan LogEntry, 100)
    for i := 0; i < workers; i++ {
        go func() {
            for entry := range in {
                time.Sleep(2 * time.Millisecond) // Simulate data enrichment
                out <- entry
            }
        }()
    }
    return out
}

// Stage 4: Save (Fan-in to 2 workers)
func saveStage(in <-chan LogEntry, workers int) {
    var wg sync.WaitGroup
    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for entry := range in {
                _ = entry // Simulate writing to database
                time.Sleep(500 * time.Microsecond)
            }
        }()
    }
    wg.Wait()
}

(3) パフォーマンス:シングルスレッドとパイプラインの比較

メトリック シングルスレッド パイプライン (4+4+2 ワーカー) 改善率
1,000万件のレコードの処理時間 約3時間 約12分 15倍
CPU使用率 5% 85% 17倍
コードの複雑度 単純 中程度
スケーラビリティ サーバーを追加 ワーカーを追加(数を変更)


3. パイプライン

▶ サンプル:数値処理パイプライン

GO
package main

import (
    "fmt"
)

// Stage 1: Generate
func generate(nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            out <- n
        }
    }()
    return out
}

// Stage 2: Square
func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            out <- n * n
        }
    }()
    return out
}

// Stage 3: Double
func double(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            out <- n * 2
        }
    }()
    return out
}

func main() {
    // Pipeline: generate → square → double
    for result := range double(square(generate(1, 2, 3, 4, 5))) {
        fmt.Println(result)
    }
    // Output: 2, 8, 18, 32, 50  (n²×2)
}
▶ 試してみよう
100%
graph LR
    A[generate] -->|chan int| B[square]
    B -->|chan int| C[double]
    C -->|chan int| D[main]
💡 ヒント: パイプラインの鍵となるのは、各ステージが <-chan T(読み取り専用チャネル)を返し、<-chan T を入力として受け取るという点です。各ステージは独自の goroutine で実行され、チャネルを介して接続されています。これが Go の CSP モデルです。



4. ファンアウト/ファンイン

▶ サンプル:ファンアウト+ファンイン

GO 📖 参照専用
package main

import (
    "fmt"
    "sync"
)

// Stage 1: Generate
func generate(nums ...int) <-chan int {
    out := make(chan int, 10)
    go func() {
        defer close(out)
        for _, n := range nums {
            out <- n
        }
    }()
    return out
}

// Stage 2: Worker (Fan-out target)
func worker(id int, in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            result := n * n
            fmt.Printf("Worker %d: %d^2 = %d\n", id, n, result)
            out <- result
        }
    }()
    return out
}

// Fan-in: merge multiple channels
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() {
    // Fan-out: one generator distributes to 3 workers
    in := generate(1, 2, 3, 4, 5, 6)
    w1 := worker(1, in)
    w2 := worker(2, in)
    w3 := worker(3, in)

    // Fan-in: merge results from 3 workers into one channel
    results := fanIn(w1, w2, w3)

    for result := range results {
        fmt.Printf("Result: %d\n", result)
    }
}
論理コード 55 行(40 行制限超過、参照専用)
100%
graph LR
    G[Generate] -->|fan-out| W1[Worker 1]
    G -->|fan-out| W2[Worker 2]
    G -->|fan-out| W3[Worker 3]
    W1 -->|fan-in| R[Results]
    W2 -->|fan-in| R
    W3 -->|fan-in| R


5. ワーカープール

▶ サンプル:ワーカープール

GO 📖 参照専用
package main

import (
    "fmt"
    "sync"
    "time"
)

type Job struct {
    ID      int
    Payload string
}

type Result struct {
    JobID  int
    Output string
    Err    error
}

func worker(id int, jobs <-chan Job, results chan<- Result, wg *sync.WaitGroup) {
    defer wg.Done()
    for job := range jobs {
        fmt.Printf("Worker %d processing job %d\n", id, job.ID)
        time.Sleep(100 * time.Millisecond) // Simulate work
        results <- Result{
            JobID:  job.ID,
            Output: fmt.Sprintf("processed by worker %d", id),
        }
    }
}

func main() {
    const numJobs = 20
    const numWorkers = 5

    jobs := make(chan Job, numJobs)
    results := make(chan Result, numJobs)

    // Start Worker Pool
    var wg sync.WaitGroup
    for i := 0; i < numWorkers; i++ {
        wg.Add(1)
        go worker(i, jobs, results, &wg)
    }

    // Send jobs
    for i := 0; i < numJobs; i++ {
        jobs <- Job{ID: i, Payload: fmt.Sprintf("data-%d", i)}
    }
    close(jobs)

    // Wait for all workers to complete
    wg.Wait()
    close(results)

    // Collect results
    for result := range results {
        fmt.Printf("Job %d → %s\n", result.JobID, result.Output)
    }
}
論理コード 46 行(40 行制限超過、参照専用)

(2) パイプライン vs ワーカープール vs ファンアウト/ファンイン

パターン 基本概念 適用可能なシナリオ
パイプライン データは複数のステージを経て流れ、各ステージでは1つの処理が行われる 処理手順が明確に定義されている(解析 → 変換 → 保存)
ワーカープール 固定数のワーカーがジョブキューからタスクを取得する タスクのバッチ処理(画像圧縮、メール送信)
ファンアウト/ファンイン 並列処理のために複数のワーカーに分散させ、その後結果を統合する ステートレスな並列計算(数値演算、データフィルタリング)


6. Or-Doneチャンネル

▶ サンプル:Or-Doneパターン

GO 📖 参照専用
package main

import (
    "fmt"
    "time"
)

// orDone wraps a channel, making it support cancellation via a done channel
func orDone(done <-chan struct{}, c <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for {
            select {
            case <-done:
                return
            case v, ok := <-c:
                if !ok {
                    return
                }
                select {
                case out <- v:
                case <-done:
                    return
                }
            }
        }
    }()
    return out
}

func generateWithCancel(done <-chan struct{}, nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            select {
            case <-done:
                return
            case out <- n:
            }
        }
    }()
    return out
}

func main() {
    done := make(chan struct{})

    // Start generator (will be canceled after reaching 5)
    nums := generateWithCancel(done, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10)

    // Consume safely via orDone
    for n := range orDone(done, nums) {
        fmt.Printf("Processing: %d\n", n)
        if n >= 5 {
            close(done) // Cancel all operations
        }
    }

    fmt.Println("Canceled")
    time.Sleep(100 * time.Millisecond) // Wait for goroutines to exit
}
論理コード 53 行(40 行制限超過、参照専用)
🔥 よくある間違い: チャネルからの受信を怠ると、ゴルーチンリークが発生する可能性があります。つまり、ジェネレータ・ゴルーチンが送信処理でブロックされてしまうのです。Or-Done パターンを使用すれば、キャンセル時にゴルーチンがリークされることはありません。 アップストリームの送信側とダウンストリームの受信側の双方がキャンセルを認識する必要がある場合は、チャネルを Or-Done でラップしてください。



7. errgroup エラーの伝播

⚙️ 前提条件:このパッケージを使用する前に、go get golang.org/x/sync/errgroup を実行してください。

▶ サンプル:errgroup

GO
package main

import (
    "fmt"
    "time"

    "golang.org/x/sync/errgroup"
)

func main() {
    // Create an errgroup
    // The first goroutine to return an error will cancel the others
    g := errgroup.Group{}

    // Start 3 goroutines
    for i := 0; i < 3; i++ {
        id := i
        g.Go(func() error {
            return doWork(id)
        })
    }

    // Wait for all goroutines to complete, get the first error
    if err := g.Wait(); err != nil {
        fmt.Printf("Task failed: %v\n", err)
    } else {
        fmt.Println("All succeeded")
    }
}

func doWork(id int) error {
    fmt.Printf("Worker %d starting\n", id)
    time.Sleep(time.Duration(id+1) * 500 * time.Millisecond)

    if id == 1 {
        return fmt.Errorf("worker %d failed", id)
    }

    fmt.Printf("Worker %d complete\n", id)
    return nil
}
▶ 試してみよう

▶ サンプル:errgroup + コンテキストタイムアウト

GO
package main

import (
    "context"
    "fmt"
    "time"

    "golang.org/x/sync/errgroup"
)

func main() {
    // errgroup with Context
    ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
    defer cancel()

    g, ctx := errgroup.WithContext(ctx)

    // Workers check if Context is canceled
    for i := 0; i < 5; i++ {
        id := i
        g.Go(func() error {
            select {
            case <-time.After(time.Duration(id+1) * time.Second):
                fmt.Printf("Worker %d complete\n", id)
                return nil
            case <-ctx.Done():
                fmt.Printf("Worker %d canceled: %v\n", id, ctx.Err())
                return ctx.Err()
            }
        })
    }

    if err := g.Wait(); err != nil {
        fmt.Printf("errgroup stopped: %v\n", err)
    }
}
▶ 試してみよう

(3) sync.WaitGroup と errgroup の比較

プロパティ sync.WaitGroup errgroup
エラーの収集 未対応 ✅ 最初のエラーを返す
伝播のキャンセル 未対応 ✅ WithContext による自動キャンセル
ゴルーチン管理 手動での追加・完了 ✅ Goのメソッドによる自動処理
タイムアウト制御 手動での実装 ✅ コンテキストを使用して簡単に実装可能


8. 完全な例:リアルタイムデータ処理パイプライン

GO
// data_pipeline.go
package main

import (
    "context"
    "fmt"
    "math/rand"
    "sync"
    "time"
)

// ---------- Data types ----------

type DataPoint struct {
    ID        int
    Value     float64
    Timestamp time.Time
}

type ProcessedData struct {
    DataPoint
    Normalized bool
    Anomaly    bool
    Category   string
}

// ---------- Pipeline Stages ----------

// Stage 1: Source (data source)
func source(ctx context.Context, count int) <-chan DataPoint {
    out := make(chan DataPoint, 100)
    go func() {
        defer close(out)
        for i := 0; i < count; i++ {
            select {
            case <-ctx.Done():
                return
            case out <- DataPoint{
                ID:        i,
                Value:     rand.Float64() * 100,
                Timestamp: time.Now(),
            }:
            }
        }
    }()
    return out
}

// Stage 2: Normalize (Fan-out)
func normalize(ctx context.Context, in <-chan DataPoint, workers int) <-chan ProcessedData {
    out := make(chan ProcessedData, 100)
    var wg sync.WaitGroup

    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for dp := range orDoneDataPoint(ctx.Done(), in) {
                select {
                case <-ctx.Done():
                    return
                case out <- ProcessedData{
                    DataPoint:  dp,
                    Normalized: true,
                    Category:   categorize(dp.Value),
                }:
                }
            }
        }()
    }

    go func() {
        wg.Wait()
        close(out)
    }()

    return out
}

// Stage 3: Detect Anomaly (Fan-out)
func detectAnomaly(ctx context.Context, in <-chan ProcessedData, workers int) <-chan ProcessedData {
    out := make(chan ProcessedData, 100)
    var wg sync.WaitGroup

    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for pd := range orDoneProcessedData(ctx.Done(), in) {
                pd.Anomaly = pd.Value > 90 || pd.Value < 10
                select {
                case <-ctx.Done():
                    return
                case out <- pd:
                }
            }
        }()
    }

    go func() {
        wg.Wait()
        close(out)
    }()

    return out
}

// Stage 4: Sink (output, Fan-in)
func sink(ctx context.Context, in <-chan ProcessedData) {
    var normalCount, anomalyCount int
    for pd := range orDoneProcessedData(ctx.Done(), in) {
        if pd.Anomaly {
            anomalyCount++
            fmt.Printf("[ANOMALY] ID=%d, Value=%.2f, Category=%s\n",
                pd.ID, pd.Value, pd.Category)
        } else {
            normalCount++
        }
    }
    fmt.Printf("\nSummary: normal=%d, anomaly=%d\n", normalCount, anomalyCount)
}

// ---------- Helper functions ----------

func orDoneDataPoint(done <-chan struct{}, in <-chan DataPoint) <-chan DataPoint {
    out := make(chan DataPoint)
    go func() {
        defer close(out)
        for {
            select {
            case <-done:
                return
            case v, ok := <-in:
                if !ok {
                    return
                }
                select {
                case out <- v:
                case <-done:
                    return
                }
            }
        }
    }()
    return out
}

func orDoneProcessedData(done <-chan struct{}, in <-chan ProcessedData) <-chan ProcessedData {
    out := make(chan ProcessedData)
    go func() {
        defer close(out)
        for {
            select {
            case <-done:
                return
            case v, ok := <-in:
                if !ok {
                    return
                }
                select {
                case out <- v:
                case <-done:
                    return
                }
            }
        }
    }()
    return out
}

func categorize(value float64) string {
    switch {
    case value < 30:
        return "low"
    case value < 70:
        return "medium"
    default:
        return "high"
    }
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
    defer cancel()

    fmt.Println("Starting real-time data processing pipeline...")

    // Build Pipeline
    data := source(ctx, 1000)
    normalized := normalize(ctx, data, 5)           // 5 normalize workers
    detected := detectAnomaly(ctx, normalized, 3)   // 3 anomaly detection workers
    sink(ctx, detected)

    fmt.Println("Pipeline processing complete")
}
💡 ヒント: 本番環境では、パイプラインの各ステージごとに、独自のバッファサイズ、再試行ロジック、および監視メトリクスを設定する必要があります。expvar または prometheus を使用すると、各ステージの処理速度、キューの長さ、エラー率を確認できます。


❓ よくある質問

Q パイプラインとワーカープールの違いは何ですか?
A パイプラインとは、複数の処理段階を経るデータフローのことで、各段階では単一のステップ(例:解析 → 変換 → 保存)が実行され、データの形式が変化する場合があります。ワーカープールは、同じジョブキューからタスクを取得する固定数のワーカーで構成され、データ形式は変更されません(例:画像圧縮)。
Q ファンアウトとファンインの背後にある基本的な概念は何ですか?
A ファンアウトとは、単一のチャネルからのデータを複数のワーカーに分散させ、並列処理を行うことです(スループットを向上させるため)。ファンインとは、複数のチャネルからの結果を単一のチャネルに統合することです(一元的な利用のため)。これら2つは通常、組み合わせて使用されます。
Q いつ「Or-Done」を使用すべきですか?
A ゴルーチンライフサイクルを制御できないチャネル(サードパーティ製ライブラリが返すチャネルなど)から読み込む場合です。「Or-Done」を使用することで、呼び出し側が処理をキャンセルした場合でも、読み取りを行うゴルーチンがチャネル操作によってブロックされることを防ぎます。単純なシナリオでは、selectdone チャネルと組み合わせて使用できます。
Q errgroupsync.WaitGroup のどちらを選べばよいですか?
A エラーを収集したり、キャンセルを伝播させたりする必要がある場合は errgroup を使用してください。単にゴルーチンが終了するのを待つだけでよい場合は、より軽量な sync.WaitGroup を使用してください。errgroupは機能面ではWaitGroupのスーパーセットですが、キャンセルを管理するために追加のgoroutineを導入します。
Q パイプライン内のチャネルのバッファサイズはどのように決定すればよいですか?
A 基本的な原則は、プロデューサーがコンシューマーによってブロックされないようにすることです。バッファが大きすぎるとメモリを浪費し、小さすぎるとゴルーチン間のコンテキストスイッチが増加します。経験則として、バッファサイズ = 生成レート × 目標レイテンシとなります。実際のサイズは、本番環境での負荷テストを通じて決定してください。
Q ゴルーチンリークを防ぐにはどうすればよいですか?
A 3つのルールがあります:(1) ゴルーチンを作成する場所には、必ず終了メカニズム(チャネルの閉じたり、doneシグナルの送信など)を設けること; (2) ライフサイクルを管理するには errgroup または WaitGroup を使用すること; (3) 制御されていないチャネルをラップするには Or-Done を使用すること。終了時期が分からない場合は、決してゴルーチンを作成してはいけません。
Q これらのパターンは組み合わせて使用できますか?
A はい、組み合わせるべきです。たとえば、パイプラインの各ステージ内で、ワーカープールとファンアウト/ファンインを組み合わせて使用することができます。また、errgroupはパイプラインのオーケストレーターとして機能し、キャンセルやエラーを一元的に管理することができます。

📖 まとめ


📝 練習問題

  1. 基本問題(難易度 ⭐):3段階の数値処理パイプラインを実装してください:生成(1~100)→ フィルタリング(偶数のみを残す)→ 合計(累積)。各段階は別々のgoroutineで実行され、チャネルを介して接続されています。

  2. 上級問題(難易度 ⭐⭐)並行ログ処理システムを実装してください。構成は、10個のパースワーカー(ファンアウト)+3個の書き込みワーカー(ファンイン)+errgroupによるエラー処理です。1,000件のログエントリをシミュレートし、統計情報(総経過時間、処理されたエントリ数、エラー数)を出力してください。

  3. 課題(難易度:⭐⭐⭐)負荷分散と動的スケーリングを備えたワーカープールを設計・実装してください。要件:(1) 最初はワーカーを5つで開始すること;(2) ジョブキューのバックログが閾値を超えた場合、ワーカーの数を自動的に増やすこと(最大20まで); (3) キューが空になった際に、ワーカーの数を(最大2まで)減らすこと;(4) Context を使用してワーカーの終了を制御すること;(5) -race を使用して、レースコンディションが発生していないことを検証すること。

Web-Tutorial.com

Web-Tutorial 技術チーム

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

100%