Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 20 additions & 1 deletion index/scorch/merge_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,25 @@ import (
index "github.com/blevesearch/bleve_index_api"
)

func TestParseMergePlannerOptionsMaxMergePlanInputSize(t *testing.T) {
s := &Scorch{
config: map[string]interface{}{
"scorchMergePlanOptions": map[string]interface{}{
"maxMergePlanInputSize": int64(256 << 20),
},
},
}

options, err := s.parseMergePlannerOptions()
if err != nil {
t.Fatal(err)
}
if options.MaxMergePlanInputSize != 256<<20 {
t.Fatalf("expected max merge plan input size %d, got %d",
256<<20, options.MaxMergePlanInputSize)
}
}

func TestObsoleteSegmentMergeIntroduction(t *testing.T) {
testConfig := CreateConfig("TestObsoleteSegmentMergeIntroduction")
err := InitTest(testConfig)
Expand Down Expand Up @@ -225,7 +244,7 @@ func setupBenchIndex(b *testing.B, name string, numBatches, docsPerBatch int) (*

func BenchmarkForceMerge(b *testing.B) {
for _, tc := range []struct {
batches int
batches int
docsPerBatch int
}{
{10, 100},
Expand Down
28 changes: 24 additions & 4 deletions index/scorch/mergeplan/merge_plan.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,15 @@ type MergePlanOptions struct {
// contain vectors that are too large.
MaxSegmentFileSize int64

// Max summed size (in bytes) of persisted segment inputs scheduled
// across the output-producing tasks in a single merge plan. A value
// <= 0 disables the limit.
//
// Merge executors may retain every task's output until the whole plan
// is introduced. Limiting each task independently does not bound that
// aggregate staging footprint, so this limit is enforced plan-wide.
MaxMergePlanInputSize int64

// The growth factor for each tier in a staircase of idealized
// segments computed by CalcBudget().
TierGrowth float64
Expand Down Expand Up @@ -238,6 +247,7 @@ func plan(segmentsIn []Segment, o *MergePlanOptions) (*MergePlan, error) {
}

rv := &MergePlan{}
var mergePlanInputSize int64

var empties []Segment
for _, eligible := range eligibles {
Expand Down Expand Up @@ -265,25 +275,32 @@ func plan(segmentsIn []Segment, o *MergePlanOptions) (*MergePlan, error) {
for startIdx := 0; startIdx < len(eligibles); startIdx++ {
roster := rosterBuf[:0]
var rosterLiveSize int64
var rosterFileSize int64 // useful for segments with vectors
var rosterFileSize int64
var rosterVectorFileSize int64

for idx := startIdx; idx < len(eligibles) && len(roster) < o.SegmentsPerMergeTask; idx++ {
eligible := eligibles[idx]
eligibleFileSize := eligible.FileSize()

if rosterLiveSize+eligible.LiveSize() >= o.MaxSegmentSize {
continue
}

if eligible.HasVector() {
efs := eligible.FileSize()
if rosterFileSize+efs >= o.MaxSegmentFileSize {
if rosterVectorFileSize+eligibleFileSize >= o.MaxSegmentFileSize {
continue
}
rosterFileSize += efs
rosterVectorFileSize += eligibleFileSize
}

if o.MaxMergePlanInputSize > 0 &&
eligibleFileSize > o.MaxMergePlanInputSize-mergePlanInputSize-rosterFileSize {
continue
}

roster = append(roster, eligible)
rosterLiveSize += eligible.LiveSize()
rosterFileSize += eligibleFileSize
}

if len(roster) > 0 {
Expand All @@ -303,6 +320,9 @@ func plan(segmentsIn []Segment, o *MergePlanOptions) (*MergePlan, error) {
// create tasks with valid merges - i.e. there should be at least 2 non-empty segments
if len(bestRoster) > 1 {
rv.Tasks = append(rv.Tasks, &MergeTask{Segments: bestRoster})
for _, segment := range bestRoster {
mergePlanInputSize += segment.FileSize()
}
}

eligibles = removeSegments(eligibles, bestRoster)
Expand Down
41 changes: 41 additions & 0 deletions index/scorch/mergeplan/merge_plan_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -762,6 +762,47 @@ func TestPlanMaxSegmentFileSize(t *testing.T) {
}
}

func TestPlanMaxMergePlanInputSizeNonVector(t *testing.T) {
segments := []Segment{
&segment{MyId: 1, MyFullSize: 100, MyLiveSize: 100, MyFileSize: 40},
&segment{MyId: 2, MyFullSize: 90, MyLiveSize: 90, MyFileSize: 40},
&segment{MyId: 3, MyFullSize: 80, MyLiveSize: 80, MyFileSize: 40},
&segment{MyId: 4, MyFullSize: 70, MyLiveSize: 70, MyFileSize: 40},
&segment{MyId: 5, MyFullSize: 60, MyLiveSize: 60, MyFileSize: 40},
&segment{MyId: 6, MyFullSize: 50, MyLiveSize: 50, MyFileSize: 40},
}
options := &MergePlanOptions{
MaxSegmentSize: 10_000,
MaxSegmentsPerTier: 1,
SegmentsPerMergeTask: 2,
TierGrowth: 2,
FloorSegmentSize: 1,
MaxMergePlanInputSize: 100,
}

plan, err := Plan(segments, options)
if err != nil {
t.Fatal(err)
}
if len(plan.Tasks) != 1 {
t.Fatalf("expected one capped task, got %d", len(plan.Tasks))
}

var inputSize int64
for _, task := range plan.Tasks {
for _, segment := range task.Segments {
if segment.HasVector() {
t.Fatal("test requires ordinary non-vector segments")
}
inputSize += segment.FileSize()
}
}
if inputSize > options.MaxMergePlanInputSize {
t.Fatalf("planned input size %d exceeds cap %d",
inputSize, options.MaxMergePlanInputSize)
}
}

func TestSingleTaskMergePlan(t *testing.T) {
o := &DefaultMergePlanOptions
o.FloorSegmentFileSize = 209715200
Expand Down
Loading