Goにおける並行性の制御:errgroup、ワーカープールの活用とバックプレッシャー

単にキャンセルするだけでなく、進行中の作業の上限を設けてください。

目次

Goにおける非同期処理のデフォルトの失敗モードは、無制限のファンアウト(Forking)です。goステートメントの1つごとに処理中のタスクが1つ増え、上限も管理する主体もありません。この記事では、その両方(上限と管理主体)を与える方法について説明します。

キャンセル(取り消し)とバウンディング(限界の設定)は、別々の制御メカニズムです。キャンセルは作業がいつ停止するかを決定し、バウンディングは一度にどのくらいの作業が存在するかを決定します。サービスが完璧にキャンセルできても、1万個のゴルーチンがそれぞれ独立して「今がデータベースを呼び出すのに最適な時だ」と判断して倒れる、という状況は起こり得ます。

光る粒子の流れを規制されたゲートを通って絞っているガラス製のファネル(漏斗)

このページ周辺の記事は、問題のもう半面を扱っています。Goのcontext.Contextを正しく使う は、作業を停止するための制御プレーンとしてcontextを扱いますが、ここでは対象は容量プレーンです——errgroup.SetLimit、セマフォチャネル、ワーカープール、有界キュー——そして、すべてのキャンセルパスが適切に接続されているように見えても、errgroup.Waitの後でも生き残る1つのリーク形状についても説明します。

5つの原語と1つの意思決定テーブル

原語(Primitive) 部分結果 失敗即時終了(Fail-fast) バックプレッシャー 動的な作業 最も適しているケース
sync.WaitGroup + errors.Join はい いいえ いいえ いいえ すべての結果が重要な、全待機バッチ
errgroup.WithContext いいえ はい いいえ いいえ 最初のエラーでリクエストを終了させるファンアウト
errgroup.SetLimit(n) いいえ はい はい(Goをブロック) いいえ 限界付きの失敗即時終了ファンアウト
セマフォチャネル はい 手動 はい(取得をブロック) はい カスタムな入場制御
ワーカープール はい 手動 はい(キューが満杯) はい 安定したストリーム、長寿命ワーカー

2つの質問が適切な行(選択肢)を選び出します。第一に:1つのタスクが失敗した場合、他のタスクの結果は欲しいですか?はいの場合、失敗即時終了(fail-fast)なグループのセマンティクスは、必要だった作業を破棄してしまいます。第二に:作業は連続して到着しますか、それとも開始前に数えられる固定バッチですか?グループはバッチを、プールとキューはストリームを処理します。

errors.Joinこそが、素の WaitGroup の行を有効なものにするものです。これにより、最初のエラーだけを残すのではなく、すべてのワーカーのエラーを収集できるようになり、Goのエラー処理アーキテクチャ の境界変換ルールと相性があります。

errgroup.Waitを生き延びるリーク

最もよく知られているerrgroupの失敗モードは、バウンディングとは関係ありません。これはキャンセルギャップであり、正しく見えるコードの後ろに潜んでいます。

g, gctx := errgroup.WithContext(ctx)
results := make(chan int) // バッファなし

go func() { // プロデューサー — このゴルーチンの所有権は誰も持っていない
    for i := 0; i < 100; i++ {
        results <- i
    }
    close(results)
}()

g.Go(func() error { // コンシューマー — グループメンバー
    for r := range results {
        if r == 5 {
            return fmt.Errorf("save failed")
        }
    }
    return nil
})

err := g.Wait() // コンシューマーのエラーで返ってくる

コンシューマーが戻ると、gctx はキャンセルされ、Wait は返ります。プロデューサーは results <- i で停止(park)しています——そして、コンテキストのキャンセルはチャネルの送信をアンブロックしません。そのゴルーチンの所有者がいないため、プロセスが存続する限りずっと停止したままになります。この形をループで実行すると、リークはちょうど線形になります:

変種 500回の反復後のライブなゴルーチン数
ctx.Done アームなしで送信 501 (+500リーク、反復ごとに1つ)
<-gctx.Done() を含む select で送信をラップ 502 (+1ベースライン)
// 修正法: キャンセル可能なタスク内のすべてのチャネル操作に
// Doneアームをつける
select {
case results <- i:
case <-gctx.Done():
    return
}

同じギャップは受信側にも存在します。ルールは機械的です:キャンセルされ得るコード内のすべてのチャネル操作に <-ctx.Done() アームをつけ、すべてのゴルーチンには、それを待つかキャンセルする所有者を持たせます。CIでの goleak は、レビューが見逃すケースを捕捉します。これは、Go Linters: コード品質のための必須ツール のツールと同じ、自動化された品質ゲートにスロットします。

