· yuta

RDBMSを自作して学ぶ:第4章 バッファプールマネージャ

#database#rdbms#go

前回はディスクとのページ単位I/Oを担うDiskManagerを実装しました。しかしDiskManagerを直接呼び出すだけでは、同じページに何度アクセスしても毎回ディスクI/Oが発生してしまいます。今回はディスクアクセスを減らすメモリキャッシュ層であるバッファプールマネージャを実装します。

なぜバッファプールが必要か

ディスクI/Oのコスト

第3章で見たとおり、メモリとディスクのアクセス速度には大きな差があります。

特性 メモリ(RAM) ディスク(SSD)
アクセス速度 〜100ns 〜100µs
速度差 基準 約1,000倍遅い

たとえば同じページに10回アクセスするクエリがあったとします。毎回DiskManagerを直接呼び出すと10回分のディスクI/Oが発生しますが、一度読み込んだページをメモリ上に保持しておけば、2回目以降はメモリアクセスだけで済みます。この「一度読んだページをメモリにキャッシュしておく」役割を担うのがバッファプールマネージャです。

バッファプールの基本的な仕組み

バッファプールは固定数の**フレーム(Frame)**というスロットを持つメモリ領域です。ページはディスクから読み込まれるとフレームにコピーされ、以降そのページへのアクセスはフレーム経由で行われます。

                     BufferPoolManager
                  ┌─────────────────────┐
                  │  Frame 0: Page 3     │
上位層 ──fetch──▶ │  Frame 1: Page 7     │──write/read──▶ Disk
                  │  Frame 2: (空き)      │
                  └─────────────────────┘

フレーム数はページ総数より少ないのが普通です。フレームが埋まった状態で新しいページを読み込むには、どれかのページをフレームから追い出す(evictする)必要があります。どのページを追い出すかを決める処理を**追い出しポリシー(Replacement Policy)**と呼びます。

バッファプールの設計

管理すべき情報

1つのフレームには次の情報を持たせます。

フィールド 役割
page ページの実データ(4096バイト)
pageID このフレームが保持しているページのID
pinCount このページを現在利用中のクライアント数
isDirty ページが読み込み後に変更されたかどうか

pinCount(ピンカウント)の役割

pinCount はバッファプール設計の要です。上位層(B+木やHeapFileなど)がページを利用している間、そのページがフレームから追い出されると困ります。そこでページを使い始めるときに pinCount をインクリメントし(ピンする)、使い終わったらデクリメントします(アンピンする)。追い出しポリシーは pinCount == 0 のフレーム、つまり誰にも使われていないフレームだけを対象にします。

isDirty(ダーティフラグ)の役割

読み込んだだけで変更していないページを追い出すときは、ディスクに書き戻す必要はありません。ページの内容を変更した場合だけ isDirty を立てておき、追い出し時にダーティなページだけをディスクに書き戻すことで、不要なディスクI/Oを減らします。

ページテーブル

フレームがどのページを保持しているかを高速に調べるため、PageID → FrameID のマップ(ページテーブル)を持ちます。これによりページがすでにキャッシュされているかどうかを O(1) で判定できます。

LRU(Least Recently Used)による追い出しポリシー

考え方

追い出しポリシーには様々な方式がありますが、本シリーズでは実装がシンプルで効果も高い**LRU(Least Recently Used)**を採用します。LRUは「最も長い間使われていないページ」を追い出し対象に選ぶ方式です。直感的には、最近使われたページは近い将来も使われる可能性が高いという経験則(局所性)に基づいています。

アンピンされた順(古い ← → 新しい)
┌───────┬───────┬───────┬───────┐
│ Page5 │ Page2 │ Page9 │ Page1 │
└───────┴───────┴───────┴───────┘

   ここが次の追い出し候補(Victim)

ピン中のページはLRUの対象外

pinCount > 0 のフレームは利用中なのでLRUの追い出し候補には含めません。あるページがアンピンされた(pinCount が0になった)時点でLRUのリストに加え、再びピンされたらリストから除外します。

実装:Replacer

インターフェースの定義

追い出しポリシーを差し替え可能にするため、インターフェースとして定義します。将来LFUなど別の方式を試す場合もBufferPoolManager側の変更は不要です。

