Skip to content
Open
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 pkg/querier/batch/batch.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,16 +61,35 @@ func NewChunkMergeIterator(it chunkenc.Iterator, chunks []chunk.Chunk, _, _ mode
return NewGenericChunkMergeIterator(it, converted)
}

// bufProvider is satisfied by any iterator carrying a reusable buffer.
type bufProvider interface {
chunkenc.Iterator
GetBuf() any
SetBuf(any)
}

// NewGenericChunkMergeIterator returns a chunkenc.Iterator that merges generic chunks together.
func NewGenericChunkMergeIterator(it chunkenc.Iterator, chunks []GenericChunk) chunkenc.Iterator {

var underlying iterator
var sharedBuf batchStream

if bp, ok := it.(bufProvider); ok {
if buf, ok := bp.GetBuf().(batchStream); ok {
sharedBuf = buf
}
}

if ia, ok := it.(*iteratorAdapter); ok {
underlying = ia.underlying
}

iter := newMergeIterator(underlying, chunks)
iter := newMergeIterator(underlying, chunks, sharedBuf)

if bp, ok := it.(bufProvider); ok && bp.GetBuf() == nil {
bp.SetBuf(iter.batchesBuf)
}

return newIteratorAdapter(iter)
}

Expand Down
36 changes: 36 additions & 0 deletions pkg/querier/batch/batch_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -171,3 +171,39 @@ func createChunks(b *testing.B, step time.Duration, numChunks, numSamplesPerChun

return result
}

type testBufCarrier struct {
chunkenc.Iterator
buf any
}

func (t *testBufCarrier) GetBuf() any { return t.buf }
func (t *testBufCarrier) SetBuf(v any) { t.buf = v }

func BenchmarkNewChunkMergeIterator_BufferCarrier(b *testing.B) {
const numSeries = 10000
chunks := createChunks(b, step, 10, 100, 3, promchunk.PrometheusXorChunk)

b.Run("no_carrier", func(b *testing.B) {
b.ReportAllocs()
for b.Loop() {
for range numSeries {
it := NewChunkMergeIterator(nil, chunks, 0, 0)
for it.Next() != chunkenc.ValNone {
}
}
}
})

b.Run("with_carrier", func(b *testing.B) {
b.ReportAllocs()
for b.Loop() {
carrier := &testBufCarrier{}
for range numSeries {
it := NewChunkMergeIterator(carrier, chunks, 0, 0)
for it.Next() != chunkenc.ValNone {
}
}
}
})
}
15 changes: 11 additions & 4 deletions pkg/querier/batch/merge.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ type mergeIterator struct {
currErr error
}

func newMergeIterator(it iterator, cs []GenericChunk) *mergeIterator {
func newMergeIterator(it iterator, cs []GenericChunk, sharedBuf batchStream) *mergeIterator {
css := partitionChunks(cs)

var c *mergeIterator
Expand All @@ -32,9 +32,16 @@ func newMergeIterator(it iterator, cs []GenericChunk) *mergeIterator {
c = mIterator.Reset(len(css))
} else {
c = &mergeIterator{
h: make(iteratorHeap, 0, len(css)),
batches: make(batchStream, 0, len(css)),
batchesBuf: make(batchStream, len(css)),
h: make(iteratorHeap, 0, len(css)),
batches: make(batchStream, 0, len(css)),
}
if sharedBuf != nil && cap(sharedBuf) >= len(css) {
c.batchesBuf = sharedBuf[:len(css)]
for i := range c.batchesBuf {
c.batchesBuf[i] = promchunk.Batch{}
}
} else {
c.batchesBuf = make(batchStream, len(css))
}
}

Expand Down
14 changes: 7 additions & 7 deletions pkg/querier/batch/merge_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,10 +20,10 @@ func TestMergeIter(t *testing.T) {
chunk4 := mkGenericChunk(t, model.TimeFromUnix(75), 100, enc)
chunk5 := mkGenericChunk(t, model.TimeFromUnix(100), 100, enc)

iter := newMergeIterator(nil, []GenericChunk{chunk1, chunk2, chunk3, chunk4, chunk5})
iter := newMergeIterator(nil, []GenericChunk{chunk1, chunk2, chunk3, chunk4, chunk5}, nil)
testIter(t, 200, newIteratorAdapter(iter), enc)

iter = newMergeIterator(iter, []GenericChunk{chunk1, chunk2, chunk3, chunk4, chunk5})
iter = newMergeIterator(iter, []GenericChunk{chunk1, chunk2, chunk3, chunk4, chunk5}, nil)
testSeek(t, 200, newIteratorAdapter(iter), enc)
})
}
Expand All @@ -33,16 +33,16 @@ func BenchmarkMergeIterator(b *testing.B) {
for i := range 10 {
chunks = append(chunks, mkGenericChunk(b, model.Time(i*25), 120, encoding.PrometheusXorChunk))
}
iter := newMergeIterator(nil, chunks)
iter := newMergeIterator(nil, chunks, nil)

for _, r := range []bool{true, false} {
b.Run(fmt.Sprintf("reuse-%t", r), func(b *testing.B) {
b.ReportAllocs()
for b.Loop() {
if r {
iter = newMergeIterator(iter, chunks)
iter = newMergeIterator(iter, chunks, nil)
} else {
iter = newMergeIterator(nil, chunks)
iter = newMergeIterator(nil, chunks, nil)
}
a := newIteratorAdapter(iter)
for a.Next() != chunkenc.ValNone {
Expand All @@ -67,10 +67,10 @@ func TestMergeHarder(t *testing.T) {
chunks = append(chunks, mkGenericChunk(t, from, samples, enc))
from = from.Add(time.Duration(offset) * time.Second)
}
iter := newMergeIterator(nil, chunks)
iter := newMergeIterator(nil, chunks, nil)
testIter(t, offset*numChunks+samples-offset, newIteratorAdapter(iter), enc)

iter = newMergeIterator(iter, chunks)
iter = newMergeIterator(iter, chunks, nil)
testSeek(t, offset*numChunks+samples-offset, newIteratorAdapter(iter), enc)
})
}
Loading