Skip to content

スケジューラ(M:N スケジューラと work-stealing)

Go が数百万の goroutine を数個の OS スレッドで回す仕組み(M:N スケジューラと work-stealing)をモデル化する。キューを 1 本にすればロックが詰まり、コアごとに分ければ偏ったとき暇なコアが遊ぶ。Go は P ごとにローカルキューを持たせてロック無しで回し、空いた P が最も混む P からキューの半分を横取りする。偏りが自動で均され、全コアが働き続ける。

この章で作るもの

os 編では CPU 1 個・実行キュー 1 本の協調スケジューラを作った。だが実機はコアが複数ある。素朴に「1 本のキューを全コアで共有」すると、そのロックの奪い合いがボトルネックになる。かといって「コアごとに独立キュー」だと、仕事が一部のコアに偏ったとき、他のコアが手持ち無沙汰になる。Go ランタイムの答えが M:N スケジューラ + work-stealing だ。

        P0                    P1                    P2
   ┌─ local ─┐           ┌─ local ─┐           ┌─ local ─┐
   │ G G G G │           │ (空)    │  ◀─steal──│ (空)    │
   └────┬────┘           └─────────┘   half    └─────────┘
        │ 偏り!                ▲___________________|
        └── P1 が半分を横取り ──┘   暇な P が忙しい P から仕事を盗む

        ┌──────────── global run queue(全 P 共有)────────────┐
        │  ローカルが溢れたぶん / 新規のあふれ                │
        └────────────────────────────────────────────────────┘
GMP モデルと work-stealing。各 P(プロセッサ)は自分のローカル実行キューを持つ。仕事は生成した P に溜まりがちで偏る。ローカルが空になった P は、まずグローバルキューを見て、無ければ他の最も忙しい P からキューの半分を横取りする。これで偏りが自動的に均される

順に見ていく。

  1. P ごとのローカルキュー: 仕事は生成した P のローカルキューに積む。ロック無しで速いが偏りやすい
  2. work-stealing で均す: 空になった P は、最も混んでいる他の P からキューの半分を盗む
  3. グローバルキューとあふれ: ローカルが溢れたら半分をグローバルへ退避し、暇な P はそこからも引く

① G と P: 走らせる仕事と、走らせるコンテキスト

まず登場人物を定義する。G は goroutine で、走らせたい仕事(ここでは work tick を消費し切ったら終わる)。P はプロセッサで、スケジューリングの文脈として自分だけのローカル実行キューを持つ。M(OS スレッド)は「P 1 つにつき常に 1 つある」と割り切って省く(本章の主題は P 間の仕事の均し方だから):

go

// State は G の状態。
type State int

const (
	Runnable State = iota // 実行待ち(どこかのキューにいる)
	Running               // いま P 上で走っている
	Done                  // 実行し終えた
)

func (s State) String() string {
	switch s {
	case Runnable:
		return "runnable"
	case Running:
		return "running"
	case Done:
		return "done"
	default:
		return "?"
	}
}

// G は goroutine。work は必要な仕事量(tick)、executed は消化済み量。
// 本モデルでは「実際の計算」はせず、work tick を消費し切ったら Done になる。
type G struct {
	ID       int
	Name     string
	work     int
	executed int
	st       State
}

// State は現在の状態を返す。
func (g *G) State() State { return g.st }

// Remaining は残りの仕事量を返す。
func (g *G) Remaining() int { return g.work }

// Executed は消化済みの仕事量を返す。
func (g *G) Executed() int { return g.executed }

// P はプロセッサ(スケジューリング文脈)。ローカル実行キューを持ち、その先頭から
// G を取り出して走らせる。ローカルキューがあることで、G の出し入れに毎回グローバルな
// ロックを取らずに済む——マルチコアでスケールする鍵。
type P struct {
	ID     int
	local  []*G // ローカル実行キュー
	ran    int  // この P が実行した総 tick(負荷の観察用)
	steals int  // この P が横取りに成功した回数
}

// QueueLen はローカルキューの長さを返す。
func (p *P) QueueLen() int { return len(p.local) }

// Ran はこの P が実行した総 tick を返す(負荷分散が効いたかの指標)。
func (p *P) Ran() int { return p.ran }

// Steals はこの P が横取りに成功した回数を返す。
func (p *P) Steals() int { return p.steals }

// QueueNames はローカルキューの G 名を先頭(次に走る)から返す。
func (p *P) QueueNames() []string {
	out := make([]string, len(p.local))
	for i, g := range p.local {
		out[i] = g.Name
	}
	return out
}

// tag は P を "P0" のように表す(トレースで横取り元を示すのに使う)。
func (p *P) tag() string { return "P" + strconv.Itoa(p.ID) }