// storage/buffer/replacer.go
package buffer

// FrameID はバッファプール内のフレームを識別するID(0始まり)
type FrameID uint32

// Replacer は追い出し対象のフレームを選ぶポリシーのインターフェース
type Replacer interface {
    // Victim は追い出すフレームを選び、Replacerの管理から取り除いて返す
    // 追い出し候補が1つもない場合は false を返す
    Victim() (FrameID, bool)
    // Pin はフレームがピンされたことを記録し、追い出し候補から除外する
    Pin(id FrameID)
    // Unpin はフレームがアンピンされたことを記録し、追い出し候補に加える
    Unpin(id FrameID)
    // Size は現在追い出し候補になっているフレーム数を返す
    Size() int
}

LRUReplacer の実装

Goの container/list(双方向連結リスト)と map を組み合わせて実装します。リストの先頭を「最も長く使われていない」、末尾を「最も最近使われた」とし、Unpin でリスト末尾に追加、Victim で先頭から取り出します。

// storage/buffer/lru_replacer.go
package buffer

import (
    "container/list"
    "sync"
)

// LRUReplacer はLeast Recently Usedポリシーで追い出し対象を選ぶReplacer実装
type LRUReplacer struct {
    mu       sync.Mutex
    list     *list.List // 先頭=最も長く使われていない, 末尾=最も最近使われた
    elements map[FrameID]*list.Element
}

// NewLRUReplacer は空のLRUReplacerを返す
func NewLRUReplacer() *LRUReplacer {
    return &LRUReplacer{
        list:     list.New(),
        elements: make(map[FrameID]*list.Element),
    }
}

// Victim は最も長く使われていないフレームを取り出す
func (r *LRUReplacer) Victim() (FrameID, bool) {
    r.mu.Lock()
    defer r.mu.Unlock()

    front := r.list.Front()
    if front == nil {
        return 0, false
    }
    id := front.Value.(FrameID)
    r.list.Remove(front)
    delete(r.elements, id)
    return id, true
}

// Pin はフレームを追い出し候補から除外する
func (r *LRUReplacer) Pin(id FrameID) {
    r.mu.Lock()
    defer r.mu.Unlock()

    if elem, ok := r.elements[id]; ok {
        r.list.Remove(elem)
        delete(r.elements, id)
    }
}

// Unpin はフレームを追い出し候補(リスト末尾)に加える
func (r *LRUReplacer) Unpin(id FrameID) {
    r.mu.Lock()
    defer r.mu.Unlock()

    if _, ok := r.elements[id]; ok {
        return // 既に追い出し候補に入っている
    }
    r.elements[id] = r.list.PushBack(id)
}

// Size は現在追い出し候補になっているフレーム数を返す
func (r *LRUReplacer) Size() int {
    r.mu.Lock()
    defer r.mu.Unlock()
    return r.list.Len()
}

実装:BufferPoolManager

構造体の定義

// storage/buffer/buffer_pool_manager.go
package buffer

import (
    "fmt"
    "sync"

    "github.com/yourname/go-rdbms/storage/disk"
    "github.com/yourname/go-rdbms/storage/page"
)

// frame はバッファプール内の1スロットを表す
type frame struct {
    page     page.Page
    pageID   page.PageID
    pinCount int
    isDirty  bool
}

// BufferPoolManager はページのメモリキャッシュを管理する
type BufferPoolManager struct {
    mu          sync.Mutex
    diskManager disk.DiskManager
    frames      []frame
    pageTable   map[page.PageID]FrameID
    freeList    []FrameID
    replacer    Replacer
}

// NewBufferPoolManager は指定サイズのフレームを持つBufferPoolManagerを返す
func NewBufferPoolManager(diskManager disk.DiskManager, poolSize int) *BufferPoolManager {
    freeList := make([]FrameID, poolSize)
    for i := range freeList {
        freeList[i] = FrameID(i)
    }

    return &BufferPoolManager{
        diskManager: diskManager,
        frames:      make([]frame, poolSize),
        pageTable:   make(map[page.PageID]FrameID),
        freeList:    freeList,
        replacer:    NewLRUReplacer(),
    }
}

空きフレームのリスト(freeList)を別途持つことで、バッファプールがまだ埋まっていない間は追い出し処理を経由せずにフレームを割り当てられるようにしています。

