Kubernetes client-go DeltaFIFOのソースコード分析

概要 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つのことを行います:

  1. 渡されたオブジェクトリストにSync/Replace DeltaTypeのDeltaを追加する
  2. その後、いくつかの削除ロジックを実行する

ここでの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関連コンポーネントのロジックを一つにまとめます。

タグ: Kubernetes client-go deltafifo Go

7月23日 16:47 投稿