Go: Goの並行処理パターン:ファンイン/ファンアウト、パイプライン、ワーカープール、Or-Done、errgroup
最終更新:2026-08-26
Goにおける並行処理は、ライブラリではなく、「ファンイン」、「ファンアウト」、「パイプライン」といったパターンに依存しています。これらのパターンを組み合わせることで、並行処理の問題の90%を解決することができます。
大量のデータを処理する必要がある場合――ジェネレータからデータを読み込み、複数のワーカーで並列に処理を行い、最後に結果を統合する――どのようにすれば洗練されたソリューションを設計できるでしょうか?このレッスンでは、Goコミュニティで最も定番の並行処理パターンを習得します。
1. 学習内容
- パイプライン・パターン
- ファンアウト/ファンイン
- 労働者プール
- Or-Doneチャンネルの簡易モード
- キャンセルの伝播
golang.org/x/sync/errgroupエラーグループ- 実践編:リアルタイムデータ処理パイプライン
2. あるデータエンジニアの実話
(1) 課題:1,000万件のログエントリを単一スレッドで処理するのに一晩かかってしまった
ファティマはデータプラットフォームエンジニアです。同社では毎日1,000万件のアクセスログが生成されており、これらをリアルタイムで処理する必要があります:
「最初のバージョンでは、forループを使用していました。レコードを読み込み→解析→データベースに書き込むという流れです。本番運用開始初日、処理速度が生成速度に追いつかないことが判明しました。ログキューが1時間あたり200万件のペースで増加していたのです。上司は『明日のレポートはこのデータにかかっている』と言いました。私はサーバーを追加しましたが、コードは依然としてシングルスレッドのままで、CPUの稼働率はわずか5%にとどまっていました。」
当時の彼女のコード:
// 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における並行処理のアプローチ:パイプライン+ファンアウト
// 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. パイプライン
▶ サンプル:数値処理パイプライン
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)
}
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. ファンアウト/ファンイン
▶ サンプル:ファンアウト+ファンイン
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)
}
}
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. ワーカープール
▶ サンプル:ワーカープール
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)
}
}
(2) パイプライン vs ワーカープール vs ファンアウト/ファンイン
| パターン | 基本概念 | 適用可能なシナリオ |
|---|---|---|
| パイプライン | データは複数のステージを経て流れ、各ステージでは1つの処理が行われる | 処理手順が明確に定義されている(解析 → 変換 → 保存) |
| ワーカープール | 固定数のワーカーがジョブキューからタスクを取得する | タスクのバッチ処理(画像圧縮、メール送信) |
| ファンアウト/ファンイン | 並列処理のために複数のワーカーに分散させ、その後結果を統合する | ステートレスな並列計算(数値演算、データフィルタリング) |
6. Or-Doneチャンネル
▶ サンプル:Or-Doneパターン
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
}
7. errgroup エラーの伝播
⚙️ 前提条件:このパッケージを使用する前に、
go get golang.org/x/sync/errgroupを実行してください。
▶ サンプル:errgroup
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 + コンテキストタイムアウト
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. 完全な例:リアルタイムデータ処理パイプライン
// 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 を使用すると、各ステージの処理速度、キューの長さ、エラー率を確認できます。
❓ よくある質問
select を done チャネルと組み合わせて使用できます。errgroup と sync.WaitGroup のどちらを選べばよいですか?errgroup を使用してください。単にゴルーチンが終了するのを待つだけでよい場合は、より軽量な sync.WaitGroup を使用してください。errgroupは機能面ではWaitGroupのスーパーセットですが、キャンセルを管理するために追加のgoroutineを導入します。doneシグナルの送信など)を設けること; (2) ライフサイクルを管理するには errgroup または WaitGroup を使用すること; (3) 制御されていないチャネルをラップするには Or-Done を使用すること。終了時期が分からない場合は、決してゴルーチンを作成してはいけません。📖 まとめ
- パイプライン:データは複数の処理段階を経て流れ、各段階ごとに1つのgoroutineが割り当てられる
- ファンアウト:1つのチャネルを複数のワーカーに分散させる(並列処理のため)
- ファンイン:複数のチャンネルを1つのチャンネルに統合すること(一元的な視聴のため)
- ワーカープール:一定数のワーカーがジョブキューからタスクを取得する
- Or-Done:キャンセル機能をサポートし、制御不能なチャネルを安全に処理する
- errgroup: ゴルーチン・グループを管理し、エラーの収集とキャンセル伝播をサポートする
- 組み合わせモード:Pipelineは内部でワーカープールとファンアウト/ファンインを併用しています
📝 練習問題
-
基本問題(難易度 ⭐):3段階の数値処理パイプラインを実装してください:生成(1~100)→ フィルタリング(偶数のみを残す)→ 合計(累積)。各段階は別々のgoroutineで実行され、チャネルを介して接続されています。
-
上級問題(難易度 ⭐⭐):並行ログ処理システムを実装してください。構成は、10個のパースワーカー(ファンアウト)+3個の書き込みワーカー(ファンイン)+errgroupによるエラー処理です。1,000件のログエントリをシミュレートし、統計情報(総経過時間、処理されたエントリ数、エラー数)を出力してください。
-
課題(難易度:⭐⭐⭐):負荷分散と動的スケーリングを備えたワーカープールを設計・実装してください。要件:(1) 最初はワーカーを5つで開始すること;(2) ジョブキューのバックログが閾値を超えた場合、ワーカーの数を自動的に増やすこと(最大20まで); (3) キューが空になった際に、ワーカーの数を(最大2まで)減らすこと;(4)
Contextを使用してワーカーの終了を制御すること;(5)-raceを使用して、レースコンディションが発生していないことを検証すること。