FetchPage:既存ページの取得

// FetchPage はpageIDのページをバッファプール上で返す
// すでにキャッシュされていればそれを、なければディスクから読み込む
func (bpm *BufferPoolManager) FetchPage(pageID page.PageID) (*page.Page, error) {
    bpm.mu.Lock()
    defer bpm.mu.Unlock()

    if frameID, ok := bpm.pageTable[pageID]; ok {
        bpm.frames[frameID].pinCount++
        bpm.replacer.Pin(frameID)
        return &bpm.frames[frameID].page, nil
    }

    frameID, err := bpm.allocateFrame()
    if err != nil {
        return nil, err
    }

    var p page.Page
    if err := bpm.diskManager.ReadPage(pageID, &p); err != nil {
        return nil, fmt.Errorf("buffer: ページ %d の読み込みに失敗しました: %w", pageID, err)
    }

    f := &bpm.frames[frameID]
    f.page = p
    f.pageID = pageID
    f.pinCount = 1
    f.isDirty = false
    bpm.pageTable[pageID] = frameID
    bpm.replacer.Pin(frameID)

    return &f.page, nil
}

キャッシュヒット時はディスクI/Oなしで即座にページを返し、ミス時のみ allocateFrame でフレームを確保してディスクから読み込みます。

NewPage:新規ページの割り当て

// NewPage はディスク上に新しいページを割り当て、バッファプールに読み込んだ状態で返す
func (bpm *BufferPoolManager) NewPage() (page.PageID, *page.Page, error) {
    bpm.mu.Lock()
    defer bpm.mu.Unlock()

    pageID, err := bpm.diskManager.AllocatePage()
    if err != nil {
        return page.InvalidPageID, nil, fmt.Errorf("buffer: 新規ページの割り当てに失敗しました: %w", err)
    }

    frameID, err := bpm.allocateFrame()
    if err != nil {
        return page.InvalidPageID, nil, err
    }

    f := &bpm.frames[frameID]
    f.page = page.Page{}
    f.page.SetPageID(pageID)
    f.pageID = pageID
    f.pinCount = 1
    f.isDirty = true
    bpm.pageTable[pageID] = frameID
    bpm.replacer.Pin(frameID)

    return pageID, &f.page, nil
}

新規ページは中身が空でも「ディスク上に確保された」時点でダーティ(未反映)として扱い、isDirty = true にしています。

UnpinPage:ページの利用終了

// UnpinPage はページの利用終了を通知する
// isDirtyがtrueの場合、ページが変更されたものとしてマークする
func (bpm *BufferPoolManager) UnpinPage(pageID page.PageID, isDirty bool) error {
    bpm.mu.Lock()
    defer bpm.mu.Unlock()

    frameID, ok := bpm.pageTable[pageID]
    if !ok {
        return fmt.Errorf("buffer: ページ %d はバッファプールにありません", pageID)
    }

    f := &bpm.frames[frameID]
    if f.pinCount <= 0 {
        return fmt.Errorf("buffer: ページ %d は既にアンピンされています", pageID)
    }

    if isDirty {
        f.isDirty = true
    }
    f.pinCount--
    if f.pinCount == 0 {
        bpm.replacer.Unpin(frameID)
    }
    return nil
}

isDirty は「trueなら立てる」だけで、falseだからといって既存のダーティフラグを下ろさない点に注意してください。あるクライアントは読むだけ、別のクライアントは書き込む、といった複数回の FetchPage/UnpinPage が重なっても、一度でも変更されたページは正しくダーティとして扱われます。

FlushPage:ディスクへの書き戻し

// FlushPage は指定ページをディスクに書き戻す
func (bpm *BufferPoolManager) FlushPage(pageID page.PageID) error {
    bpm.mu.Lock()
    defer bpm.mu.Unlock()
    return bpm.flushFrame(pageID)
}

func (bpm *BufferPoolManager) flushFrame(pageID page.PageID) error {
    frameID, ok := bpm.pageTable[pageID]
    if !ok {
        return fmt.Errorf("buffer: ページ %d はバッファプールにありません", pageID)
    }

    f := &bpm.frames[frameID]
    if !f.isDirty {
        return nil
    }
    if err := bpm.diskManager.WritePage(pageID, &f.page); err != nil {
        return fmt.Errorf("buffer: ページ %d の書き戻しに失敗しました: %w", pageID, err)
    }
    f.isDirty = false
    return nil
}

