Skip to content

APIサーバとinformer

実装: orchestration/apiserver/ / 実行: go test ./orchestration/apiserver/

この編のコントローラはずっと「現状を数える」と言ってきた。だがどうやって知るのか。毎回全件を問い合わせれば重く、たまにしか聞かなければ遅い。答えは手元に写しを持つことで、一度だけ全件を取り、以降は変更だけを流してもらう。そして流れは切れる。写しは古くなりうる。この編が level-triggered にこだわってきた理由が、ここで分かる。

この章で作るもの

ここまでのすべての章は、ある仕掛けの上に成り立っていた。コントローラは現状を知っている、という前提だ。調整ループは Pod の数を数えたし、スケジューラは空きを見たし、DaemonSetはノードの一覧を見た。だが、どうやって知るのか。

素朴には、毎回全件を問い合わせればよい。だが調整は頻繁に走るし、コントローラは何十個もある。全員が全件を毎秒問い合わせたら、API サーバが潰れる。かといって、たまにしか問い合わせないと、障害に気づくのが遅くなる。

答えは、手元に写しを持つことだった。最初に一度だけ全件を取り、以降は変更だけを流してもらって写しを更新する。読むときは写しを見るので、API サーバには触らない。これが informer で、この編のコントローラが「現状を数える」と言っていたときに実際に見ていたものになる。

そして、この仕組みには弱点がある。流れが切れると、切れている間の変更を取りこぼす。

  置き場(唯一の真実)              写し(informer)

  版 5: web-1, web-2       ──全件──→  版 5: web-1, web-2
                            ──差分──→
  版 6: web-3 を追加        ──差分──→  版 6: web-1, web-2, web-3

        watch が切れる            ✕

  版 7: web-1 を削除          届かない    版 6 のまま
  版 8: web-4 を追加          届かない    web-1 がまだ居るように見える

        張り直す                 ──差分──→  版 8 に追いつく
        (履歴が残っていれば。無ければ全件を取り直す)
一度だけ全件、以降は差分。流れが切れている間、写しは古いまま取り残される

順に見ていく。

  1. 一度だけ全件、以降は差分: 毎回全件は重く、差分だけでは始点が無い
  2. 読みは写しから: API サーバに触らないので、何度読んでも安い
  3. 写しは古くなる: だから level-triggered でなければ、この上では動けない

① 唯一の真実と、その履歴

まず、置き場を作る:

go

// Object は保存される1つの資源。Version は変更のたびに上がる。
type Object struct {
	Kind    string
	Name    string
	Value   string
	Version int
}

// Event は1回の変更。watch はこれを流す。
type Event struct {
	Type    string // "added" / "modified" / "deleted"
	Object  Object
	Version int // このイベント時点の全体の版
}

// Store は API サーバが持つ唯一の真実。ここが正で、写しは常にこの後を追う。
type Store struct {
	objs    map[string]*Object
	version int
	history []Event // 変更の履歴。watch が途中から流すために持つ

	Reads   int // 全件の読み出し回数(重い操作)
	Watches int // watch を張り直した回数
}

// NewStore は空の置き場を作る。
func NewStore() *Store { return &Store{objs: map[string]*Object{}} }

// Version は今の版を返す。
func (s *Store) Version() int { return s.version }

// Put は資源を作るか書き換える。版が上がり、履歴に積まれる。
func (s *Store) Put(kind, name, value string) Object {
	s.version++
	kt := "modified"
	o, ok := s.objs[kind+"/"+name]
	if !ok {
		kt = "added"
		o = &Object{Kind: kind, Name: name}
		s.objs[kind+"/"+name] = o
	}
	o.Value = value
	o.Version = s.version
	s.history = append(s.history, Event{Type: kt, Object: *o, Version: s.version})
	return *o
}

// Delete は資源を消す。これも版が上がり、履歴に積まれる。
func (s *Store) Delete(kind, name string) bool {
	o, ok := s.objs[kind+"/"+name]
	if !ok {
		return false
	}
	s.version++
	delete(s.objs, kind+"/"+name)
	s.history = append(s.history, Event{Type: "deleted", Object: *o, Version: s.version})
	return true
}

// List は全件を返す。重い操作で、写しを最初に埋めるときだけ呼ぶ。
func (s *Store) List(kind string) []Object {
	s.Reads++
	var keys []string
	for k, o := range s.objs {
		if o.Kind == kind {
			keys = append(keys, k)
		}
	}
	sort.Strings(keys)
	out := make([]Object, 0, len(keys))
	for _, k := range keys {
		out = append(out, *s.objs[k])
	}
	return out
}