関連する2つの形状について、名前を付けておく価値があります。プロデューサーがグループメンバーである場合、同じコードはリークではなくデッドロック(deadlock)を生成します——Wait は停止した送りを待って永遠にブロックしますが、少なくともこれは音を立てて失敗します。そして、select で送信するが、コンシューマーが終了したチャネルから受信するプロデューサーは同じようにリークします。Done アームは両側(送信側と受信側)に必要です。

errgroup.SetLimit: 限界付きの失敗即時終了ファンアウト

SetLimit は、errgroupを入場制御されたグループに変えます。限界を超える各 Go 呼び出しは、スロットが解放されるまでブロックします:

g, ctx := errgroup.WithContext(ctx)
g.SetLimit(8)
for _, item := range items {
    g.Go(func() error {
        return process(ctx, item)
    })
}
err := g.Wait()

本番環境では3つのセマンティクスが重要です。Go はブロックするため、上記のループはプロデューサーにバックプレッシャーを適用します——それは通常望ましいことですが、呼び出し元のゴルーチンがスタック(停止)する可能性があることを意味します。したがって、大規模なバッチを投入するリクエストハンドラは、ループ全体周囲に独自のタイムアウト予算が必要です。SetLimit(0) はすべての Go 呼び出しを永遠にブロックします。そして、ゴルーチンがアクティブな間に限界を変更するとパニックします——限界はグループの存続期間中に固定されます。

SetLimit はグループのセマンティクスを継承します:最初のエラーがコンテキストをキャンセルし、Wait は最初のエラーを返して、進行中の作業の結果を破棄します。最初の失敗で操作全体を失敗させるのが正しい場合に適した原語です——1ページのデータ取得、独立したヘルスチェックのセットの呼び出し、リクエストのシャードへの扇出し(fanning out)。

セマフォチャネル: バックプレッシャーを機能として

容量 n のバッファ付きチャネルは、依存関係なしのセマフォです:

sem := make(chan struct{}, 8)
var wg sync.WaitGroup
for _, item := range items {
    wg.Add(1)
    sem <- struct{}{} // 8つのタスクが進行中の場合ブロック
    go func() {
        defer wg.Done()
        defer func() { <-sem }()
        _ = process(ctx, item)
    }()
}
wg.Wait()

取得は go ステートメントの前に起こるため、ゴルーチンはスロットが存在する場合にのみ作成されます。バッファサイズが契約です:それは同時に、並行処理の上限であり、待機中の入場者たちのキューでもあります。

この原語がその場所を正当化するのは、グループが表現しない入場制御が必要なときです——コンテキスト認識の取得、下流ごとのプール、または失敗時の部分結果:

select {
case sem <- struct{}{}: // スロットを取得
case <-ctx.Done():      // 入る前に呼び出し元が諦めた
    return ctx.Err()
}

同じワークロードで測定すると、セマフォチャネルと SetLimit(8) は1.5%以内で完了します——それらの間の選択は、パフォーマンスではなくセマンティクスです。

ワーカープール: バッチではなくストリーム向け

長寿命ワーカーの固定プールは、入場(キュー)と実行(ワーカー)を分離します:

func pool(ctx context.Context, workers int) chan<- func() {
    jobs := make(chan func())
    var wg sync.WaitGroup
    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for {
                select {
                case f, ok := <-jobs:
                    if !ok {
                        return
                    }
                    f()
                case <-ctx.Done():
                    return
                }
            }
        }()
    }
    return jobs
}

プールは、連続して到着する作業に適しています——キューコンシューマー、ポーリングループ、リクエストワーカー——バッチごとのグループが何度も作成され破棄されるのを防ぐためです。Go 1.25以降、sync.WaitGroup.Go は共通の形状に対して Add/Done ペアリングのボイラープレート(定型コード)を削除します。シャットダウンの契約は正しくするべき部分です:jobs をクローズすると、キューが空になった後にワーカーを終了させます;ctx をキャンセルすると、キューに入った作業を放棄します。どちらも正当ですが、ドキュメント化された動作として機能するのは1つだけです。

有界キュー: ブロックするか、破棄するか

バッファ付きチャネルは、上限のあるキューでもあります。バッファが満杯になると、プロデューサーがブロックします——そして、その瞬間が設計上の判断です。ブロックは、キューに供給する誰かに上流へプレッシャーを伝播させます;破棄は、失われた作業の代価にロードを削減します。モニタリングパイプラインは破棄できます;注文パイプラインはブロックするか、永続ストレージに溢出(spill)しなければなりません。