Plocal(ローカルキュー)を持つのが要。G の出し入れがローカルに閉じるので、コアごとにロック無しで回せる。これがマルチコアでスケールする土台だ。ran(実行 tick)は、あとで「負荷が均せたか」を見る指標になる。

② スケジューラ: ローカルに積み、あふれたらグローバルへ

Scheduler は複数の P と、全 P が共有するグローバルキューを持つ。Go は goroutine を生成して P0 のローカルキューに積み、「1 つの goroutine から生まれた仕事は同じ P に溜まる」を模す。ローカルが上限に達したら、半分をグローバルへ退避する(spill):

go

// Kind はスケジュールのトレース1行の種類。
type Kind string

const (
	KindSpawn  Kind = "spawn"  // goroutine を生成し、どこかのキューに積んだ
	KindRun    Kind = "run"    // P が G を1量子(quantum)走らせた
	KindDone   Kind = "done"   // G が仕事を消化し切って終了した
	KindSteal  Kind = "steal"  // 暇な P が他の P から仕事を横取りした
	KindGlobal Kind = "global" // グローバルキューから仕事を引いた
	KindIdle   Kind = "idle"   // 走らせる G が無く、P が遊んだ
	KindSpill  Kind = "spill"  // ローカルキューが溢れ、半分をグローバルへ退避した
)

// Event はトレース1行。At はその出来事が起きた論理時刻(ラウンド開始時刻)。
type Event struct {
	At   int
	P    int // 対象プロセッサ番号(スケジューラ全体の出来事は -1)
	Kind Kind
	G    string // 対象 goroutine(該当しなければ空)
	N    int    // run の tick / steal・global の本数
}

// Scheduler は複数の P を回す M:N スケジューラ。1 ラウンド = 全 P が1量子ずつ
// 同時に進む、というモデルで並行実行を決定的に表す。
type Scheduler struct {
	ps       []*P
	global   []*G // グローバル実行キュー(全 P が共有)
	quantum  int  // 1 回の実行で走らせる tick 数(協調的プリエンプションの粒度)
	localCap int  // ローカルキューの上限。超えると半分をグローバルへ退避
	clock    int
	nextGID  int
	all      []*G // 生成した全 G(検証・観察用)
	trace    []Event
}

// NewScheduler は numP 個のプロセッサを持つスケジューラを作る。quantum は
// 1 回の実行で走らせる tick 数(小さいほど頻繁に切り替わる)。
func NewScheduler(numP, quantum int) *Scheduler {
	if numP < 1 {
		numP = 1
	}
	if quantum < 1 {
		quantum = 1
	}
	ps := make([]*P, numP)
	for i := range ps {
		ps[i] = &P{ID: i}
	}
	return &Scheduler{ps: ps, quantum: quantum, localCap: 6, nextGID: 1}
}

// Go は goroutine を生成し、P0 のローカルキューに積む。1 つの goroutine から
// 生まれた仕事は同じ P に溜まりがち——これが偏りを生み、work-stealing が要る理由。
func (s *Scheduler) Go(name string, work int) *G { return s.GoOn(0, name, work) }

// GoOn は指定した P のローカルキューに goroutine を積む。ローカルが上限に達して
// いたら、半分をグローバルキューへ退避してから積む(Go の runqput のあふれ処理)。
func (s *Scheduler) GoOn(pid int, name string, work int) *G {
	if pid < 0 || pid >= len(s.ps) {
		pid = 0
	}
	if work < 1 {
		work = 1
	}
	g := &G{ID: s.nextGID, Name: name, work: work, st: Runnable}
	s.nextGID++
	s.all = append(s.all, g)

	p := s.ps[pid]
	if len(p.local) >= s.localCap {
		n := len(p.local) / 2
		s.global = append(s.global, p.local[:n]...)
		p.local = append([]*G(nil), p.local[n:]...)
		s.trace = append(s.trace, Event{At: s.clock, P: pid, Kind: KindSpill, N: n})
	}
	p.local = append(p.local, g)
	s.trace = append(s.trace, Event{At: s.clock, P: pid, Kind: KindSpawn, G: name})
	return g
}

なぜ P0 に固めるのか。実際の Go でも、ある goroutine が go f() で生んだ子は、親が乗っている P のローカルキューに入る。だから仕事は自然と偏る。この偏りを次の work-stealing が均す。それを見せるために、あえて偏った積み方をしている。

③ Step と work-stealing: 暇な P が仕事を盗む

