From aba0093d3077ed78734bf6af63ed1de29cb0f71c Mon Sep 17 00:00:00 2001 From: Andrew Batz Date: Fri, 24 Jul 2026 12:41:55 -0700 Subject: [PATCH] fix(scorch): serialize merge task execution Execute merge plan tasks sequentially, matching upstream behavior. Concurrent tasks multiply transient disk staging because every in-flight output remains live and ineligible for garbage collection until the whole plan is introduced. Serial execution bounds staging to one task at a time and allows cancellation to take effect between tasks. Merge planner size and task-count defaults remain unchanged. --- index/scorch/merge.go | 28 ++++++++++++---------------- 1 file changed, 12 insertions(+), 16 deletions(-) diff --git a/index/scorch/merge.go b/index/scorch/merge.go index 2038644a9..ee77ffcba 100644 --- a/index/scorch/merge.go +++ b/index/scorch/merge.go @@ -329,38 +329,34 @@ func (s *Scorch) planMergeAtSnapshot(ctx context.Context, return nil } - // execute all merge tasks concurrently - var wg sync.WaitGroup + // Execute merge tasks SEQUENTIALLY (as upstream does). Concurrent tasks + // multiply transient on-disk staging: every in-flight task's output files + // are live (ineligible for GC) until the whole plan lands, and observed + // zapx merge output can exceed the planner's predicted sizes by well over + // an order of magnitude. Serial execution bounds staging to one task's + // footprint at a time and lets a cancel take effect between tasks. results := make([]*mergeTaskResult, len(tasks)) var mergeErrs []error - var errMu sync.Mutex var allFilenames []string - var filenamesMu sync.Mutex for taskIdx, task := range tasks { atomic.AddUint64(&s.stats.TotFileMergePlanTasksSegments, uint64(len(task.Segments))) - wg.Add(1) - go func(taskIdx int, task *mergeplan.MergeTask) { - defer wg.Done() - + select { + case <-cancelCh: + mergeErrs = append(mergeErrs, segment.ErrClosed) + default: result, err := s.executeMergeTask(task, cancelCh, taskIdx) if err != nil { - errMu.Lock() mergeErrs = append(mergeErrs, err) - errMu.Unlock() - return + continue } - if result != nil { results[taskIdx] = result - filenamesMu.Lock() allFilenames = append(allFilenames, result.filenames...) - filenamesMu.Unlock() } - }(taskIdx, task) + } } - wg.Wait() // check for errors from merge workers if len(mergeErrs) > 0 {