Skip to content

Commit 525b363

Browse files
committed
code cleanup
1 parent 6492578 commit 525b363

2 files changed

Lines changed: 5 additions & 96 deletions

File tree

index/scorch/introducer.go

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -375,7 +375,7 @@ func (s *Scorch) introduceMerge(nextMerge *segmentMerge) {
375375
deletedSinceItr := deletedSince.Iterator()
376376
for deletedSinceItr.HasNext() {
377377
oldDocNum := deletedSinceItr.Next()
378-
newDocNum := nextMerge.oldNewDocNums[segmentID][oldDocNum]
378+
newDocNum := segSnapAtMerge.oldNewDocIDs[oldDocNum]
379379
newSegmentDeleted[segSnapAtMerge.workerID].Add(uint32(newDocNum))
380380
}
381381
}
@@ -384,7 +384,6 @@ func (s *Scorch) introduceMerge(nextMerge *segmentMerge) {
384384
// obsolete segments wrt root in meantime, whatever
385385
// segments left behind in old map after processing
386386
// the root segments would be the obsolete segment set
387-
delete(nextMerge.old, segmentID)
388387
delete(nextMerge.mergedSegHistory, segmentID)
389388
} else if root.segment[i].LiveSize() > 0 {
390389
// this segment is staying
@@ -417,20 +416,22 @@ func (s *Scorch) introduceMerge(nextMerge *segmentMerge) {
417416
// before the newMerge introduction, need to clean the newly
418417
// merged segment wrt the current root segments, hence
419418
// applying the obsolete segment contents to newly merged segment
420-
for segID, ss := range nextMerge.mergedSegHistory {
419+
for _, ss := range nextMerge.mergedSegHistory {
421420
obsoleted := ss.oldSegment.DocNumbersLive()
422421
if obsoleted != nil {
423422
obsoletedIter := obsoleted.Iterator()
424423
for obsoletedIter.HasNext() {
425424
oldDocNum := obsoletedIter.Next()
426-
newDocNum := nextMerge.oldNewDocNums[segID][oldDocNum]
425+
newDocNum := ss.oldNewDocIDs[oldDocNum]
427426
newSegmentDeleted[ss.workerID].Add(uint32(newDocNum))
428427
}
429428
}
430429
}
431430

432431
skipped := true
433432
for i, newMergedSegment := range nextMerge.new {
433+
// checking if this newly merged segment is worth keeping based on
434+
// obsoleted doc count since the merge intro started
434435
if newMergedSegment != nil &&
435436
newMergedSegment.Count() > newSegmentDeleted[i].GetCardinality() {
436437
stats := newFieldStats()

index/scorch/merge.go

Lines changed: 0 additions & 92 deletions
Original file line numberDiff line numberDiff line change
@@ -319,7 +319,6 @@ func (s *Scorch) planMergeAtSnapshot(ctx context.Context,
319319
if segSnapshot.LiveSize() == 0 {
320320
atomic.AddUint64(&s.stats.TotFileMergeSegmentsEmpty, 1)
321321
oldMap[segSnapshot.id] = nil
322-
// mergedSegHistory[segSnapshot.id] = nil
323322
delete(mergedSegHistory, segSnapshot.id)
324323
} else {
325324
segmentsToMerge = append(segmentsToMerge, segSnapshot.segment)
@@ -387,8 +386,6 @@ func (s *Scorch) planMergeAtSnapshot(ctx context.Context,
387386

388387
sm := &segmentMerge{
389388
id: []uint64{newSegmentID},
390-
old: oldMap,
391-
oldNewDocNums: oldNewDocNums,
392389
mergedSegHistory: mergedSegHistory,
393390
new: []segment.Segment{seg},
394391
newCount: seg.Count(),
@@ -454,8 +451,6 @@ type mergedSegmentHistory struct {
454451

455452
type segmentMerge struct {
456453
id []uint64
457-
old map[uint64]*SegmentSnapshot
458-
oldNewDocNums map[uint64][]uint64
459454
new []segment.Segment
460455
mergedSegHistory map[uint64]*mergedSegmentHistory
461456
notifyCh chan *mergeTaskIntroStatus
@@ -532,8 +527,6 @@ func (s *Scorch) mergeSegmentBasesParallel(snapshot *IndexSnapshot, flushableObj
532527

533528
sm := &segmentMerge{
534529
id: newMergedSegmentIDs,
535-
old: make(map[uint64]*SegmentSnapshot, numSegments),
536-
oldNewDocNums: make(map[uint64][]uint64, numSegments),
537530
new: newMergedSegments,
538531
mergedSegHistory: make(map[uint64]*mergedSegmentHistory, numSegments),
539532
notifyCh: make(chan *mergeTaskIntroStatus),
@@ -543,10 +536,6 @@ func (s *Scorch) mergeSegmentBasesParallel(snapshot *IndexSnapshot, flushableObj
543536
for i, flushable := range flushableObjs {
544537
for j, idx := range flushable.sbIdxs {
545538
ss := snapshot.segment[idx]
546-
sm.old[ss.id] = ss
547-
sm.oldNewDocNums[ss.id] =
548-
newDocNumsSet[i][j]
549-
550539
// oldSegmentSnapshot.id -> {threadIdx, oldSegmentSnapshot, docIDs}
551540
sm.mergedSegHistory[ss.id] = &mergedSegmentHistory{
552541
workerID: uint64(i),
@@ -581,87 +570,6 @@ func (s *Scorch) mergeSegmentBasesParallel(snapshot *IndexSnapshot, flushableObj
581570
return newSnapshot, newMergedSegmentIDs, nil
582571
}
583572

584-
// perform a merging of the given SegmentBase instances into a new,
585-
// persisted segment, and synchronously introduce that new segment
586-
// into the root
587-
func (s *Scorch) mergeSegmentBases(snapshot *IndexSnapshot,
588-
sbs []segment.Segment, sbsDrops []*roaring.Bitmap,
589-
sbsIndexes []int) (*IndexSnapshot, uint64, error) {
590-
atomic.AddUint64(&s.stats.TotMemMergeBeg, 1)
591-
592-
memMergeZapStartTime := time.Now()
593-
594-
atomic.AddUint64(&s.stats.TotMemMergeZapBeg, 1)
595-
596-
newSegmentID := atomic.AddUint64(&s.nextSegmentID, 1)
597-
filename := zapFileName(newSegmentID)
598-
path := s.path + string(os.PathSeparator) + filename
599-
600-
newDocNums, _, err :=
601-
s.segPlugin.Merge(sbs, sbsDrops, path, s.closeCh, s)
602-
603-
atomic.AddUint64(&s.stats.TotMemMergeZapEnd, 1)
604-
605-
memMergeZapTime := uint64(time.Since(memMergeZapStartTime))
606-
atomic.AddUint64(&s.stats.TotMemMergeZapTime, memMergeZapTime)
607-
if atomic.LoadUint64(&s.stats.MaxMemMergeZapTime) < memMergeZapTime {
608-
atomic.StoreUint64(&s.stats.MaxMemMergeZapTime, memMergeZapTime)
609-
}
610-
611-
if err != nil {
612-
atomic.AddUint64(&s.stats.TotMemMergeErr, 1)
613-
return nil, 0, err
614-
}
615-
616-
seg, err := s.segPlugin.Open(path)
617-
if err != nil {
618-
atomic.AddUint64(&s.stats.TotMemMergeErr, 1)
619-
return nil, 0, err
620-
}
621-
622-
// update persisted stats
623-
atomic.AddUint64(&s.stats.TotPersistedItems, seg.Count())
624-
atomic.AddUint64(&s.stats.TotPersistedSegments, 1)
625-
626-
sm := &segmentMerge{
627-
id: []uint64{newSegmentID},
628-
old: make(map[uint64]*SegmentSnapshot, len(sbsIndexes)),
629-
oldNewDocNums: make(map[uint64][]uint64, len(sbsIndexes)),
630-
new: []segment.Segment{seg},
631-
notifyCh: make(chan *mergeTaskIntroStatus),
632-
}
633-
634-
for i, idx := range sbsIndexes {
635-
ss := snapshot.segment[idx]
636-
sm.old[ss.id] = ss
637-
sm.oldNewDocNums[ss.id] = newDocNums[i]
638-
}
639-
640-
select { // send to introducer
641-
case <-s.closeCh:
642-
_ = seg.DecRef()
643-
return nil, 0, segment.ErrClosed
644-
case s.merges <- sm:
645-
}
646-
647-
// blockingly wait for the introduction to complete
648-
var newSnapshot *IndexSnapshot
649-
introStatus := <-sm.notifyCh
650-
if introStatus != nil && introStatus.indexSnapshot != nil {
651-
newSnapshot = introStatus.indexSnapshot
652-
atomic.AddUint64(&s.stats.TotMemMergeSegments, uint64(len(sbs)))
653-
atomic.AddUint64(&s.stats.TotMemMergeDone, 1)
654-
if introStatus.skipped {
655-
// close the segment on skipping introduction.
656-
_ = newSnapshot.DecRef()
657-
_ = seg.Close()
658-
newSnapshot = nil
659-
}
660-
}
661-
662-
return newSnapshot, newSegmentID, nil
663-
}
664-
665573
func (s *Scorch) ReportBytesWritten(bytesWritten uint64) {
666574
atomic.AddUint64(&s.stats.TotFileMergeWrittenBytes, bytesWritten)
667575
}

0 commit comments

Comments
 (0)