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
5 changes: 4 additions & 1 deletion config/frac_version.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,9 @@ const (

// BinaryDataV5 - token blocks have zone maps (eng letters presense) and doc frequencies for heavy tokens
BinaryDataV5

// BinaryDataV6 - keep offsets separate in token block
BinaryDataV6
)

const CurrentFracVersion = BinaryDataV5
const CurrentFracVersion = BinaryDataV6
83 changes: 77 additions & 6 deletions frac/sealed/token/block_loader.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,9 @@ type Block struct {
Offsets []uint32
FreqIndexes []uint16 // indexes of tokens which have doc freqs (frequencies)
Freqs []uint32 // frequencies of certain tokens (how many docs have this token included at least once)

// TODO(cheb0) delete this field and convert V0..V5 to a new in-memory format when data all clusters have V6 fractions
FracVer config.BinaryDataVersion
}

func (b *Block) Size() int {
Expand All @@ -46,6 +49,8 @@ func (b Block) Pack(dst []byte, buf []uint32) []byte {
dst = binary.LittleEndian.AppendUint32(dst, uint32(len(b.Payload)))
dst = append(dst, b.Payload...)

dst = packer.CompressDeltaBitpackUint32(dst, b.Offsets, buf)

if len(b.FreqIndexes) > 0 {
dst = packer.CompressDeltaBitpackUint16(dst, b.FreqIndexes, buf)
dst = packer.CompressDeltaBitpackUint32(dst, b.Freqs, buf)
Expand All @@ -55,16 +60,53 @@ func (b Block) Pack(dst []byte, buf []uint32) []byte {
}

func (b *Block) Unpack(data []byte, fracVer config.BinaryDataVersion, unpackBuf *UnpackBuffer) error {
b.FracVer = fracVer

if fracVer >= config.BinaryDataV6 {
unpackBuf.Reset(fracVer)
return b.unpackV6(data, unpackBuf)
}

if fracVer >= config.BinaryDataV5 {
unpackBuf.Reset(fracVer)
return b.unpackV5(data, unpackBuf)
}
return b.unpackV1(data)
}

func (b *Block) unpackV1(data []byte) error {
b.Payload = append([]byte{}, data...)
return b.parseTokenPayload(b.Payload)
func (b *Block) unpackV6(data []byte, buf *UnpackBuffer) error {
if len(data) < util.SizeOfUint32 {
return fmt.Errorf("token block too short: %d bytes", len(data))
}
flags := data[0]
data = data[1:]

// token payload
payloadLen := binary.LittleEndian.Uint32(data[:util.SizeOfUint32])
data = data[util.SizeOfUint32:]
if uint32(len(data)) < payloadLen {
return fmt.Errorf("invalid token block payload length: %d, data len %d", payloadLen, len(data))
}

payload := data[:payloadLen]
data = data[payloadLen:]

b.Payload = append(b.Payload[:0], payload...)

// offsets
var err error
data, buf.decompressedUint32, err = packer.DecompressDeltaBitpackUint32(data, buf.decompressedUint32, buf.compressed)
if err != nil {
return err
}
b.Offsets = append(b.Offsets, buf.decompressedUint32...)

err = b.unpackFreqs(data, buf, flags)
if err != nil {
return err
}

return nil
}

func (b *Block) unpackV5(data []byte, buf *UnpackBuffer) error {
Expand All @@ -85,11 +127,22 @@ func (b *Block) unpackV5(data []byte, buf *UnpackBuffer) error {

b.Payload = append(b.Payload[:0], payload...)

if err := b.parseTokenPayload(payload); err != nil {
if err := b.parseTokenPayloadV5(payload); err != nil {
return err
}

err := b.unpackFreqs(data, buf, flags)
if err != nil {
return err
}

return nil
}

func (b *Block) unpackFreqs(data []byte, buf *UnpackBuffer, flags byte) error {
if flags&1 > 0 {
buf.decompressedUint32 = buf.decompressedUint32[:0]

var err error
data, buf.decompressedUint16, err = packer.DecompressDeltaBitpackUint16(data, buf.decompressedUint16, buf.compressed)
if err != nil {
Expand All @@ -103,11 +156,16 @@ func (b *Block) unpackV5(data []byte, buf *UnpackBuffer) error {
}
b.Freqs = append(b.Freqs, buf.decompressedUint32...)
}

return nil
}

func (b *Block) parseTokenPayload(data []byte) error {
func (b *Block) unpackV1(data []byte) error {
b.Payload = append([]byte{}, data...)
return b.parseTokenPayloadV5(b.Payload)
}

// parseTokenPayloadV5 derives offsets from tokens payload. Only used for v1...v5 legacy fractions.
func (b *Block) parseTokenPayloadV5(data []byte) error {
b.Offsets = b.Offsets[:0]

var offset uint32
Expand All @@ -129,6 +187,10 @@ func (b *Block) parseTokenPayload(data []byte) error {
}

func (b *Block) Len() int {
if b.FracVer >= config.BinaryDataV6 {
return len(b.Offsets) - 1
}

return len(b.Offsets)
}

Expand All @@ -147,6 +209,15 @@ func (b *Block) GetFreq(index int) uint32 {
}

func (b *Block) GetToken(index int) []byte {
if b.FracVer >= config.BinaryDataV6 {
return b.Payload[b.Offsets[index]:b.Offsets[index+1]]
}

return b.getTokenV5(index)
}

//go:noinline
func (b *Block) getTokenV5(index int) []byte {
offset := b.Offsets[index]
l := binary.LittleEndian.Uint32(b.Payload[offset:])
offset += uint32(util.SizeOfUint32) // skip val length
Expand Down
52 changes: 19 additions & 33 deletions frac/sealed/token/block_loader_test.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
package token

import (
"encoding/binary"
"testing"

"github.com/stretchr/testify/assert"
Expand All @@ -11,14 +10,12 @@ import (
)

func TestBlock_PackUnpack_NoFreq(t *testing.T) {
src := Block{
Payload: packTokenPayload([]byte("foo"), []byte("bar")),
}
src := buildTokenPayload([]byte("foo"), []byte("bar"))

var buf []uint32
packed := src.Pack(nil, buf)
var dst Block
require.NoError(t, dst.Unpack(packed, config.BinaryDataV5, &UnpackBuffer{}))
require.NoError(t, dst.Unpack(packed, config.BinaryDataV6, &UnpackBuffer{}))

assert.Equal(t, 2, dst.Len())
assert.Equal(t, []byte("foo"), dst.GetToken(0))
Expand All @@ -29,16 +26,14 @@ func TestBlock_PackUnpack_NoFreq(t *testing.T) {
}

func TestBlock_PackUnpack_WithFreq(t *testing.T) {
src := Block{
Payload: packTokenPayload([]byte("dog"), []byte("cat"), []byte("horse"), []byte("duck")),
FreqIndexes: []uint16{0, 2},
Freqs: []uint32{100, 200},
}
src := buildTokenPayload([]byte("dog"), []byte("cat"), []byte("horse"), []byte("duck"))
src.FreqIndexes = []uint16{0, 2}
src.Freqs = []uint32{100, 200}

var buf []uint32
packed := src.Pack(nil, buf)
var dst Block
require.NoError(t, dst.Unpack(packed, config.BinaryDataV5, &UnpackBuffer{}))
require.NoError(t, dst.Unpack(packed, config.BinaryDataV6, &UnpackBuffer{}))

assert.Equal(t, src.Payload, dst.Payload)

Expand All @@ -48,33 +43,19 @@ func TestBlock_PackUnpack_WithFreq(t *testing.T) {
assert.Equal(t, uint32(0), dst.GetFreq(3))
}

func TestBlock_Unpack_Legacy(t *testing.T) {
legacy := packTokenPayload([]byte("legacy"))

var dst Block
require.NoError(t, dst.Unpack(legacy, config.BinaryDataV4, &UnpackBuffer{}))

assert.Equal(t, legacy, dst.Payload)
assert.Equal(t, []uint32{0}, dst.Offsets)
assert.Empty(t, dst.FreqIndexes)
assert.Empty(t, dst.Freqs)
}

func TestBlock_UnpackBufferReuse(t *testing.T) {
src := Block{
Payload: packTokenPayload([]byte("a"), []byte("b")),
FreqIndexes: []uint16{1},
Freqs: []uint32{64},
}
src := buildTokenPayload([]byte("a"), []byte("b"))
src.FreqIndexes = []uint16{1}
src.Freqs = []uint32{64}

var packBuf []uint32
packed := src.Pack(nil, packBuf)

unpackBuf := &UnpackBuffer{}

var dst1, dst2 Block
require.NoError(t, dst1.Unpack(packed, config.BinaryDataV5, unpackBuf))
require.NoError(t, dst2.Unpack(packed, config.BinaryDataV5, unpackBuf))
require.NoError(t, dst1.Unpack(packed, config.BinaryDataV6, unpackBuf))
require.NoError(t, dst2.Unpack(packed, config.BinaryDataV6, unpackBuf))

assert.Equal(t, dst1.FreqIndexes, dst2.FreqIndexes)
assert.Equal(t, dst1.Freqs, dst2.Freqs)
Expand All @@ -83,11 +64,16 @@ func TestBlock_UnpackBufferReuse(t *testing.T) {
assert.Equal(t, uint32(64), dst2.GetFreq(1))
}

func packTokenPayload(tokens ...[]byte) []byte {
func buildTokenPayload(tokens ...[]byte) Block {
var payload []byte
var offsets []uint32
offsets = append(offsets, 0)
for _, tok := range tokens {
payload = binary.LittleEndian.AppendUint32(payload, uint32(len(tok)))
offsets = append(offsets, offsets[len(offsets)-1]+uint32(len(tok)))
payload = append(payload, tok...)
}
return payload
return Block{
Payload: payload,
Offsets: offsets,
}
}
12 changes: 8 additions & 4 deletions indexwriter/blocks.go
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
package indexwriter

import (
"encoding/binary"
"iter"
"math"
"unsafe"

"github.com/ozontech/seq-db/config"
"github.com/ozontech/seq-db/frac/sealed/lids"
"github.com/ozontech/seq-db/frac/sealed/seqids"
"github.com/ozontech/seq-db/frac/sealed/token"
Expand Down Expand Up @@ -56,6 +56,7 @@
blockIdx uint32
blockSize int
)
block.payload.FracVer = config.BinaryDataV6

var (
currentTID uint32
Expand Down Expand Up @@ -123,9 +124,12 @@
}
}

tokenIndex := uint32(len(block.payload.Offsets))
block.payload.Offsets = append(block.payload.Offsets, uint32(len(block.payload.Payload)))
block.payload.Payload = binary.LittleEndian.AppendUint32(block.payload.Payload, uint32(len(tok)))
offsets := block.payload.Offsets
if len(offsets) == 0 {
offsets = append(offsets, 0)
}
tokenIndex := uint32(len(offsets) - 1)
block.payload.Offsets = append(offsets, offsets[len(offsets)-1]+uint32(len(tok)))

Check failure on line 132 in indexwriter/blocks.go

View workflow job for this annotation

GitHub Actions / lint

appendAssign: append result not assigned to the same slice (gocritic)
block.payload.Payload = append(block.payload.Payload, tok...)

if len(tlids) >= tokenFreqAbsThreshold {
Expand Down
Loading