Skip to content
Merged
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
220 changes: 166 additions & 54 deletions db/etl/buffers.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,15 +33,15 @@ import (
)

const (
//SliceBuffer - just simple slice w
// SliceBuffer - just simple slice w
SortableSliceBuffer = iota
//SortableAppendBuffer - map[k] [v1 v2 v3]
// SortableAppendBuffer - map[k] [v1 v2 v3]
SortableAppendBuffer
// SortableOldestAppearedBuffer - buffer that keeps only the oldest entries.
// if first v1 was added under key K, then v2; only v1 will stay
SortableOldestAppearedBuffer

//BufIOSize - 128 pages | default is 1 page | increasing over `64 * 4096` doesn't show speedup on SSD/NVMe, but show speedup in cloud drives
// BufIOSize - 128 pages | default is 1 page | increasing over `64 * 4096` doesn't show speedup on SSD/NVMe, but show speedup in cloud drives
BufIOSize = 128 * 4096

entryLocSize = 16 // sizeof(entryLoc): insertionOrder(4) + offset(4) + keyLen(4) + valLen(4)
Expand Down Expand Up @@ -79,23 +79,64 @@ func writeSortedEntries(w io.Writer, entries []sortableBufferEntry) error {

var BufferOptimalSize = dbg.EnvDataSize("ETL_OPTIMAL", 256*datasize.MB) /* var because we want to sometimes change it from tests or command-line flags */

// etlSmallBufRAM (BufferOptimalSize/8) bounds the flush threshold so a full
// set of domain/history/index flush collectors (~17 per batch writer) stays
// around 512 MB when all run full. Pooled buffers start empty and grow with
// the data they actually see; grown capacity survives reuse (Reset preserves
// cap), so hot collectors amortize growth while never-full ones stay small.
var etlSmallBufRAM = dbg.EnvDataSize("ETL_SMALL", BufferOptimalSize/8)
var SmallSortableBuffers = NewAllocator(&sync.Pool{
New: func() any {
return NewSortableBuffer(etlSmallBufRAM)
},
})
var etlLargeBufRAM = BufferOptimalSize
var LargeSortableBuffers = NewAllocator(&sync.Pool{
New: func() any {
return NewSortableBuffer(etlLargeBufRAM)
},
})
// etlSmallBufRAM (BufferOptimalSize/8) bounds the flush threshold:
// 3_domains * 2 + 3_history * 1 + 4_indices * 2 = 17 etl collectors,
// 17*(256Mb/8) = 544Mb for all collectors combined. Buffers pool their
// chunks — see dataChunks below.
var (
etlSmallBufRAM = dbg.EnvDataSize("ETL_SMALL", BufferOptimalSize/8)
SmallSortableBuffers = NewAllocator(&sync.Pool{
New: func() any {
return NewSortableBuffer(etlSmallBufRAM)
},
})
)

var (
etlLargeBufRAM = BufferOptimalSize
LargeSortableBuffers = NewAllocator(&sync.Pool{
New: func() any {
return NewSortableBuffer(etlLargeBufRAM)
},
})
)

const (
// sortableBuffer stores key/value bytes in chunks of a power-of-two size, so
// entryLoc.offset can pack the chunk index with the offset inside the chunk
// and splitting the two is a shift and a mask. 1MB is also the least a
// collector can hold once it takes a chunk at all.
dataChunkBits = 20
dataChunkSize = 1 << dataChunkBits // 1MB

// The chunk index takes what is left of a positive int32, so one buffer
// addresses at most maxDataChunks*dataChunkSize bytes (~2GB); nextChunk
// panics past that. NewSortableBuffer's MaxInt32 bound on optimalSize
// does not fully rule this out, since Put can grow the buffer past
// optimalSize before CheckFlushSize is checked.
maxDataChunks = math.MaxInt32>>dataChunkBits + 1
)

// dataChunks are shared by all sortableBuffer instances: a buffer takes chunks as
// it fills and gives them back on Reset, instead of pinning its peak size forever.
var dataChunks = sync.Pool{New: func() any {
c := make([]byte, dataChunkSize)
return &c
}}

func getDataChunk() []byte { return *dataChunks.Get().(*[]byte) }

// isPooledChunk reports whether c came from the pool. An oversized entry gets a
// private chunk instead, and handing that one back would let a later
// getDataChunk give an unrelated buffer a chunk of the wrong size.
func isPooledChunk(c []byte) bool { return len(c) == dataChunkSize }

func putDataChunk(c []byte) {
if !isPooledChunk(c) {
return
}
dataChunks.Put(&c)
}

type Buffer interface {
// Put does copy `k` and `v`
Expand Down Expand Up @@ -123,9 +164,10 @@ var (
_ Buffer = &oldestEntrySortableBuffer{}
)

// entryLoc stores the location of a key/value pair within sortableBuffer.data.
// Key occupies data[offset : offset+keyLen], value follows at data[offset+max(0,keyLen) : ...+valLen].
// keyLen/valLen of -1 indicates nil.
// entryLoc stores the location of a key/value pair inside sortableBuffer.
// offset packs the chunk index and the offset inside that chunk:
// idx<<dataChunkBits | off. Key occupies chunk[off : off+keyLen], value follows
// right after it. keyLen/valLen of -1 indicates nil.
type entryLoc struct {
insertionOrder int32 // enables stable sort via unstable SortFunc
offset int32
Expand All @@ -143,16 +185,45 @@ func NewSortableBuffer(bufferOptimalSize datasize.ByteSize) *sortableBuffer {
}

type sortableBuffer struct {
entries []entryLoc
data []byte
entries []entryLoc
// chunks hold the key/value bytes. Growing by chunk instead of by one big
// slice keeps Put from re-allocating and copying everything collected so
// far. All chunks are dataChunkSize, except the private chunk an entry
// larger than that gets. cur is the chunk being filled.
chunks [][]byte
cur []byte
curBase int32 // packed location of cur's first byte: curIdx<<dataChunkBits
curOff int32
chunkBytes int
optimalSize int
}

// nextChunk starts a chunk able to hold n bytes. An entry never straddles
// chunks, so Get can hand out direct references.
func (b *sortableBuffer) nextChunk(n int) {
if len(b.chunks) >= maxDataChunks {
panic(fmt.Sprintf("etl: sortableBuffer exceeded %d chunks", maxDataChunks))
}
if n > dataChunkSize {
b.cur = make([]byte, n)
} else {
b.cur = getDataChunk()
}
b.chunks = append(b.chunks, b.cur)
b.curBase = int32(len(b.chunks)-1) << dataChunkBits //nolint:gosec
b.curOff = 0
b.chunkBytes += len(b.cur)
}

// entryData returns e's bytes: the key, immediately followed by the value.
func (b *sortableBuffer) entryData(e *entryLoc) []byte {
return b.chunks[e.offset>>dataChunkBits][e.offset&(dataChunkSize-1):]
}

// Put adds key and value to the buffer. These slices will not be accessed later,
// so no copying is necessary
func (b *sortableBuffer) Put(k, v []byte) {
e := entryLoc{
offset: int32(len(b.data)), //nolint:gosec
keyLen: int32(len(k)), //nolint:gosec
valLen: int32(len(v)), //nolint:gosec
insertionOrder: int32(len(b.entries)), //nolint:gosec
Expand All @@ -163,11 +234,26 @@ func (b *sortableBuffer) Put(k, v []byte) {
if v == nil {
e.valLen = -1
}
if n := len(k) + len(v); n > 0 {
off := b.curOff
if int(off)+n > len(b.cur) {
b.nextChunk(n)
off = 0
}
data := b.cur[off:]
copy(data, k)
copy(data[len(k):], v)
e.offset = b.curBase | off
b.curOff = off + int32(n) //nolint:gosec
}
b.entries = append(b.entries, e)
b.data = append(append(b.data, k...), v...)
}

func (b *sortableBuffer) Size() int { return len(b.data) + len(b.entries)*entryLocSize }
// Size counts the stored bytes, the tails wasted by the chunks already filled,
// and entryLocSize bytes of metadata per entry.
func (b *sortableBuffer) Size() int {
return b.chunkBytes - (len(b.cur) - int(b.curOff)) + len(b.entries)*entryLocSize
}

func (b *sortableBuffer) Len() int {
return len(b.entries)
Expand All @@ -176,49 +262,69 @@ func (b *sortableBuffer) Len() int {
func (b *sortableBuffer) Get(i int) ([]byte, []byte) {
e := &b.entries[i]
kLen, vLen := int(e.keyLen), int(e.valLen)
keyOffset := int(e.offset)
valOffset := keyOffset
if kLen > 0 {
valOffset += kLen
}
var key, val []byte
if kLen > 0 {
key = b.data[keyOffset : keyOffset+kLen]
} else if kLen == 0 {
if kLen == 0 {
key = []byte{}
}
if vLen > 0 {
val = b.data[valOffset : valOffset+vLen]
} else if vLen == 0 {
if vLen == 0 {
val = []byte{}
}
if kLen <= 0 && vLen <= 0 {
return key, val
}
data := b.entryData(e)
if kLen > 0 {
key = data[:kLen:kLen]
data = data[kLen:]
}
if vLen > 0 {
val = data[:vLen:vLen]
}
return key, val
}

// Prealloc sizes the entries slice. predictDataSize only reserves room in the
// chunks slice for the chunk pointers; the chunks themselves are still taken
// one at a time, which is what keeps an idle buffer from holding its peak.
func (b *sortableBuffer) Prealloc(predictKeysAmount, predictDataSize int) Buffer {
if cap(b.entries) < predictKeysAmount {
b.entries = make([]entryLoc, 0, predictKeysAmount)
}
if cap(b.data) < predictDataSize {
b.data = make([]byte, 0, predictDataSize)
if n := predictDataSize/dataChunkSize + 1; cap(b.chunks) < n {
Comment thread
AskAlexSharov marked this conversation as resolved.
b.chunks = slices.Grow(b.chunks, n)
}
return b
}

func (b *sortableBuffer) Reset() {
b.entries = b.entries[:0]
b.data = b.data[:0]
for i, c := range b.chunks {
putDataChunk(c)
b.chunks[i] = nil
}
b.chunks = b.chunks[:0]
b.cur, b.curBase, b.curOff = nil, 0, 0
b.chunkBytes = 0
}
func (b *sortableBuffer) SizeLimit() int { return b.optimalSize }
func (b *sortableBuffer) Sort() {
data := b.data
cmp := func(a, b entryLoc) int {
aKey := data[a.offset : a.offset+max(a.keyLen, 0)]
bKey := data[b.offset : b.offset+max(b.keyLen, 0)]
if c := bytes.Compare(aKey, bKey); c != 0 {
chunks := b.chunks
// Key extraction stays inside cmp: pdqsortCmpFunc calls the comparator
// indirectly, so a separate closure never inlines and costs a call per key.
cmp := func(x, y entryLoc) int {
var xk, yk []byte
if x.keyLen > 0 {
off := x.offset & (dataChunkSize - 1)
xk = chunks[x.offset>>dataChunkBits][off : off+x.keyLen]
}
if y.keyLen > 0 {
off := y.offset & (dataChunkSize - 1)
yk = chunks[y.offset>>dataChunkBits][off : off+y.keyLen]
}
if c := bytes.Compare(xk, yk); c != 0 {
return c
}
return int(a.insertionOrder - b.insertionOrder) // StableSort: preserve insertion order for duplicate keys
return int(x.insertionOrder - y.insertionOrder) // StableSort: preserve insertion order for duplicate keys
}
if slices.IsSortedFunc(b.entries, cmp) {
return
Expand All @@ -235,28 +341,28 @@ func (b *sortableBuffer) Write(w io.Writer) error {
for i := range b.entries {
e := &b.entries[i]
kLen, vLen := int(e.keyLen), int(e.valLen)
keyOffset := int(e.offset)
valOffset := keyOffset
if kLen > 0 {
valOffset += kLen
var data []byte
if kLen > 0 || vLen > 0 {
data = b.entryData(e)
}
// write key
n := binary.PutVarint(numBuf[:], int64(e.keyLen))
if _, err := w.Write(numBuf[:n]); err != nil {
return err
}
if kLen > 0 {
if _, err := w.Write(b.data[keyOffset : keyOffset+kLen]); err != nil {
if _, err := w.Write(data[:kLen]); err != nil {
return err
}
data = data[kLen:]
}
// write value
n = binary.PutVarint(numBuf[:], int64(e.valLen))
if _, err := w.Write(numBuf[:n]); err != nil {
return err
}
if vLen > 0 {
if _, err := w.Write(b.data[valOffset : valOffset+vLen]); err != nil {
if _, err := w.Write(data[:vLen]); err != nil {
return err
}
}
Expand Down Expand Up @@ -294,6 +400,7 @@ func (b *appendSortableBuffer) SizeLimit() int { return b.optimalSize }
func (b *appendSortableBuffer) Len() int {
return len(b.entries)
}

func (b *appendSortableBuffer) Sort() {
b.sortedBuf = b.sortedBuf[:0]
if cap(b.sortedBuf) < len(b.entries) {
Expand All @@ -316,11 +423,13 @@ func (b *appendSortableBuffer) Swap(i, j int) {
func (b *appendSortableBuffer) Get(i int) ([]byte, []byte) {
return b.sortedBuf[i].key, b.sortedBuf[i].value
}

func (b *appendSortableBuffer) Reset() {
b.sortedBuf = nil
b.entries = make(map[string][]byte)
b.size = 0
}

func (b *appendSortableBuffer) Prealloc(predictKeysAmount, predictDataSize int) Buffer {
b.entries = make(map[string][]byte, predictKeysAmount) // maps have no cap(), always recreate
if cap(b.sortedBuf) < predictKeysAmount {
Expand Down Expand Up @@ -392,11 +501,13 @@ func (b *oldestEntrySortableBuffer) Swap(i, j int) {
func (b *oldestEntrySortableBuffer) Get(i int) ([]byte, []byte) {
return b.sortedBuf[i].key, b.sortedBuf[i].value
}

func (b *oldestEntrySortableBuffer) Reset() {
b.sortedBuf = nil
b.entries = make(map[string][]byte)
b.size = 0
}

func (b *oldestEntrySortableBuffer) Prealloc(predictKeysAmount, predictDataSize int) Buffer {
b.entries = make(map[string][]byte, predictKeysAmount) // maps have no cap(), always recreate
if cap(b.sortedBuf) < predictKeysAmount {
Expand All @@ -408,6 +519,7 @@ func (b *oldestEntrySortableBuffer) Prealloc(predictKeysAmount, predictDataSize
func (b *oldestEntrySortableBuffer) Write(w io.Writer) error {
return writeSortedEntries(w, b.sortedBuf)
}

func (b *oldestEntrySortableBuffer) CheckFlushSize() bool {
return b.size >= b.optimalSize
}
Expand Down
Loading
Loading