allocateFrame:フレームの確保

空きフレームがあればそれを使い、なければLRUに従って追い出し候補を選び、ダーティなら書き戻してから再利用します。

// allocateFrame は空きフレームを確保する
// 空きがなければreplacerに追い出しを依頼する
func (bpm *BufferPoolManager) allocateFrame() (FrameID, error) {
    if len(bpm.freeList) > 0 {
        frameID := bpm.freeList[len(bpm.freeList)-1]
        bpm.freeList = bpm.freeList[:len(bpm.freeList)-1]
        return frameID, nil
    }

    victimID, ok := bpm.replacer.Victim()
    if !ok {
        return 0, fmt.Errorf("buffer: 追い出し可能なページがありません(全ページがピン中です)")
    }

    victim := &bpm.frames[victimID]
    if victim.isDirty {
        if err := bpm.diskManager.WritePage(victim.pageID, &victim.page); err != nil {
            return 0, fmt.Errorf("buffer: 追い出し対象ページ %d の書き戻しに失敗しました: %w", victim.pageID, err)
        }
    }
    delete(bpm.pageTable, victim.pageID)

    return victimID, nil
}

全フレームがピンされている状態で新しいページを要求すると、追い出せるフレームが1つもないためエラーを返します。これはバッファプールサイズが小さすぎるか、上位層がページをアンピンし忘れていることを示すシグナルです。

テストを書く

// storage/buffer/buffer_pool_manager_test.go
package buffer_test

import (
    "os"
    "testing"

    "github.com/yourname/go-rdbms/storage/buffer"
    "github.com/yourname/go-rdbms/storage/disk"
)

// setupBufferPoolManager はテスト用の一時ファイルでBufferPoolManagerを作成する
func setupBufferPoolManager(t *testing.T, poolSize int) (*buffer.BufferPoolManager, func()) {
    t.Helper()

    f, err := os.CreateTemp("", "test-buffer-*.db")
    if err != nil {
        t.Fatal(err)
    }
    f.Close()

    dm, err := disk.NewFileDiskManager(f.Name())
    if err != nil {
        os.Remove(f.Name())
        t.Fatal(err)
    }

    bpm := buffer.NewBufferPoolManager(dm, poolSize)
    cleanup := func() {
        dm.Close()
        os.Remove(f.Name())
    }
    return bpm, cleanup
}

func TestBufferPoolManager_NewPageAndFetch(t *testing.T) {
    bpm, cleanup := setupBufferPoolManager(t, 3)
    defer cleanup()

    pageID, p, err := bpm.NewPage()
    if err != nil {
        t.Fatal(err)
    }
    copy(p.Data(), []byte("hello, buffer pool!"))
    if err := bpm.UnpinPage(pageID, true); err != nil {
        t.Fatal(err)
    }

    fetched, err := bpm.FetchPage(pageID)
    if err != nil {
        t.Fatal(err)
    }
    if string(fetched.Data()[:19]) != "hello, buffer pool!" {
        t.Errorf("データ不一致: got %q", fetched.Data()[:19])
    }
    bpm.UnpinPage(pageID, false)
}

func TestBufferPoolManager_EvictsLeastRecentlyUsed(t *testing.T) {
    bpm, cleanup := setupBufferPoolManager(t, 2)
    defer cleanup()

    id1, p1, _ := bpm.NewPage()
    copy(p1.Data(), []byte("page1"))
    bpm.UnpinPage(id1, true)

    id2, p2, _ := bpm.NewPage()
    copy(p2.Data(), []byte("page2"))
    bpm.UnpinPage(id2, true)

    // id1を使うことで最近使った扱いにする(id2がLRUの先頭に残る)
    bpm.FetchPage(id1)
    bpm.UnpinPage(id1, false)

    // 3ページ目を確保すると、直近アクセスのないid2が追い出される
    id3, p3, _ := bpm.NewPage()
    copy(p3.Data(), []byte("page3"))
    bpm.UnpinPage(id3, true)

    // id1は引き続きバッファプール上にあるはず
    if _, err := bpm.FetchPage(id1); err != nil {
        t.Fatalf("id1はバッファプールに残っているはず: %v", err)
    }
    bpm.UnpinPage(id1, false)

    // id2はディスクから再読込されるが、追い出し時に書き戻されているのでデータは保持されている
    fetched2, err := bpm.FetchPage(id2)
    if err != nil {
        t.Fatal(err)
    }
    if string(fetched2.Data()[:5]) != "page2" {
        t.Errorf("id2のデータが失われている: got %q", fetched2.Data()[:5])
    }
    bpm.UnpinPage(id2, false)
}

