From d4f20eed33ba42a468abd3940b57e25f57c1e17d Mon Sep 17 00:00:00 2001 From: Alexey Sharov Date: Thu, 27 Aug 2026 12:24:43 +0700 Subject: [PATCH] db/etl: give both merges the prefix, the sift and the in-place advance etl has two k-way merges: the one mergeSortFiles runs over data providers, and the one that reads a buffer's sorted chunks. Neither needs to be the container/heap shape heap.go says it was copy-pasted from - a full bytes.Compare on every comparison, two of them per level, and for the provider heap a pop followed by a push per entry. Both now compare a cached 8-byte big-endian key prefix first, so a comparison reads the key only when two cursors agree on it. Both sift by sinking a hole to a leaf and climbing back, one compare a level rather than two, which also serves as the heapify step. The provider heap keeps its root in place while loadFunc runs and takes the next key there, so one sift replaces the pop and the push, and the initial fill is one heapInit over k appended elements rather than k pushes. The five variants - today, and each trick added in turn - were built side by side and checked to give byte-identical output before any was measured. Provider heap alone, 200k keys over k runs (M4): k4 -48%, k8 -55%, k16 -46%. The chunk merge, on EPYC: PutSortLoad/random_500k -19.7% and random_100k -9.9%, of which the prefix is -9.6% and the sift -5.3%. End to end: MergeSortFiles/file_only_100k -17.3%, SortableBufferLoadOnly/random_500k -18.0%, sorted_500k -17.0%. Keys arriving in order are unaffected either way, since chunks already in order end to end skip the heap. --- db/etl/collector.go | 24 +++++-- db/etl/etl_test.go | 49 +++++++++++++ db/etl/heap.go | 171 +++++++++++++++++++++++--------------------- 3 files changed, 156 insertions(+), 88 deletions(-) diff --git a/db/etl/collector.go b/db/etl/collector.go index 6157b93521e..6521001ab13 100644 --- a/db/etl/collector.go +++ b/db/etl/collector.go @@ -307,17 +307,19 @@ func mergeSortFiles(logPrefix string, providers []dataProvider, loadFunc simpleL } } - h := &Heap{} - heapInit(h) + h := &Heap{elems: make([]*HeapElem, 0, len(providers))} for i, provider := range providers { if key, value, err := provider.Next(); err == nil { - heapPush(h, &HeapElem{key, value, i}) + e := &HeapElem{Value: value, TimeIdx: i} + e.setKey(key) + h.elems = append(h.elems, e) } else /* we must have at least one entry per file */ { eee := fmt.Errorf("%s: error reading first readers: n=%d current=%d provider=%s err=%w", logPrefix, len(providers), i, provider, err) panic(eee) } } + heapInit(h) var prevK, prevV []byte @@ -331,7 +333,9 @@ func mergeSortFiles(logPrefix string, providers []dataProvider, loadFunc simpleL } } - element := heapPop(h) + // The root stays in the heap while loadFunc runs, then takes its + // provider's next key in place: one sift instead of a pop and a push. + element := h.elems[0] provider := providers[element.TimeIdx] // SortableOldestAppearedBuffer must guarantee that only 1 oldest value of key will appear @@ -363,9 +367,15 @@ func mergeSortFiles(logPrefix string, providers []dataProvider, loadFunc simpleL } } - if element.Key, element.Value, err = provider.Next(); err == nil { - heapPush(h, element) - } else if !errors.Is(err, io.EOF) { + key, value, err := provider.Next() + switch { + case err == nil: + element.setKey(key) + element.Value = value + heapFixRoot(h) + case errors.Is(err, io.EOF): + heapPopRoot(h) + default: return fmt.Errorf("%s: error while reading next element from disk: %w", logPrefix, err) } } diff --git a/db/etl/etl_test.go b/db/etl/etl_test.go index 75ece88ebaa..3a9730fdbb3 100644 --- a/db/etl/etl_test.go +++ b/db/etl/etl_test.go @@ -2265,3 +2265,52 @@ func TestCollectRejectsOversizedKey(t *testing.T) { require.Error(t, err) require.Contains(t, err.Error(), "exceeds") } + +// BenchmarkProviderHeap drives mergeSortFiles' heap on its own: k sorted runs +// of keys, advanced the way the load loop does, with no file or collector +// around it. +func BenchmarkProviderHeap(b *testing.B) { + const keyLen = 32 + for _, k := range []int{4, 8, 16} { + for _, count := range []int{200_000} { + b.Run(fmt.Sprintf("k%d_n%d", k, count), func(b *testing.B) { + runs := make([][][]byte, k) + for i := range runs { + runs[i] = make([][]byte, 0, count/k) + for j := range count / k { + key := make([]byte, keyLen) + x := uint64(i*count+j) * 6364136223846793005 //nolint:gosec + binary.BigEndian.PutUint64(key, x) + binary.BigEndian.PutUint64(key[8:], x^0xdeadbeef) + runs[i] = append(runs[i], key) + } + slices.SortFunc(runs[i], bytes.Compare) + } + val := make([]byte, 128) + b.ResetTimer() + for b.Loop() { + at := make([]int, k) + h := &Heap{elems: make([]*HeapElem, 0, k)} + for i := range runs { + e := &HeapElem{Value: val, TimeIdx: i} + e.setKey(runs[i][0]) + h.elems = append(h.elems, e) + at[i] = 1 + } + heapInit(h) + for h.Len() > 0 { + e := h.elems[0] + i := e.TimeIdx + if at[i] < len(runs[i]) { + e.setKey(runs[i][at[i]]) + at[i]++ + heapFixRoot(h) + continue + } + heapPopRoot(h) + } + } + }) + } + } +} diff --git a/db/etl/heap.go b/db/etl/heap.go index 9cfd84fe662..9cc9006aa38 100644 --- a/db/etl/heap.go +++ b/db/etl/heap.go @@ -18,116 +18,109 @@ package etl import ( "bytes" + "encoding/binary" "slices" ) +// HeapElem is the entry a provider currently sits on. pfx caches the first 8 +// key bytes, so a comparison reaches Key only when two providers agree there. type HeapElem struct { + pfx uint64 Key []byte Value []byte TimeIdx int } +func (e *HeapElem) setKey(k []byte) { e.Key, e.pfx = k, keyPrefix(k) } + +// Heap orders providers by the key each one sits on. Providers are numbered in +// the order their files were written, so the lower TimeIdx wins a tie and +// equal keys come out in the order they went in. type Heap struct { elems []*HeapElem } -func (h *Heap) Len() int { - return len(h.elems) -} +func (h *Heap) Len() int { return len(h.elems) } -func (h *Heap) Less(i, j int) bool { - if c := bytes.Compare(h.elems[i].Key, h.elems[j].Key); c != 0 { +func lessElem(a, b *HeapElem) bool { + if a.pfx != b.pfx { + return a.pfx < b.pfx + } + if c := bytes.Compare(a.Key, b.Key); c != 0 { return c < 0 } - return h.elems[i].TimeIdx < h.elems[j].TimeIdx -} - -func (h *Heap) Swap(i, j int) { - h.elems[i], h.elems[j] = h.elems[j], h.elems[i] -} - -func (h *Heap) Push(x *HeapElem) { - h.elems = append(h.elems, x) -} - -func (h *Heap) Pop() *HeapElem { - old := h.elems - n := len(old) - 1 - x := old[n] - //old[n].Key, old[n].Value, old[n].TimeIdx = nil, nil, 0 - old[n] = nil - h.elems = old[0:n] - return x + return a.TimeIdx < b.TimeIdx } -// ------ Copy-Paste of `container/heap/heap.go` without interface conversion - -// Init establishes the heap invariants required by the other routines in this package. -// Init is idempotent with respect to the heap invariants -// and may be called whenever the heap invariants may have been invalidated. -// The complexity is O(n) where n = h.Len(). +// heapInit orders elems appended in provider order. func heapInit(h *Heap) { - // heapify - n := h.Len() - for i := n/2 - 1; i >= 0; i-- { - down(h, i, n) + for i := len(h.elems)/2 - 1; i >= 0; i-- { + h.siftRoot(i) } } -// Push pushes the element x onto the heap. -// The complexity is O(log n) where n = h.Len(). -func heapPush(h *Heap, x *HeapElem) { - h.Push(x) - up(h, h.Len()-1) -} - -// Pop removes and returns the minimum element (according to Less) from the heap. -// The complexity is O(log n) where n = h.Len(). -// Pop is equivalent to Remove(h, 0). -func heapPop(h *Heap) *HeapElem { - n := h.Len() - 1 - h.Swap(0, n) - down(h, 0, n) - return h.Pop() -} +// heapFixRoot restores the order after the root provider moved to its next key. +func heapFixRoot(h *Heap) { h.siftRoot(0) } -func up(h *Heap, j int) { - for { - i := (j - 1) / 2 // parent - if i == j || !h.Less(j, i) { - break - } - h.Swap(i, j) - j = i +// heapPopRoot drops the root, whose provider has no more keys. +func heapPopRoot(h *Heap) { + last := len(h.elems) - 1 + h.elems[0] = h.elems[last] + h.elems[last] = nil + h.elems = h.elems[:last] + if last > 0 { + h.siftRoot(0) } } -func down(h *Heap, i0, n int) bool { - i := i0 +// siftRoot restores the heap under i, whose element changed: sink the hole to +// a leaf taking the smaller child, then climb back until the old element fits. +// One compare a level rather than the two a top-down sift needs, and the +// element that just moved on usually holds a larger key now, so the hole +// nearly always reaches a leaf. The climb back is also the heapify step. +func (h *Heap) siftRoot(i int) { + x, top, n := h.elems[i], i, len(h.elems) for { - j1 := 2*i + 1 - if j1 >= n || j1 < 0 { // j1 < 0 after int overflow + l := 2*i + 1 + if l >= n { break } - j := j1 // left child - if j2 := j1 + 1; j2 < n && h.Less(j2, j1) { - j = j2 // = 2*i + 2 // right child + if r := l + 1; r < n && lessElem(h.elems[r], h.elems[l]) { + l = r } - if !h.Less(j, i) { + h.elems[i] = h.elems[l] + i = l + } + for i > top { + p := (i - 1) / 2 + if lessElem(h.elems[p], x) { break } - h.Swap(i, j) - i = j + h.elems[i] = h.elems[p] + i = p } - return i > i0 + h.elems[i] = x } // ------ the merge over a sortableBuffer's sorted chunks +// keyPrefix is the first 8 bytes of k, zero-padded. Big-endian, so comparing +// two prefixes as integers orders them the way bytes.Compare would. +func keyPrefix(k []byte) uint64 { + if len(k) >= 8 { + return binary.BigEndian.Uint64(k) + } + var pad [8]byte + copy(pad[:], k) + return binary.BigEndian.Uint64(pad[:]) +} + // merger walks already-sorted chunks in key order, a cursor per chunk under a -// heap of chunk ids. +// heap of chunk ids. The prefixes live apart from the cursors because the heap +// reads them on nearly every comparison and they stay in L1. type merger struct { heap []int32 // chunk ids, ordered by their cursor's key + pfx []uint64 // each cursor's key prefix, by chunk id cur []cursor // by chunk id // Chunks already in order end to end, which ascending keys produce, are @@ -147,6 +140,7 @@ type cursor struct { func (m *merger) rewind(chunks []dataChunk) { clear(m.cur) // a shorter run would leave the old cursors pinning their chunks m.cur = slices.Grow(m.cur[:0], len(chunks))[:len(chunks)] + m.pfx = slices.Grow(m.pfx[:0], len(chunks))[:len(chunks)] for i := range chunks { m.cur[i] = cursor{ents: chunks[i].entries(), buf: chunks[i].buf} } @@ -204,13 +198,14 @@ func (m *merger) next() ([]byte, entryLoc, bool) { func (m *merger) release() { clear(m.cur) - m.cur, m.heap = m.cur[:0], m.heap[:0] + m.cur, m.pfx, m.heap = m.cur[:0], m.pfx[:0], m.heap[:0] m.chunk, m.concat = 0, false } func (m *merger) load(id int32) { c := &m.cur[id] c.key = keyOf(c.buf, c.ents[c.at]) + m.pfx[id] = keyPrefix(c.key) } // chunksInOrder reports whether every chunk's last key comes before the next @@ -236,26 +231,40 @@ func (m *merger) chunksInOrder() bool { // less orders two cursors by the key they sit on. Chunks fill in insertion // order, so the lower id wins a tie and equal keys keep the order they went in. func (m *merger) less(x, y int32) bool { + if px, py := m.pfx[x], m.pfx[y]; px != py { + return px < py + } if r := bytes.Compare(m.cur[x].key, m.cur[y].key); r != 0 { return r < 0 } return x < y } -// siftRoot restores the heap under i, whose element changed. +// siftRoot restores the heap under i, whose cursor just moved: sink the hole +// to a leaf taking the smaller child, then climb back until the old value +// fits. One compare a level instead of the usual two, and the cursor that just +// won usually holds a larger key, so the hole nearly always reaches a leaf. +// The climb back is also what makes it serve as the heapify step. func (m *merger) siftRoot(i int) { + x, top := m.heap[i], i for { - s, l, r := i, 2*i+1, 2*i+2 - if l < len(m.heap) && m.less(m.heap[l], m.heap[s]) { - s = l + l := 2*i + 1 + if l >= len(m.heap) { + break } - if r < len(m.heap) && m.less(m.heap[r], m.heap[s]) { - s = r + if r := l + 1; r < len(m.heap) && m.less(m.heap[r], m.heap[l]) { + l = r } - if s == i { - return + m.heap[i] = m.heap[l] + i = l + } + for i > top { + p := (i - 1) / 2 + if m.less(m.heap[p], x) { + break } - m.heap[i], m.heap[s] = m.heap[s], m.heap[i] - i = s + m.heap[i] = m.heap[p] + i = p } + m.heap[i] = x }