// Since は版 from より後の変更を返す。watch が途中から追いつくために使う。
//
// 履歴が残っている範囲でしか答えられない。古すぎる版から聞かれたら、
// 差分では答えられないので、全件を取り直してもらうしかない。
func (s *Store) Since(from int) ([]Event, bool) {
	if len(s.history) == 0 {
		return nil, true
	}
	oldest := s.history[0].Version
	if from < oldest-1 {
		return nil, false // 古すぎる。差分では追いつけない
	}
	var out []Event
	for _, e := range s.history {
		if e.Version > from {
			out = append(out, e)
		}
	}
	return out, true
}

// Compact は古い履歴を捨てる。無限には持てないので、いつか捨てる。
// 捨てた後は、それより古い版から追いつくことができなくなる。
func (s *Store) Compact(keep int) {
	if len(s.history) > keep {
		s.history = s.history[len(s.history)-keep:]
	}
}

Store が唯一の真実で、version が変更のたびに単調に増える。この版が、写しが「どこまで追いついているか」を表す目盛りになる。

history に変更を積んでいるのが、差分で追いつくための仕掛けになる。写しが版5まで見ていると分かっていれば、版6以降のイベントだけを渡せばよい。全件を送らなくて済む。

だが Compact があることに注意が要る。履歴は無限には持てないので、いつか古い分を捨てる。捨てた後は、それより古い版から聞かれても差分では答えられない。Since が偽を返すのがその場合で、そのときは全件を取り直してもらうしかない。実物の etcd も履歴を圧縮するので、長く切れていた watch は同じ理由で失敗する。

② 写しを保つ

informer 側はこうなる:

go

// Informer は手元の写し。読むときはここを見るので、API サーバに触らない。
type Informer struct {
	kind    string
	store   *Store
	cache   map[string]Object
	version int // どこまで追いついているか

	Connected bool // watch が繋がっているか
	Resyncs   int  // 全件を取り直した回数

	Log []string
}

// NewInformer は写しを作る。この時点ではまだ空で、Start で埋まる。
func NewInformer(s *Store, kind string) *Informer {
	return &Informer{kind: kind, store: s, cache: map[string]Object{}}
}

// Start は最初に一度だけ全件を取り、そこから watch を張る。
//
// この「一度だけ全件、以降は差分」が肝になる。毎回全件を取れば正確だが
// 重い。差分だけを追えば軽いが、始点が要る。だから最初に一度だけ全件を取る。
func (i *Informer) Start() {
	for _, o := range i.store.List(i.kind) {
		i.cache[o.Name] = o
	}
	i.version = i.store.Version()
	i.Connected = true
	i.store.Watches++
	i.Resyncs++
	i.logf("全件を読み込んで watch を張った(版 " + itoa(i.version) + ")")
}

// Disconnect は watch が切れた状態にする。以降の変更は届かなくなる。
func (i *Informer) Disconnect() {
	i.Connected = false
	i.logf("watch が切れた。この間の変更は届かない")
}

// Sync は届いた変更を写しに反映する。繋がっていなければ何もしない。
//
// 繋がっていない間に起きた変更は、ここでは反映されない。だから写しは
// 古くなる。読み手は古い写しを見て判断することになる。
func (i *Informer) Sync() {
	if !i.Connected {
		return
	}
	events, ok := i.store.Since(i.version)
	if !ok {
		// 履歴が捨てられていて追いつけない。全件を取り直すしかない。
		i.logf("履歴が古すぎて差分で追いつけない。全件を取り直す")
		i.cache = map[string]Object{}
		i.Start()
		return
	}
	for _, e := range events {
		if e.Type == "deleted" {
			delete(i.cache, e.Object.Name)
			continue
		}
		i.cache[e.Object.Name] = e.Object
	}
	i.version = i.store.Version()
}

// Reconnect は watch を張り直し、切れている間の変更に追いつく。
func (i *Informer) Reconnect() {
	i.Connected = true
	i.store.Watches++
	i.logf("watch を張り直した")
	i.Sync()
}

// List は写しから読む。API サーバには触らないので、何度呼んでも安い。
func (i *Informer) List() []Object {
	names := make([]string, 0, len(i.cache))
	for n := range i.cache {
		names = append(names, n)
	}
	sort.Strings(names)
	out := make([]Object, 0, len(names))
	for _, n := range names {
		out = append(out, i.cache[n])
	}
	return out
}

// Get は写しから1件返す。
func (i *Informer) Get(name string) (Object, bool) {
	o, ok := i.cache[name]
	return o, ok
}

// Stale は写しが最新から遅れているかを返す。
func (i *Informer) Stale() bool { return i.version < i.store.Version() }

// Lag は何版ぶん遅れているかを返す。
func (i *Informer) Lag() int { return i.store.Version() - i.version }

Start が最初に一度だけ全件を読み、そこから watch を張る。List は写しから読むだけなので、API サーバには触らない。テストで、100回読んでも全件読み出しが増えないことを固定した。

