From c71accf511d5c68e2502c2e623f5a5511928317a Mon Sep 17 00:00:00 2001 From: AJ Roetker Date: Fri, 24 Jul 2026 15:07:12 -0700 Subject: [PATCH] fix(scorch): cap aggregate merge plan input bytes --- index/scorch/merge_test.go | 21 +++++++++++- index/scorch/mergeplan/merge_plan.go | 28 +++++++++++++--- index/scorch/mergeplan/merge_plan_test.go | 41 +++++++++++++++++++++++ 3 files changed, 85 insertions(+), 5 deletions(-) diff --git a/index/scorch/merge_test.go b/index/scorch/merge_test.go index 682e7efb1..21b7b6f7a 100644 --- a/index/scorch/merge_test.go +++ b/index/scorch/merge_test.go @@ -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) @@ -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}, diff --git a/index/scorch/mergeplan/merge_plan.go b/index/scorch/mergeplan/merge_plan.go index 0f9539b38..0e025a035 100644 --- a/index/scorch/mergeplan/merge_plan.go +++ b/index/scorch/mergeplan/merge_plan.go @@ -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 @@ -238,6 +247,7 @@ func plan(segmentsIn []Segment, o *MergePlanOptions) (*MergePlan, error) { } rv := &MergePlan{} + var mergePlanInputSize int64 var empties []Segment for _, eligible := range eligibles { @@ -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 { @@ -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) diff --git a/index/scorch/mergeplan/merge_plan_test.go b/index/scorch/mergeplan/merge_plan_test.go index 1f272d579..65e5eb24b 100644 --- a/index/scorch/mergeplan/merge_plan_test.go +++ b/index/scorch/mergeplan/merge_plan_test.go @@ -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