func TestBufferPoolManager_AllPinnedReturnsError(t *testing.T) {
    bpm, cleanup := setupBufferPoolManager(t, 1)
    defer cleanup()

    if _, _, err := bpm.NewPage(); err != nil {
        t.Fatal(err)
    }

    // プールサイズ1で、既存ページをアンピンせずに新規ページを確保しようとするとエラー
    if _, _, err := bpm.NewPage(); err == nil {
        t.Error("全フレームがピン中なのにエラーにならなかった")
    }
}

func TestBufferPoolManager_FlushPage(t *testing.T) {
    bpm, cleanup := setupBufferPoolManager(t, 2)
    defer cleanup()

    pageID, p, _ := bpm.NewPage()
    copy(p.Data(), []byte("flush test"))
    bpm.UnpinPage(pageID, true)

    if err := bpm.FlushPage(pageID); err != nil {
        t.Fatal(err)
    }
}
go test ./storage/buffer/... -v

期待される出力:

=== RUN   TestBufferPoolManager_NewPageAndFetch
--- PASS: TestBufferPoolManager_NewPageAndFetch (0.00s)
=== RUN   TestBufferPoolManager_EvictsLeastRecentlyUsed
--- PASS: TestBufferPoolManager_EvictsLeastRecentlyUsed (0.00s)
=== RUN   TestBufferPoolManager_AllPinnedReturnsError
--- PASS: TestBufferPoolManager_AllPinnedReturnsError (0.00s)
=== RUN   TestBufferPoolManager_FlushPage
--- PASS: TestBufferPoolManager_FlushPage (0.00s)
PASS
ok      github.com/yourname/go-rdbms/storage/buffer   0.010s

ここまでのコード全体像

本章が完了した時点のパッケージ構成です。

storage/
├── page/
│   └── page.go                    ← Page型・PageID型・ページヘッダ操作
├── disk/
│   ├── disk_manager.go            ← DiskManagerインターフェース
│   └── file_disk_manager.go       ← ファイルベースの実装
└── buffer/
    ├── replacer.go                ← Replacerインターフェース
    ├── lru_replacer.go            ← LRUによる追い出しポリシー実装
    └── buffer_pool_manager.go     ← ページキャッシュ本体

BufferPoolManager の設計で意識したポイント

① ページテーブルでO(1)検索する

map[page.PageID]FrameID を使うことで、あるページがキャッシュ済みかどうかを定数時間で判定できます。線形探索にしてしまうとフレーム数が増えるほど FetchPage が遅くなり、キャッシュを持つ意味が薄れてしまいます。

② pinCountで安全な追い出しを保証する

上位層が利用中のページは pinCount > 0 である限りLRUの対象外です。これにより、あるゴルーチンがページを読み書きしている最中に、別のゴルーチンの FetchPage 呼び出しによってそのページが追い出されてしまう事故を防いでいます。

③ isDirtyで不要な書き込みを避ける

読み込んだだけで変更していないページは追い出し時にディスクへ書き戻す必要がありません。isDirty フラグでこれを区別することで、無駄なディスクI/Oを削減しています。

④ Replacerをインターフェースとして分離する

追い出しポリシーをインターフェースにしたことで、BufferPoolManager 自体はLRUの実装詳細を知りません。将来Clock-SweepやLFUなど別のアルゴリズムを試す際も、Replacer を実装した新しい型を渡すだけで済みます。

まとめ

次章では、このバッファプールの上にテーブルの実データを格納するHeapFile / SlottedPageを実装し、可変長のレコードをページ内でどう管理するかを見ていきます。