-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathcompaction.go
More file actions
150 lines (125 loc) · 3.68 KB
/
Copy pathcompaction.go
File metadata and controls
150 lines (125 loc) · 3.68 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
package stratago
import (
"fmt"
"os"
"path/filepath"
"time"
"github.com/thomazdavis/stratago/sstable"
)
const CompactionThreshold = 4
// compactionWorker runs in the background and checks for compaction opportunities
func (db *StrataGo) compactionWorker() {
defer db.wg.Done()
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for {
select {
case <-ticker.C:
if err := db.RunCompaction(); err != nil {
fmt.Printf("Compaction failed: %v\n", err)
}
case <-db.closeChan:
return
}
}
}
// getTier returns a bucket index based on file size
func getTier(size int64) int {
mb := int64(1024 * 1024)
if size < 10*mb {
return 0 // Tier 0: < 10MB
} else if size < 50*mb {
return 1 // Tier 1: 10MB - 50MB
} else if size < 250*mb {
return 2 // Tier 2: 50MB - 250MB
} else if size < 1024*mb {
return 3 // Tier 3: 250MB - 1GB
}
return 4 // Tier 4: > 1GB
}
// RunCompaction executes a Size-Tiered compaction job
func (db *StrataGo) RunCompaction() error {
// Select the files
filesToCompact, startIndex, currentTier := db.selectFilesForCompaction()
// If we didn't find any valid group, abort gracefully
if len(filesToCompact) == 0 {
return nil
}
fmt.Printf("Starting compaction on Tier %d (Merging %d files)...\n", currentTier, len(filesToCompact))
// Iterators (Newest to Oldest)
var iters []*sstable.Iterator
for i := len(filesToCompact) - 1; i >= 0; i-- {
iter, err := filesToCompact[i].NewIterator()
if err != nil {
return fmt.Errorf("failed to create iterator: %w", err)
}
iters = append(iters, iter)
}
var newestTimestamp int64
fmt.Sscanf(filepath.Base(filesToCompact[len(filesToCompact)-1].Path()), "data_%d.sst", &newestTimestamp)
mergedSSTName := fmt.Sprintf("data_%d.sst", newestTimestamp)
mergedSSTPath := filepath.Join(db.dataDir, mergedSSTName)
builder, err := sstable.NewBuilder(mergedSSTPath)
if err != nil {
return err
}
if err := sstable.Merge(iters, builder); err != nil {
return fmt.Errorf("merge failed: %w", err)
}
newReader, err := sstable.NewReader(mergedSSTPath)
if err != nil {
return err
}
db.mu.Lock()
newReaders := make([]*sstable.Reader, 0, len(db.sstReaders)-CompactionThreshold+1)
newReaders = append(newReaders, db.sstReaders[:startIndex]...)
newReaders = append(newReaders, newReader)
newReaders = append(newReaders, db.sstReaders[startIndex+CompactionThreshold:]...)
db.sstReaders = newReaders
db.mu.Unlock()
// Delete the old files from disk
for _, r := range filesToCompact {
oldPath := r.Path()
r.Close()
os.Remove(oldPath)
}
for _, it := range iters {
it.Close()
}
fmt.Println("Compaction complete!")
return nil
}
// selectFilesForCompaction scans the current SSTables and finds a contiguous group
// of files in the same size tier. Returns the files, their starting index, and the tier
func (db *StrataGo) selectFilesForCompaction() ([]*sstable.Reader, int, int) {
db.mu.RLock()
defer db.mu.RUnlock()
var filesToCompact []*sstable.Reader
var startIndex int
currentTier := -1
var currentGroup []*sstable.Reader
var groupStartIndex int
for i, r := range db.sstReaders {
stat, err := os.Stat(r.Path())
if err != nil {
continue
}
tier := getTier(stat.Size())
if tier == currentTier {
currentGroup = append(currentGroup, r)
// If we hit our threshold of contiguous files in the same tier
if len(currentGroup) == CompactionThreshold {
filesToCompact = currentGroup
startIndex = groupStartIndex
return filesToCompact, startIndex, currentTier
}
} else {
// Reset the group because the tier changed
currentTier = tier
currentGroup = []*sstable.Reader{r}
groupStartIndex = i
}
}
// Return nil if no valid group was found
return nil, 0, -1
}