select {
case queue <- event: // 承認された
default:            // キュー満杯: ここに破棄ポリシーがある
    dropped.Add(1)
}

ポリシーが何であれ、明示的かつ測定可能なものにしてください。沈黙の default ブランチは、イベントが消える場所です。

バウンディングのコスト — 測定による

同じワークロード、1つあたり5msのシミュレーションI/Oを行う10,000個のタスク、Go 1.27.1:

戦略 ピーク時のゴルーチン数 実時間(Wall time)
タスクごとに無制限な go ~10,000 起動 18 ms
errgroup.SetLimit(8) 8 6.9 s
セマフォチャネル (cap 8) 8 6.6 s

無制限の実行は実時間で勝利します。実験は何もバックプレッシャーを押し返さないからです——1万個のタイマーゴルーチンは単に並行して睡眠します。まさにそれが罠です。無制限なファンアウトのコストは、あなたのプロセス内のCPUであることは決してありません。それは1万の同時接続、バーストの下でタイムアウトし始める下流、そしてあなたが呼び出している相手のキュー深さに従うメモリグラフです。実際の依存関係に対して、無制限の実行は18 msでは完了しません——タイムアウトします。バウンディングされた実行が6.6秒を支払うのは、10,000タスク ÷ 8スロット × 5 ms が回避できない6.25秒のシリアル化であり、上限がポイントだからです:制御されないバーストを予測可能な6.25秒の排他(drain)に変えます。

上限がわずかのスロットより高い場合、1つの改良が重要です:上限を下流が並行して持続できる能力以下に保ち、その制限(接続プールのサイズ、レートリミット、ワーカー容量)から導き出してください。「適度な」並行性についての感覚ではなく。

バウンディングされたコードの測定とテスト

3つのシグナルが大部分をカバーします。ゴルーチン数のトレンド(Prometheus Goコレクタ経由でエクスポートされる /sched/goroutines:goroutines)は、リークを傾きとして捕捉します。スケジューラ遅延(/sched/latencies:seconds)は飽和を捕捉します——CPU時間を待っているランナブルなゴルーチンは、数が関係なく実際のプレッシャーです。2つのゴルーチンダンプのpprof差分は、停止したゴルーチンがどこにいるかを特定します。

テストのために、バウンディングされたワーカーは Testing Concurrent Go Code with synctest が作られたケースです:フェイクタイムにより、10,000タスクの完全な排他がミリ秒で実行され、synctest.Wait は静穏性(quiescence)のための「sleepして期待する」に代わります、そしてワーカーが永続的にブロックされたままであればバブルは速やかに失敗します——上記のリーク再現は本番環境ではなく、テストバブル内で死にます。

サービススケーリングでは、同じバウンディングの問題はサービスの間に再登場し、そこでセマフォはレートリミットになり、サーキットブレイカーが下流を守ります——Goにおけるサーキットブレイカーパターン はそのレイヤーをカバーし、GoマイクロサービスによるAI/MLオーケストレーション は、オーケストレーション層が取るキューバックの形式を示します。

選択方法

flowchart TD A[並行して実行する作業] --> B{固定バッチまたは連続ストリーム?} B -- 連続ストリーム --> C[ワーカープールまたは有界キュー] B -- 固定バッチ --> D{最初の失敗時に
部分結果が必要?} D -- はい --> E[WaitGroup + errors.Join
またはセマフォチャネル] D -- いいえ --> F[errgroup.WithContext] F --> G{並行処理の上限が必要?} G -- はい --> H[errgroup.SetLimit] G -- いいえ --> I[限界なしのグループ
バッチが証明可能なほど小さい場合のみ] C --> J{キュー満杯: ブロックまたは破棄?} J -- ブロック --> K[上流へのバックプレッシャー] J -- 破棄 --> L[ロードを削減、破棄数をカウント]

テーブルとフローチャートは1つの習慣に圧縮されます:go を書く前に、上限と所有者の名前を付けます。上限は SetLimit、セマフォ、または有界キューです;所有者は Wait、すべてのチャネル操作での ctx.Done アーム、またはその両方です。それら2つの名前を一目で答えられるゴルーチン生成サイトは、午前3時に誰かをページ(呼び出し)することのないサイトです。

関連リンク

購読する

システム、インフラ、AIエンジニアリングの新記事をお届けします。