心臓部。Step は 1 ラウンド進める。全 P が index 順に「必要なら仕事を確保 → 1 量子走らせる」を 1 回ずつ行う(ラウンド制で並行実行を決定的に表す)。仕事の確保がポイントだ。ローカルが空なら、まずグローバルキューを見て、それも空なら steal で他の最も混んでいる P からキューの半分を横取りする:

go

// Step は1ラウンド進める: 全 P が「必要なら仕事を確保 → 1 量子走らせる」を
// index 順に1回ずつ行う。走らせる仕事がどこにも無ければ false(完了)。
func (s *Scheduler) Step() bool {
	if s.drained() {
		return false
	}
	at := s.clock
	for _, p := range s.ps {
		// ローカルが空なら、まずグローバル、次に他 P から横取りして仕事を確保する。
		if len(p.local) == 0 {
			if !s.fromGlobal(p, at) && !s.steal(p, at) {
				s.trace = append(s.trace, Event{At: at, P: p.ID, Kind: KindIdle})
				continue
			}
		}
		// ローカル先頭の G を1量子走らせる。
		g := p.local[0]
		p.local = p.local[1:]
		g.st = Running
		n := s.quantum
		if g.work < n {
			n = g.work
		}
		g.work -= n
		g.executed += n
		p.ran += n
		if g.work == 0 {
			g.st = Done
			s.trace = append(s.trace, Event{At: at, P: p.ID, Kind: KindRun, G: g.Name, N: n})
			s.trace = append(s.trace, Event{At: at, P: p.ID, Kind: KindDone, G: g.Name})
		} else {
			// 量子を使い切ったがまだ残る → ローカルキューの末尾へ戻す(協調的な切り替え)。
			g.st = Runnable
			p.local = append(p.local, g)
			s.trace = append(s.trace, Event{At: at, P: p.ID, Kind: KindRun, G: g.Name, N: n})
		}
	}
	s.clock += s.quantum // ラウンドぶんの時間が(並行に)経過した
	return true
}

// fromGlobal はグローバルキューから仕事を1バッチ引く。全 P で均されるよう、
// おおよそ「全体 / P 数 + 1」本を取る(1 つの P が総取りしないように)。
func (s *Scheduler) fromGlobal(p *P, at int) bool {
	if len(s.global) == 0 {
		return false
	}
	n := len(s.global)/len(s.ps) + 1
	if n > len(s.global) {
		n = len(s.global)
	}
	batch := s.global[:n]
	s.global = s.global[n:]
	p.local = append(p.local, batch...)
	s.trace = append(s.trace, Event{At: at, P: p.ID, Kind: KindGlobal, N: n})
	return true
}

// steal は最も混んでいる他の P からローカルキューの半分を横取りする。
// 決定的にするため、同数なら index の小さい P を選ぶ。半分盗めない(2 本未満の)
// P は対象にしない——盗むコストに見合わないから。
func (s *Scheduler) steal(thief *P, at int) bool {
	var victim *P
	for _, p := range s.ps {
		if p == thief || len(p.local) < 2 {
			continue
		}
		if victim == nil || len(p.local) > len(victim.local) {
			victim = p
		}
	}
	if victim == nil {
		return false
	}
	n := len(victim.local) / 2
	batch := victim.local[:n]
	victim.local = append([]*G(nil), victim.local[n:]...)
	thief.local = append(thief.local, batch...)
	thief.steals++
	s.trace = append(s.trace, Event{At: at, P: thief.ID, Kind: KindSteal, G: victim.tag(), N: n})
	return true
}

// drained は全ローカルキュー・グローバルキューが空(= 走らせる仕事が無い)かを返す。
func (s *Scheduler) drained() bool {
	if len(s.global) > 0 {
		return false
	}
	for _, p := range s.ps {
		if len(p.local) > 0 {
			return false
		}
	}
	return true
}

// Run は全 goroutine が終わるまでラウンドを回し、トレースを返す。
func (s *Scheduler) Run() []Event {
	for s.Step() {
	}
	return s.trace
}

steal が均衡の要だ。盗むのが半分なのは、盗みすぎず・盗まなさすぎずの勘所だ。全部盗むと今度は相手が空になって盗み返しが起き、少ししか盗まないと何度も盗みに行く羽目になる。半分なら、盗んだ側も盗まれた側も次の仕事を持てる。quantum を使い切った G をローカル末尾へ戻すのは、os 編の yield に当たる協調的な切り替え点だ。

動かす

下のデモは、この M:N スケジューラをそのままブラウザで動かしている。3 つの P があり、goroutine は全部 P0 に積まれる(偏り)。「1手すすめる」で 1 ラウンド進むと、P1・P2 が空になった瞬間に P0 から仕事を半分横取りする様子が見える。各 P の実行量(下のバー)が、偏った投入にもかかわらずだんだん揃っていく。work-stealing が負荷を均しているのが分かるはずだ。