この安さが、この編のすべてを支えている。調整ループが「毎回まるごと現状を数え直す」と言えたのは、数え直しが安いからだ。もし毎回 API サーバに全件を聞いていたら、level-triggered は成立しない。重すぎて、イベント駆動にせざるを得なくなる。level-triggered という設計判断と、写しを持つという実装は、互いに支え合っている。

Sync が差分を反映する。Since が失敗したときは全件を取り直す。この取り直しが、実物では resync と呼ばれる処理になる。テストで、履歴が捨てられていると取り直し、残っていれば差分で足りることを固定した。

③ 写しは古くなる

DisconnectReconnect が、この章でいちばん大事な部分になる。

watch が切れている間、変更は届かない。写しは切れた時点のまま止まる。読み手はその古い写しを見て判断する。消えたはずの Pod がまだ居るように見えるし、増えたはずの Pod は見えない。テストで、切れている間に消した Pod が写しには残り続けることを固定した。

ここで、この編がずっと level-triggered にこだわってきた理由がつながる。調整ループの章で「イベントを取りこぼしても、次に現状を数え直せば追いつく」と書いた。あれは理屈の上の話ではなく、この土台の性質そのものだった。写しは実際に古くなるし、実際にイベントを取りこぼす。取りこぼしを前提にしていない設計は、この上では動かない。

張り直したときの振る舞いも、同じことを言っている。Reconnect は途中の経過を再生しない。切れている間に web-3 が作られて更新されて、という経過があっても、届くのは最終の姿だけだ。テストで、2回変わったものが最終の値になることを固定した。イベントの列を正確に追うのではなく、最終の状態に追いつく。これも level-triggered の形になっている。

動かす

下のデモは、置き場と写しを並べて見る。繋がっている間は差分が届いて追いつく。「watch を切る」を押してから Pod を作ったり消したりすると、写しだけが取り残されて赤く印がつく。張り直すと追いつくが、切れている間の変更が多すぎると履歴が足りず、全件を取り直すことになる。

デモAPIサーバとinformer真実 版2 / 写し 版2
watch: 繋がっている ・ 全件読み出し 1 回 ・ 履歴は直近 4 件だけ保持
置き場(唯一の真実)版 2
web-1running
web-2running
写し(informer)版 2
web-1running
web-2running
写しが真実に追いついている。読むのは写しなので、何度読んでも置き場には触らない
起きたこと
(まだ何も起きていない)

左が唯一の真実、右がコントローラの手元にある写し。繋がっている間は、変更が差分で届いて写しが追いつく。 「watch を切る」を押してから Pod を作ったり消したりすると、写しだけが取り残される。消えたはずのものが 写しには残り、読み手はそれを見て判断する。張り直せば追いつくが、履歴は直近 4 件しか 残っていないので、切れている間に変更が多すぎると差分では追いつけず、全件を取り直すことになる。

設計の観点

  • 安く読めることが設計を決める: 写しから読めるから level-triggered が成立する。読みが高ければ、イベント駆動にせざるを得ず、取りこぼしに弱くなる
  • 写しは常に遅れる前提で書く: 読んだ値がすでに古いことがある。書き込むときは版を確認する(楽観ロック)など、別の手当てが要る
  • 履歴は有限: 長く切れていた写しは差分で追いつけない。取り直しの経路を必ず用意しておく
  • 最終状態に追いつけばよい: 途中の経過を再生しない。何が起きたかでなく、今どうであるかだけを見る形になっている
  • レプリケーションと同じ形: 一度スナップショットを取り、以降はログを流して追いつく。写しが遅れることも、ログが尽きたら取り直すことも同じ
  • 真実は1つに集める: すべての状態が置き場に集まっているから、どのコントローラも同じものを見て判断できる。分散した状態を突き合わせる必要がない

対照と実例

毎回問い合わせる写しを持つ(informer)
読みの負荷コントローラの数に比例最初の1回だけ
反応の速さ問い合わせの間隔ぶん遅れる変更が届き次第
正確さ常に最新遅れることがある
切れたとき次の問い合わせで直る張り直すまで古いまま
必要な設計どちらでもよいlevel-triggered が必須

裏どり:

  • API サーバと etcd: 状態は etcd に保存され、API サーバがその前段で認証・認可・admission を通す
  • watch と resourceVersion: 変更は版を目印に差分として流れる。版が古すぎると Expired になり、全件から取り直す
  • informer と lister: クライアント側で写しを保ち、コントローラは lister 経由で写しから読む
  • 定期的な resync: 実物の informer は、取りこぼしに備えて一定間隔で全件を流し直す設定を持つ

簡略化したこと

  • etcd を作らない: 実物は etcd が保存と watch を担う
  • 合意なし: 実物の etcd は Raft で複数台に複製する
  • 通信なし: watch は HTTP のストリームで流れる。ここでは関数を呼ぶだけ
  • 並行なし: 複数の informer が同時に動く状況は扱わない
  • workqueue なし: 実物は変更をキューに積み、重複をまとめてから reconcile を呼ぶ

参考資料