概要 Queueインターフェース DeltaFIFO 要素の追加・削除・更新 - queueActionLocked() Pop() Replace()
概要
ソースコードのバージョン情報
- プロジェクト: kubernetes
- ブランチ: master
- 最新コミットID: d25d741c
- 日付: 2021-09-26
前回の記事「Kubernetes client-goソースコード分析 - 序章」で、カスタムコントローラーが関与するclient-goコンポーネントの全体的なワークフローについて触れました。概要は以下の図のようになります。
DeltaFIFOは上記の重要なコンポーネントの一つです。本日は、client-goにおけるDeltaFIFO関連のコードを詳しく調査します。
Queueインターフェース
workqueueにあるキューの概念と同様に、ここにもキューがあります。Queueインターフェースはclient-go/tools/cacheパッケージのfifo.goファイルに定義されており、以下のメソッドがあります:
type Queue interface {
Store
Pop(PopProcessFunc) (interface{}, error) // ブロックし、要素がpopできるようになるか、キューが閉じられるまで待機する
AddIfNotPresent(interface{}) error
HasSynced() bool
Close()
}
ここではStoreインターフェースが埋め込まれており、その定義は以下の通りです:
type Store interface {
Add(obj interface{}) error
Update(obj interface{}) error
Delete(obj interface{}) error
List() []interface{}
ListKeys() []string
Get(obj interface{}) (item interface{}, exists bool, err error)
GetByKey(key string) (item interface{}, exists bool, err error)
Replace([]interface{}, string) error
Resync() error
}
Storeインターフェースのメソッドは比較的直感的です。Storeの実装は多数存在しますが、Queueで使用されているのはどれでしょうか。
Queueインターフェースの実装はFIFOとDeltaFIFOの2つのタイプです。Informerで使用されているのはDeltaFIFOであり、DeltaFIFOはFIFOに依存していないため、以下ではDeltaFIFOがどのように実装されているかを直接見ていきましょう。
DeltaFIFO
- client-go/tools/cache/delta_fifo.go:97
type DeltaFIFO struct {
lock sync.RWMutex
cond sync.Cond
items map[string]Deltas
queue []string // このqueueには重複する要素はなく、上記のitemsのキーと一致します
populated bool
initialPopulationCount int
keyFunc KeyFunc // 上記のmapで使用するキーを構築するために使用されます
knownObjects KeyListerGetter // すべてのキーを検索するために使用されます
closed bool
emitDeltaTypeReplaced bool
}
ここにはDeltasという型があり、具体的な定義を見てみましょう:
type Deltas []Delta
type Delta struct {
Type DeltaType
Object interface{}
}
type DeltaType string
const (
Added DeltaType = "Added"
Updated DeltaType = "Updated"
Deleted DeltaType = "Deleted"
Replaced DeltaType = "Replaced"
Sync DeltaType = "Sync"
)
Delta構造体は、DeltaType(文字列)と、そのDeltaが発生した具体的なオブジェクトを保存しています。
DeltaFIFOは主にキューとマップを内部で管理しており、直感的に以下のように表現できます:
DeltaFIFOのNew関数はNewDeltaFIFOWithOptions()です。
- client-go/tools/cache/delta_fifo.go:218
func NewDeltaFIFOWithOptions(opts DeltaFIFOOptions) *DeltaFIFO {
if opts.KeyFunction == nil {
opts.KeyFunction = MetaNamespaceKeyFunc
}
fifo := &DeltaFIFO{
items: map[string]Deltas{},
queue: []string{},
keyFunc: opts.KeyFunction,
knownObjects: opts.KnownObjects,
emitDeltaTypeReplaced: opts.EmitDeltaTypeReplaced,
}
fifo.cond.L = &fifo.lock
return fifo
}
要素の追加・削除・更新 - queueActionLocked()
DeltaFIFOのAdd()などのメソッドは、メソッド本体が非常に短いことがわかります。概ね以下のようになっています:
func (fifo *DeltaFIFO) Add(item interface{}) error {
fifo.lock.Lock()
defer fifo.lock.Unlock()
fifo.populated = true
return fifo.queueActionLocked(Added, item)
}
内部のロジックは、対応するDeltaTypeを渡してqueueActionLocked()メソッドを呼び出すだけです。前述したように、DeltaTypeはAdded、Updated、Deletedなどの文字列です。そこで、まずqueueActionLocked()メソッドの実装を見てみましょう。
- client-go/tools/cache/delta_fifo.go:409
func (fifo *DeltaFIFO) queueActionLocked(deltaType DeltaType, item interface{}) error {
id, err := fifo.KeyOf(item) // このオブジェクトのキーを計算する
if err != nil {
return KeyError{item, err}
}
oldDeltas := fifo.items[id] // itemsマップから現在のDeltasを取得する
newDeltas := append(oldDeltas, Delta{deltaType, item}) // Deltaを構築し、Deltasに追加する(つまり[]Deltaに追加する)
newDeltas = dedupDeltas(newDeltas) // もし最近のDeltaが重複している場合、後者のDeltaを保持する;現在のバージョンではDeletedの重複シナリオのみを処理している
if len(newDeltas) > 0 { // 理論的にはnewDeltasの長さは必ず0より大きい
if _, exists := fifo.items[id]; !exists {
fifo.queue = append(fifo.queue, id) // idが存在しない場合、キューに追加する
}
fifo.items[id] = newDeltas // idが既に存在する場合、itemsマップ内の対応するキーのDeltasのみを更新する
fifo.cond.Broadcast()
} else { // 理論的にはここは実行されない
if oldDeltas == nil {
klog.Errorf("id=%qのdedupDeltasが不可能です: oldDeltas=%#+v, item=%#+v; 無視します", id, oldDeltas, item)
return nil
}
klog.Errorf("id=%qのdedupDeltasが不可能です: oldDeltas=%#+v, item=%#+v; DeltaFIFOの不変条件を破り、空のDeltasを格納します", id, oldDeltas, item)
fifo.items[id] = newDeltas
return fmt.Errorf("id=%qのdedupDeltasが不可能です: oldDeltas=%#+v, item=%#+v; DeltaFIFOの不変条件を破り、空のDeltasを格納しました", id, oldDeltas, item)
}
return nil
}
ここまで来ると、Add()、Delete()、Update()、Get()などの関数は非常に明確になります。対応する変更タイプのobjをキューに追加するだけです。
Pop()
Popは、要素の追加または更新順序に従って1つの要素(Deltas)を順次返します。キューが空の場合はブロックされます。また、Popプロセスはキューから要素を1つ削除してから返すため、処理に失敗した場合はAddIfNotPresent()メソッドを通じてこの要素をキューに再追加する必要があります。
Popのパラメータはtype PopProcessFunc func(interface{}) error型のprocessで、Pop()関数内でキューの最初の要素をキューから削除し、processに渡して処理します。処理に失敗した場合は再入隊されますが、このDeltasと対応するエラー情報が返されます。
- client-go/tools/cache/delta_fifo.go:515
func (fifo *DeltaFIFO) Pop(processFunc PopProcessFunc) (interface{}, error) {
fifo.lock.Lock()
defer fifo.lock.Unlock()
for { // このループは実際には意味がなく、以下の!okと組み合わせて発生しない問題を解決するだけ
for len(fifo.queue) == 0 { // 空の場合、このループに入る
if fifo.closed { // キューが閉じられている場合は直接返す
return nil, ErrFIFOClosed
}
fifo.cond.Wait() // 待機
}
id := fifo.queue[0] // queueにはkeyが格納されている
fifo.queue = fifo.queue[1:] // queueからこのkeyを削除する
depth := len(fifo.queue)
if fifo.initialPopulationCount > 0 { // 初回のReplace()で挿入された要素数
fifo.initialPopulationCount--
}
item, ok := fifo.items[id] // items map[string]DeltasからDeltasを1つ取得する
if !ok { // 理論的には見つからないことはないため、上記のforの入れ子を導入しましたが、少し良くない感じがします
klog.Errorf("信じられません! %qはf.queueにありましたが、f.itemsにはありません; 無視します。", id)
continue
}
delete(fifo.items, id) // itemsマップからもこの要素を削除する
// キューの長さが10を超え、1つの要素の処理時間が0.1秒を超えるとログを記録する; キューの長さは理論的には変化しないはずです。なぜなら、要素を処理している間はブロックされるため、新しい要素が入ってこないからです
if depth > 10 {
trace := utiltrace.New("DeltaFIFO Pop Process",
utiltrace.Field{Key: "ID", Value: id},
utiltrace.Field{Key: "Depth", Value: depth},
utiltrace.Field{Key: "Reason", Value: "遅いイベントハンドラがキューをブロックしている"})
defer trace.LogIfLong(100 * time.Millisecond)
}
err := processFunc(item) // PopProcessFuncに渡して処理する
if e, ok := err.(ErrRequeue); ok { // 再入隊が必要な場合はキューに戻す
fifo.addIfNotPresent(id, item)
err = e.Err
}
// このDeltasとエラー情報を返す
return item, err
}
}
Pop()の実際の呼び出しシーンを見てみましょう:
- client-go/tools/cache/controller.go:181
func (c *controller) processLoop() {
for {
obj, err := c.config.Queue.Pop(PopProcessFunc(c.config.Process))
if err != nil {
if err == ErrFIFOClosed {
return
}
if c.config.RetryOnError {
c.config.Queue.AddIfNotPresent(obj) // 実際にはPop内部で既にAddIfNotPresentを呼んでいるため、ここも少し冗長かもしれません; より堅牢かもしれません
}
}
}
}
ここでまだ疑問があります。process関数はどのように実装されているのでしょうか?sharedIndexInformer内のprocess関数のロジックを見てみましょう(別の記事「Kubernetes client-go Informerソースコード分析」でこのメソッドを再度詳しく説明します):
- client-go/tools/cache/shared_informer.go:537
func (s *sharedIndexInformer) HandleDeltas(obj interface{}) error {
s.blockDeltas.Lock()
defer s.blockDeltas.Unlock()
// このループは古いものから新しいものへのプロセスです
for _, d := range obj.(Deltas) {
switch d.Type {
case Sync, Replaced, Added, Updated: // 以下のcaseはDeleted
s.cacheMutationDetector.AddObject(d.Object)
if old, exists, err := s.indexer.Get(d.Object); err == nil && exists {
// indexerを更新する
if err := s.indexer.Update(d.Object); err != nil {
return err
}
isSync := false
switch {
case d.Type == Sync:
isSync = true
case d.Type == Replaced:
if accessor, err := meta.Accessor(d.Object); err == nil {
if oldAccessor, err := meta.Accessor(old); err == nil {
isSync = accessor.GetResourceVersion() == oldAccessor.GetResourceVersion()
}
}
}
// 更新通知
s.processor.distribute(updateNotification{oldObj: old, newObj: d.Object}, isSync)
} else {
// objをindexerに追加する
if err := s.indexer.Add(d.Object); err != nil {
return err
}
// 追加通知
s.processor.distribute(addNotification{newObj: d.Object}, false)
}
case Deleted: // 削除の場合、indexerからobjを削除する
if err := s.indexer.Delete(d.Object); err != nil {
return err
}
// 削除メッセージを発行する
s.processor.distribute(deleteNotification{oldObj: d.Object}, false)
}
}
return nil
}
Replace()
Replace()は単純に2つのことを行います:
- 渡されたオブジェクトリストにSync/Replace DeltaTypeのDeltaを追加する
- その後、いくつかの削除ロジックを実行する
ここでのReplace()プロセスは、新しい[]Deltasを渡すと理解できます。もしこの要素が既にDeltaFIFO内に存在する場合、Sync/Replaceアクションが追加されます。逆に、DeltaFIFO内に余分なDeltasがある場合、それはapiserverとの接続が失われたため、実際には既に削除されているが、削除アクションがwatchされなかったオブジェクトである可能性があります。そのため、直接DeletedのDeltaが追加されます。
func (fifo *DeltaFIFO) Replace(items []interface{}, _ string) error {
fifo.lock.Lock()
defer fifo.lock.Unlock()
keys := make(sets.String, len(items)) // list内の各itemのkeyを保存するために使用
// 古いコードとの互換性ロジック
action := Sync
if fifo.emitDeltaTypeReplaced {
action = Replaced
}
for _, item := range items { // 各itemの後にSync/Replacedアクションを追加する
key, err := fifo.KeyOf(item)
if err != nil {
return KeyError{item, err}
}
keys.Insert(key)
if err := fifo.queueActionLocked(action, item); err != nil {
return fmt.Errorf("オブジェクトをキューに追加できませんでした: %v", err)
}
}
if fifo.knownObjects == nil {
queuedDeletions := 0
for k, oldItem := range fifo.items { // fifo.items内の古い要素を削除する
if keys.Has(k) {
continue
}
var deletedObj interface{}
if n := oldItem.Newest(); n != nil { // Deltasが空でない場合、戻り値がある
deletedObj = n.Object
}
queuedDeletions++
// 削除をマークする; apiserverとの接続が失われたため、削除状態がタイムリーに取得されなかったシナリオのため、ここではDeletedFinalStateUnknownタイプです
if err := fifo.queueActionLocked(Deleted, DeletedFinalStateUnknown{k, deletedObj}); err != nil {
return err
}
}
if !fifo.populated {
fifo.populated = true
fifo.initialPopulationCount = keys.Len() + queuedDeletions
}
return nil
}
knownKeys := fifo.knownObjects.ListKeys() // keyは例えば"default/pod_1"のような文字列です
queuedDeletions := 0
for _, k := range knownKeys {
if keys.Has(k) {
continue
}
// 新しいリストに存在しない古い要素を削除対象としてマークする
deletedObj, exists, err := fifo.knownObjects.GetByKey(k)
if err != nil {
deletedObj = nil
klog.Errorf("キー%vの検索中に予期せぬエラー%vが発生しました。オブジェクトなしでDeleteFinalStateUnknownマーカーを配置します", err, k)
} else if !exists {
deletedObj = nil
klog.Infof("キー%vは既知のオブジェクトストアに存在しません。オブジェクトなしでDeleteFinalStateUnknownマーカーを配置します", k)
}
queuedDeletions++
// 削除アクションを追加する; apiserverとの接続が失われたなどのシナリオにより、削除イベントがwatchされなかった場合、DeletedFinalStateUnknownタイプになります
if err := fifo.queueActionLocked(Deleted, DeletedFinalStateUnknown{k, deletedObj}); err != nil {
return err
}
}
if !fifo.populated {
fifo.populated = true
fifo.initialPopulationCount = keys.Len() + queuedDeletions
}
return nil
}
ここにはknownObjects属性があります。Replace()ロジックを完全に理解するには、knownObjectsが何であるかを理解する必要があります。
knownObjects属性の初期化を追跡すると、cache型で実装されたStoreが参照されていることがわかります。cacheはIndexerを実装するcacheであり、Indexerのソースコード分析は別の記事「Kubernetes client-go Indexer / ThreadSafeStoreソースコード分析」で見ることができます。
- client-go/tools/cache/store.go:258
func NewStore(keyFunc KeyFunc) Store {
return &cache{
cacheStorage: NewThreadSafeStore(Indexers{}, Indices{}),
keyFunc: keyFunc,
}
}
ここではStoreとして使用されており、Indexerではありません。NewStore()関数を呼び出す際に渡されるパラメータは以下の通りです:
clientState := NewStore(DeletionHandlingMetaNamespaceKeyFunc)
// DeletedFinalStateUnknownオブジェクトのキー取得問題を処理する
func DeletionHandlingMetaNamespaceKeyFunc(obj interface{}) (string, error) {
if d, ok := obj.(DeletedFinalStateUnknown); ok {
return d.Key, nil
}
return MetaNamespaceKeyFunc(obj)
}
したがって、knownObjectsはcache型のインスタンスを通じて、Indexerと同様のメカニズムを使用し、内部のThreadSafeStoreを通じてキュー内のすべての要素のキーを検索する機能を実現しています。
DeltaFIFOとIndexerの間にはInformerというブリッジがあります。ここではsharedIndexInformerの*HandleDeltas()*メソッドを簡単に触れましたが、後でInformerのロジックを詳しく分析し、最終的にカスタムコントローラーとclient-go関連コンポーネントのロジックを一つにまとめます。