イベントループ(epoll 風の I/O 多重化)
1 本のスレッドで何千もの接続を同時にさばく、Node.js・nginx・Redis の心臓部を Go でモデル化する。1 接続 = 1 スレッドでは、接続の数だけスレッドが要って破綻する。イベントループは接続を指す番号(FD)を全部ノンブロッキングにし、準備できたものを epoll にまとめて尋ね、そこだけを 1 本のスレッドで処理する。止まるのは全接続をまとめて待つ間だけだ。
この章で作るもの
サーバの仕事は「接続を待ち、来たら読み、返事を書く」。1 本なら何も難しくない。難しいのは同時に何千本もあるときだ。素朴な方法は「1 接続 = 1 スレッド」だが、read はデータが来るまでそのスレッドを止める(ブロックする)ので、接続の数だけスレッドが要る。1 万接続で 1 万スレッド、コンテキストスイッチとスタックのメモリで破綻する。これが有名な C10K 問題だ。
答えは、待ち方を根本から変えること。FD(ファイルディスクリプタ。接続やファイルを OS が指すための番号)を全部ノンブロッキングにして「待たない」ようにし、代わりに epoll に「監視したい FD 全部のうち、どれか準備できたら教えて」と一括で尋ねる。返ってきた「準備できた FD」だけを 1 本のスレッドで順に処理する。
registered FDs ┌──── epoll_wait ────┐ dispatch(1 スレッド)
┌ c1: r ┐ │ │ ┌───────────────┐
┌ c2: r ┐ ──────────▶│ ready {c1:r, c3:r} │──────▶│ c1.handler │
┌ c3: r ┐ │ (準備できたものだけ)│ │ c3.handler │
└───────┘ └────────────────────┘ └───────┬───────┘
▲ │
└───────────────── また epoll_wait へ ◀─────────────────┘
(ready 無しなら次の到着まで眠る)順に見ていく。
- ノンブロッキング I/O = 待たない約束:
read/writeはスレッドを止めず、進められなければ即ErrWouldBlock(EAGAIN)を返す。ブロックした瞬間、他の全接続が巻き添えで止まる - readiness 多重化 = 一括問い合わせ: 多数の FD を 1 回の
Wait(epoll_wait)で走査し、準備できたものだけ返す。関心(Interest)でマスクする - level-triggered とバックプレッシャ: 未読が残る限り Readable は報告され続ける。送りたいデータが残ったら Writable を登録し、空いたら書き足す
① ノンブロッキング FD: 待たない約束
すべての出発点は「read/write が決してブロックしない」ことだ。受信バッファが空なら、データを待つ代わりに即 ErrWouldBlock(実機の EAGAIN)を返す。送信バッファが満杯でも同じ。スレッドを止めない、という約束がイベントループを可能にする:
// ErrWouldBlock は、ノンブロッキング FD への読み書きが「今は進められない」
// ことを表す。実機の EAGAIN / EWOULDBLOCK に当たる。スレッドを止める代わりに
// 即座にこれを返し、呼び手は「準備できたら教えて」とイベントループに委ねる。
// このエラーが event loop の全ての出発点——決してブロックしない、という約束だ。
var ErrWouldBlock = errors.New("eventloop: would block")
// Interest は、その FD について通知してほしいイベントの種類。ビットマスク。
// 読み出せるデータを待つなら Readable、書き込む余地を待つなら Writable。
type Interest uint8
const (
Readable Interest = 1 << iota // 読み出せるデータが届いている
Writable // 送信バッファに書き込む余地がある
)
// Has は f を含むかを返す。
func (i Interest) Has(f Interest) bool { return i&f != 0 }
// String は "r" / "w" / "rw" / "-" で表す。
func (i Interest) String() string {
s := ""
if i.Has(Readable) {
s += "r"
}
if i.Has(Writable) {
s += "w"
}
if s == "" {
return "-"
}
return s
}
// FD は 1 本の接続(ソケット)を表す疑似ファイルディスクリプタ。in は
// カーネル受信バッファ(到着済み・未読)、out は送信バッファの空き容量。
// 実メモリ・実ソケット無しで、ノンブロッキング I/O の意味論だけを取り出す。
type FD struct {
id int
name string
in []byte // 受信済みで未読のバイト列
out int // 送信バッファの空き(あと何バイト書けるか)
sent []byte // これまで書き出したバイト列(観察用)
closed bool
}
// ID は FD 番号を返す。
func (f *FD) ID() int { return f.id }
// Name は接続名を返す。
func (f *FD) Name() string { return f.name }
// Sent はこの FD にこれまで書き出した全バイトの写しを返す。
func (f *FD) Sent() []byte { return append([]byte(nil), f.sent...) }
// Buffered は受信バッファに残っている未読バイト数を返す(観察用)。
func (f *FD) Buffered() int { return len(f.in) }
// Read はノンブロッキング読み出し。到着済みバイトを最大 max だけ返す。
// 受信バッファが空なら、スレッドを止めず即 ErrWouldBlock を返す——ここが肝。
func (f *FD) Read(max int) ([]byte, error) {
if f.closed {
return nil, errors.New("eventloop: read on closed fd")
}
if len(f.in) == 0 {
return nil, ErrWouldBlock // 実機の read() が EAGAIN を返す状況
}
n := max
if n > len(f.in) {
n = len(f.in)
}
b := append([]byte(nil), f.in[:n]...)
f.in = f.in[n:]
return b, nil
}
// Write はノンブロッキング書き込み。送信バッファの空き out までしか書けず、
// 空きが無ければ ErrWouldBlock。書けた分だけ返す(部分書き込みがありうる)。
// 「書ききれなかった残り」をどう扱うかが、後述のバックプレッシャの話になる。
func (f *FD) Write(data []byte) (int, error) {
if f.closed {
return 0, errors.New("eventloop: write on closed fd")
}
if f.out == 0 {
return 0, ErrWouldBlock // 送信バッファが満杯。相手が受け取るまで書けない
}
n := len(data)
if n > f.out {
n = f.out
}
f.out -= n
f.sent = append(f.sent, data[:n]...)
return n, nil
}
// ready は、いま発火している Interest を返す。受信バッファに未読があれば
// Readable、送信バッファに空きがあれば Writable。条件が続く限り毎回報告される
// = level-triggered(epoll の既定)。
func (f *FD) ready() Interest {
var r Interest
if len(f.in) > 0 {
r |= Readable
}
if f.out > 0 {
r |= Writable
}
return r
}
// deliver は「ネットワークからバイトが到着した」を表す(受信バッファに積む)。
// 外の世界(loop.Deliver)から呼ぶ。
func (f *FD) deliver(b []byte) { f.in = append(f.in, b...) }
// drain は「送信バッファに空きが n だけ戻った」を表す(相手が受信した)。
func (f *FD) drain(n int) { f.out += n }ready() に注目してほしい。FD の状態(未読があるか / 書く余地があるか)から、いま発火している Interest を導出している。これが次の epoll が読み取る「準備できているか」の実体だ。未読が残る限り Readable が立ち続ける = level-triggered(条件が続く限り毎回報告される。epoll の既定)。
② epoll = 多数の FD への一括 readiness 問い合わせ
Poller が epoll オブジェクトのモデルだ。監視したい FD をあらかじめ登録しておき、Wait(epoll_wait)1 回で「今 ready な FD 全部」をまとめて受け取る。多数の接続を 1 スレッドで見張れるのは、この一括問い合わせのおかげだ:
// Ready は Wait が返す 1 件: どの FD で、どの Interest が発火したか。
type Ready struct {
FD int
Events Interest
}
// Poller は epoll オブジェクトのモデル。関心のある FD をあらかじめ登録して
// おき、Wait 1 回で「今 ready な FD 全部」をまとめて受け取る。多数の接続を
// 1 スレッドで見張れるのは、この一括問い合わせのおかげだ。
//
// 素朴な select(2) は毎回「見たい FD の集合」を丸ごと渡し直し、カーネルは
// 全部を線形に走査した。epoll は関心を一度登録したら再利用でき、準備できた
// FD だけを返すので、接続数が増えても効率が落ちにくい。ここではその「登録
// しておく → 準備できたものだけ返る」という性質をモデル化する。
type Poller struct {
fds map[int]*FD
interest map[int]Interest
order []int // 登録順。走査順を決定的にする
}
// NewPoller は空の Poller を作る。
func NewPoller() *Poller {
return &Poller{fds: map[int]*FD{}, interest: map[int]Interest{}}
}
// Add は FD を関心付きで登録する(epoll_ctl ADD)。
func (p *Poller) Add(f *FD, in Interest) {
if _, ok := p.fds[f.id]; !ok {
p.order = append(p.order, f.id)
}
p.fds[f.id] = f
p.interest[f.id] = in
}
// Modify は登録済み FD の関心を差し替える(epoll_ctl MOD)。
// 「送りたいデータが溜まったら Writable を足し、送り終えたら外す」のに使う。
// これを怠って Writable を張りっぱなしにすると、送信バッファが空いている限り
// 毎回 ready と報告され、ループが空回りする(busy loop)——実務の頻出バグ。
func (p *Poller) Modify(id int, in Interest) {
if _, ok := p.fds[id]; ok {
p.interest[id] = in
}
}
// Remove は FD を監視対象から外す(epoll_ctl DEL)。接続を閉じたとき。
func (p *Poller) Remove(id int) {
delete(p.fds, id)
delete(p.interest, id)
for i, x := range p.order {
if x == id {
p.order = append(p.order[:i], p.order[i+1:]...)
break
}
}
}
// Wait は今この瞬間に ready な FD を、登録順(決定的)に返す。各 FD の現在の
// readiness と登録した関心の積(AND)を取り、空でなければ 1 件立てる。
// これが epoll_wait の中核——ブロッキングを除けば、監視中の FD を走査して
// 準備できたものだけ返す操作そのものだ。level-triggered(条件が続く限り毎回)。
func (p *Poller) Wait() []Ready {
var out []Ready
for _, id := range p.order {
ev := p.fds[id].ready() & p.interest[id]
if ev != 0 {
out = append(out, Ready{FD: id, Events: ev})
}
}
return out
}
// Len は監視中の FD 数を返す。
func (p *Poller) Len() int { return len(p.order) }
// Watching は監視中の FD 番号を登録順に返す(観察用)。
func (p *Poller) Watching() []int {
return append([]int(nil), p.order...)
}
// InterestOf は登録された関心を返す。未登録なら 0。
func (p *Poller) InterestOf(id int) Interest { return p.interest[id] }素朴な select(2) は毎回「見たい FD の集合」を丸ごと渡し直し、カーネルは全部を線形に走査した。epoll は関心を一度登録したら再利用でき、準備できた FD だけを返すので、接続数が増えても効率が落ちにくい。ここが「なぜ epoll なのか」の核心だ。
Wait は各 FD の readiness と登録した関心(Interest)の AND を取る。だから、送信バッファが空いていても Writable を待っていなければ報告されない。逆に、Writable を待ちっぱなしにすると、空いている限り毎回 ready と返ってきてループが空回りする(busy loop)。実務の頻出バグで、Modify で関心を出し入れして防ぐ。
③ イベントループ: poll → dispatch を回す
ここが心臓部。Loop は poller に「どれか準備できたか」を尋ね、準備できた FD のハンドラだけを呼ぶ。これを繰り返すのがループの全て。ready が 1 つも無ければ、次の到着まで論理時計を飛ばして「眠る」。実機で epoll_wait がブロックして CPU を手放すのに当たる:
// Phase はトレース 1 行の種類。イベントループの 1 周がどの段階かを表す。
type Phase string
const (
PhasePoll Phase = "poll" // epoll_wait が ready 集合を返した
PhaseDispatch Phase = "dispatch" // ある FD のハンドラを呼んだ
PhaseIdle Phase = "idle" // ready が無く、次の到着まで眠った(epoll_wait のブロック)
)
// Step はトレース 1 行。At はその周が起きた論理時刻。
type Step struct {
At int
Phase Phase
FD int // dispatch のとき対象 FD 番号。それ以外は -1
Note string
}
// worldEvent は「外の世界」で起きること(ネットワーク到着 / 送信バッファ解放)。
// 実機では非同期に起きるが、ここでは tick で予定して決定化する。
type worldEvent struct {
at int
seq int // 同時刻の適用順を決定的にする
fn func()
}
// Loop は 1 本のスレッドで多数の FD をさばくイベントループ(リアクタ)。
// poller に「どれか準備できたか」を尋ね、準備できた FD のハンドラだけを呼ぶ。
// これを繰り返すのがループの全て——決してどれか 1 本の接続で立ち止まらない。
type Loop struct {
poller *Poller
fds map[int]*FD
handlers map[int]func(*FD, Interest)
nextID int
clock int
world []worldEvent
worldSeq int
trace []Step
}
// NewLoop は空のイベントループを作る。
func NewLoop() *Loop {
return &Loop{
poller: NewPoller(),
fds: map[int]*FD{},
handlers: map[int]func(*FD, Interest){},
nextID: 1,
}
}
// Open は新しい接続(FD)を作る。sendBuf は送信バッファの初期容量
// (相手が受け取る前に書ける最大バイト数)。
func (l *Loop) Open(name string, sendBuf int) *FD {
f := &FD{id: l.nextID, name: name, out: sendBuf}
l.nextID++
l.fds[f.id] = f
return f
}
// Register は FD をイベントループに登録する。in で待つイベントを、handler で
// ready 時の処理を渡す。handler はノンブロッキングに Read/Write し、決して
// ブロックしてはならない(ブロックした瞬間、他の全接続が巻き添えで止まる)。
func (l *Loop) Register(f *FD, in Interest, handler func(*FD, Interest)) {
l.poller.Add(f, in)
l.handlers[f.id] = handler
}
// Watch は登録済み FD の関心を差し替える(送りたいデータができたら Writable を
// 足す/送り終えたら外す)。
func (l *Loop) Watch(f *FD, in Interest) { l.poller.Modify(f.id, in) }
// CloseFD は接続を閉じ、監視対象から外す。
func (l *Loop) CloseFD(f *FD) {
f.closed = true
l.poller.Remove(f.id)
delete(l.handlers, f.id)
}
// Deliver は「tick 時にこの FD へ data が届く」と予定する(ネットワーク到着)。
func (l *Loop) Deliver(at int, f *FD, data []byte) {
b := append([]byte(nil), data...)
l.schedule(at, func() { f.deliver(b) })
}
// FreeWrite は「tick 時にこの FD の送信バッファが n だけ空く」と予定する
// (相手がそこまで受信し、こちらが続きを書けるようになる)。
func (l *Loop) FreeWrite(at int, f *FD, n int) {
l.schedule(at, func() { f.drain(n) })
}
func (l *Loop) schedule(at int, fn func()) {
l.world = append(l.world, worldEvent{at: at, seq: l.worldSeq, fn: fn})
l.worldSeq++
}
// applyWorld は clock 時点までに予定された外界イベントを、時刻 → 登録順で
// 適用する(決定的)。
func (l *Loop) applyWorld() {
var due, rest []worldEvent
for _, e := range l.world {
if e.at <= l.clock {
due = append(due, e)
} else {
rest = append(rest, e)
}
}
sort.SliceStable(due, func(i, j int) bool {
if due[i].at != due[j].at {
return due[i].at < due[j].at
}
return due[i].seq < due[j].seq
})
for _, e := range due {
e.fn()
}
l.world = rest
}
// nextWorldAt は未適用の外界イベントのうち最も早い tick を返す。無ければ -1。
func (l *Loop) nextWorldAt() int {
next := -1
for _, e := range l.world {
if next == -1 || e.at < next {
next = e.at
}
}
return next
}
// Tick はイベントループの 1 周: 外界を反映 → poll → ready を順に dispatch。
// ready が 1 件も無ければ、次の到着まで clock を進めて「眠る」(epoll_wait が
// ブロックして CPU を手放すのに当たる)。もう予定が無ければ false(ループ終了)。
func (l *Loop) Tick() bool {
l.applyWorld()
ready := l.poller.Wait()
if len(ready) == 0 {
// 準備できた FD が無い。次の到着まで眠る——ここで初めてスレッドは
// 止まるが、止まるのは「全接続まとめて待つ epoll_wait の中」であって、
// どれか 1 本の read の中ではない。os 編の idle 空転と同じ発想。
next := l.nextWorldAt()
if next == -1 {
return false // 予定なし → 何も起こらない。ループ終了
}
l.trace = append(l.trace, Step{At: l.clock, Phase: PhaseIdle, FD: -1,
Note: "wait until t=" + strconv.Itoa(next)})
l.clock = next
return true
}
// epoll_wait が返す ready 集合。1 回の問い合わせで「今さばける FD 全部」。
l.trace = append(l.trace, Step{At: l.clock, Phase: PhasePoll, FD: -1,
Note: readyNote(ready)})
// 準備できた FD だけを順に処理する。ハンドラはノンブロッキングなので、
// この for が終われば 1 周分の仕事が終わる——1 スレッドで多重化できる核心。
for _, r := range ready {
f := l.fds[r.FD]
h := l.handlers[r.FD]
if h == nil || f.closed {
continue
}
l.trace = append(l.trace, Step{At: l.clock, Phase: PhaseDispatch, FD: r.FD,
Note: r.Events.String()})
h(f, r.Events)
}
l.clock++
return true
}
// Run はループを回し続け、これ以上何も起きなくなったら止めてトレースを返す。
func (l *Loop) Run() []Step {
for l.Tick() {
}
return l.trace
}
// Clock は現在の論理時刻を返す。
func (l *Loop) Clock() int { return l.clock }
// Trace はこれまでのトレースを返す。
func (l *Loop) Trace() []Step { return l.trace }
// Poller は内部の Poller を返す(観察用)。
func (l *Loop) Poller() *Poller { return l.poller }
// readyNote は ready 集合を "ready {fd1:r, fd2:rw}" の形の文字列にする。
func readyNote(rs []Ready) string {
parts := make([]string, len(rs))
for i, r := range rs {
parts[i] = "fd" + strconv.Itoa(r.FD) + ":" + r.Events.String()
}
return "ready {" + strings.Join(parts, ", ") + "}"
}ここで os 編の協調スケジューラと同じ骨格が見える。どれか 1 本の接続で立ち止まらず、準備できたものへ制御を回し続ける。眠るのは「全接続まとめて待つ epoll_wait の中」だけで、Tick の PhaseIdle は os 編の idle 空転(HLT)とまったく同じ発想だ。違うのは切り替えの合図で、あちらは yield、こちらは「FD の readiness」になっている。
ハンドラがノンブロッキングであることが絶対条件だ。もしハンドラの中で本当にブロックする処理(同期的な DB クエリ、重い計算)をすれば、その間ループ全体が止まり、他の全接続が待たされる。「イベントループを塞ぐな(Don't block the event loop)」が Node.js の第一戒律なのはこのためだ。
動かす
下のデモは、この epoll 風イベントループをそのままブラウザで動かしている(Go 実装の考え方を JS に移植)。3 本の接続に、時刻をずらしてデータが届く。「1手すすめる」で 1 周ぶんだ。epoll_wait が今 ready な接続を返し、ループがそれらを順に処理(読んでエコー)する様子を追える。1 本のループカーソルが、準備できた接続を渡り歩くのが見えるはずだ。同じことをブロッキング方式でやれば接続の数だけスレッドが要る。それが 1 本で済んでいる。
epoll_wait → ready {c1:r}
設計の観点: I/O モデルとトレードオフ
- ブロッキング(スレッド per 接続)vs イベント駆動: 前者はコードが素直(上から下へ read→処理→write)だが、接続数だけスレッドが要り C10K で破綻する。後者は 1 スレッドで多重化できるが、処理を細切れのコールバック/状態機械に分解する必要があり、コードが「ひっくり返る」(制御の反転)
- なぜ select ではなく epoll/kqueue か:
select/pollは毎回全 FD 集合を渡し直し O(N) で走査する。epoll(Linux)/kqueue(BSD/macOS)は関心を一度登録して再利用し、準備できた FD だけを返すので、アイドルな接続が大量にあっても効率が落ちない。これが大規模サーバの前提技術 - level-triggered vs edge-triggered: LT は「条件が続く限り毎回通知」で扱いやすい。ET は「状態が変化した瞬間だけ通知」で通知回数は減るが、一度に読めるだけ読み切る責任が生まれる(読み残すと次の通知が来ない)。nginx は ET を使い込む
- イベントループを塞ぐな: ハンドラ内でブロックすると全接続が止まる。重い処理はワーカスレッドプール(libuv の
uv_queue_work)や別プロセスへ逃がす。CPU バウンドな仕事はイベントループの苦手分野 - バックプレッシャ: 相手が遅くて送信バッファが詰まったら、書けるまで Writable を待ち、掃けたら関心を外す。これを怠るとメモリに送信待ちが溜まり続ける
メリット・デメリットと実例
| モデル | 多重化 | コードの素直さ | CPU バウンド | 実例 |
|---|---|---|---|---|
| ブロッキング(スレッド per 接続) | 弱い(接続数=スレッド数) | 素直(上から下へ) | 耐える | 素朴な Java/Go 以前のサーバ、Apache prefork |
| イベントループ(epoll/kqueue) | 強い(1 スレッドで多数) | 崩れる(コールバック/状態機械) | 弱い(塞ぐと全滅) | Node.js、nginx、Redis、libuv |
| goroutine(M:N + ネット poller) | 強い(ブロッキング風に書ける) | 素直(見た目は同期) | 耐える | Go の net パッケージ |
| スレッドプール + ブロッキング | 中(プール分だけ) | 素直 | 耐える | 従来型 Java(Servlet)、DB コネクションプール |
裏どり:
- nginx: マスタ + ワーカプロセス構成で、各ワーカが epoll(Linux)/kqueue(BSD)のイベントループを回す。少数プロセスで大量接続をさばく設計。edge-triggered を使う
- Redis: 単一スレッドのイベントループ(
aeライブラリ、epoll/kqueue/select を抽象化)。「シングルスレッドなのに速い」のは、メモリ内処理が短く、I/O 多重化で待ちを潰しているから(6.0 以降 I/O のマルチスレッド化あり) - Node.js / libuv: 1 本のイベントループ(libuv)で JS を回す。だから「イベントループを塞ぐな」が鉄則で、CPU 重い処理は Worker やスレッドプールへ逃がす。ファイル I/O は裏でスレッドプールを使う
- Go の goroutine: 見た目はブロッキング(
conn.Readで待つ)だが、ランタイムのネットワークポーラが内部で epoll/kqueue を使い、待つ goroutine を外して他を走らせる。イベントループの複雑さをランタイムが隠し、開発者は同期的に書ける。本章の仕組みが goroutine の下でも動いている
簡略化したこと
- 実ソケット・syscall なし: FD はバイト列を持つ構造体。
epoll_create/epoll_ctl/epoll_waitの意味論だけを取り出す。到着は tick で予定して決定化する - level-triggered のみ: edge-triggered は扱わない(違いは設計の観点の節で説明)。実機の epoll は両対応
- 単一スレッド固定: マルチスレッド epoll・
SO_REUSEPORTによる負荷分散・ワーカプールは省略 - タイマ/シグナルなし:
epollに混ぜるtimerfd/signalfd/eventfdは扱わない。監視対象は read/write の readiness のみ - 接続の accept は省略: リスニングソケットから新接続を受ける流れ(それ自体が「読める」イベント)は割愛し、接続は最初から用意する
参考資料
- Dan Kegel, "The C10K problem" — 1 万接続をどうさばくか、という問題設定の原典
man epoll/man epoll_ctl/man epoll_wait— Linux の readiness 通知 API。kqueue は BSD/macOS- libuv — Node.js の裏で動くイベントループ。設計ドキュメントが読みやすい
- 実装: foundations/eventloop