デモscheduler(M:N + work-stealing)ラウンド 0
1 / 11
P0盗んだ回数 0
ローカルキュー
ABCDEF
実行量 0
P1盗んだ回数 0
ローカルキュー
(空 — 盗みに行く)
実行量 0
P2盗んだ回数 0
ローカルキュー
(空 — 盗みに行く)
実行量 0
global queue
(空)

goroutine を全部 P0 に積んだ(偏り)。P1・P2 は空っぽ

各キューの先頭 = 次に走る goroutine実行量バー: 偏った投入でも横取りで揃っていく

設計の観点: なぜ M:N + work-stealing か

  • 1:1 vs M:N: 「1 goroutine = 1 OS スレッド」(1:1)は単純だが、スレッドは生成が重く(数 KB〜MB のスタック)、数万本で破綻する。M:N は軽量な G を少数の M に多重化するので、goroutine を数十万本作っても安い。C10K をイベントループで解くのと同じ動機を、言語ランタイムに埋め込んだ形
  • なぜローカルキューか: 単一グローバルキューはスケールしない。全コアが 1 つのロックを奪い合う。P ごとのローカルキューにすれば、大半の操作がロック無しで済む。その代償(偏り)を work-stealing で埋める
  • work-stealing の性質: 盗む側(暇なコア)がコストを払う。忙しいコアは自分の仕事に集中でき、余ったコアが自律的に仕事を探しに行く。分散設計の定番で、Java の ForkJoinPool・Rust の Tokio・TBB も同じ発想
  • 公平性と飢餓: ローカル優先だとグローバルキューの G が後回しになりうる。実機の Go は約 61 回に 1 回グローバルを優先して見て飢餓を防ぐ。「速さ(ローカル優先)」と「公平さ(たまに全体を見る)」のトレードオフ
  • プリエンプション: 協調点(関数呼び出し)でしか切り替わらないと、タイトなループが P を占有して他を飢えさせる。Go 1.14 以降は非同期プリエンプション(シグナルで割り込む)を足した。os 編で触れた話がここでも効く

メリット・デメリットと実例

方式並列度生成コスト負荷分散実例
1:1(スレッド=タスク)出る高い(スタック重い)OS 任せpthread、Java 旧来スレッド
単一キュー M:Nロック競合で頭打ち低い中央キュー素朴なスレッドプール
M:N + work-stealing(本章)出る低い自律的に均衡Go、Java ForkJoinPool、Rust Tokio、Erlang
イベントループ(1 スレッド)出ない(1 コア)低い不要Node.js、nginx、Redis

裏どり:

  • Go ランタイム: まさに GMP + work-stealing。runtime/proc.goschedule / findrunnable / runqsteal が実物。GOMAXPROCS が P の数(並列度)。Go 1.14 で非同期プリエンプションを追加
  • Java ForkJoinPool: 各ワーカスレッドが自分の deque を持ち、空になったら他から末尾を盗む。parallelStream() の裏でも動く、work-stealing の代表例
  • Rust Tokio / Erlang BEAM: どちらも軽量タスク(async task / プロセス)を少数スレッドに多重化し、work-stealing で均す。M:N は現代の並行ランタイムの共通解
  • 一方 Node.js は逆張り: 1 スレッドのイベントループで多重化する。CPU 並列が要らない I/O 中心の用途では、スケジューラの複雑さを持たないこちらが単純で速い

簡略化したこと

  • M(OS スレッド)を省略: 「P 1 つに M が常に 1 つ」の前提。syscall で M がブロックしたとき P を別 M に渡す handoff・spinning M・M の生成/破棄は扱わない
  • 横取りは決定的: 実機はランダムな P から複数回試みる。ここでは最混雑 P 固定(再現性を優先)
  • グローバルの定期チェック省略: 実機の「約 61 回に 1 回グローバル優先」は入れない(仕組みは設計の観点の節で説明)
  • 仕事は tick 数: G は実際の計算をせず work tick を消費し切ったら終了。チャネル待ち・I/O 待ちでのブロックや goroutine 間の依存は扱わない
  • 非同期プリエンプション無し: 切り替えは量子境界のみ。シグナル割り込みによる強制プリエンプションは省略(os 編で対比を説明)

参考資料

  • Dmitry Vyukov, "Scalable Go Scheduler Design Doc" — Go の GMP スケジューラ設計の原典
  • Go ランタイムの runtime/proc.goschedule / findrunnable / runqsteal の実物
  • Blumofe & Leiserson, "Scheduling Multithreaded Computations by Work Stealing"(1999) — work-stealing の理論的原典
  • 実装: foundations/scheduler