Skip to content
Draft
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
28 changes: 19 additions & 9 deletions pkg/storage/indexheader/binary_writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,10 @@ const (
// BinaryFormatV1 represents first version of index-header file.
BinaryFormatV1 = 1

// BinaryFormatV2 represents the second version of the index-header file,
// which contains only the symbols table.
BinaryFormatV2 = 2

indexTOCLen = 6*8 + crc32.Size
BinaryTOCLen = 2*8 + crc32.Size // 16 (2 x uint64) + 4 (CRC32)
// HeaderLen represents number of bytes reserved of index header for header.
Expand Down Expand Up @@ -69,7 +73,7 @@ type BinaryTOC struct {
}

// WriteBinary build index-header file from the pieces of index in object storage.
func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, filename string) (err error) {
func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, filename string, writeV2 bool) (err error) {
ir, indexVersion, err := newChunkedIndexReader(ctx, bkt, id)
if err != nil {
return errors.Wrap(err, "new index reader")
Expand All @@ -79,7 +83,7 @@ func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, f
// Buffer for copying and encbuffers.
// This also will control the size of file writer buffer.
buf := make([]byte, 32*1024)
bw, err := newBinaryWriter(tmpFilename, buf)
bw, err := newBinaryWriter(tmpFilename, buf, writeV2)
if err != nil {
return errors.Wrap(err, "new binary index header writer")
}
Expand All @@ -101,12 +105,14 @@ func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, f
return errors.Wrap(err, "flush")
}

if err := ir.CopyPostingsOffsets(bw.PostingOffsetsWriter(), buf); err != nil {
return err
}
if !writeV2 {
if err := ir.CopyPostingsOffsets(bw.PostingOffsetsWriter(), buf); err != nil {
return err
}

if err := bw.f.Flush(); err != nil {
return errors.Wrap(err, "flush")
if err := bw.f.Flush(); err != nil {
return errors.Wrap(err, "flush")
}
}

if err := bw.WriteTOC(); err != nil {
Expand Down Expand Up @@ -243,7 +249,7 @@ type binaryWriter struct {
crc32 hash.Hash
}

func newBinaryWriter(fn string, buf []byte) (w *binaryWriter, err error) {
func newBinaryWriter(fn string, buf []byte, writeV2 bool) (w *binaryWriter, err error) {
df, err := fileutil.OpenDir(filepath.Dir(fn))
if err != nil {
return nil, err
Expand Down Expand Up @@ -274,7 +280,11 @@ func newBinaryWriter(fn string, buf []byte) (w *binaryWriter, err error) {

w.buf.Reset()
w.buf.PutBE32(MagicIndex)
w.buf.PutByte(BinaryFormatV1)
if writeV2 {
w.buf.PutByte(BinaryFormatV2)
} else {
w.buf.PutByte(BinaryFormatV1)
}

return w, w.f.Write(w.buf.Get())
}
Expand Down
6 changes: 6 additions & 0 deletions pkg/storage/indexheader/header.go
Original file line number Diff line number Diff line change
Expand Up @@ -117,21 +117,27 @@ const (

var (
errInvalidIndexHeaderSection = errors.New(fmt.Sprintf("invalid index-header section; must be one of: %s", SectionPostingsOffsetsTable))
errInvalidIndexHeaderVersion = errors.New("bucket reader must be enabled to write v2 index-header")
)

type BucketReaderConfig struct {
Enabled bool `yaml:"enabled" category:"experimental"`
BucketIndexSections string `yaml:"index_sections" category:"experimental"`
WriteV2IndexHeader bool `yaml:"write_v2_index_header" category:"experimental"`

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Introducing a new temporary flag so that we can keep enabling bucket reader without enabling v2 if necessary

}

func (cfg *BucketReaderConfig) RegisterFlagsWithPrefix(f *flag.FlagSet, prefix string) {
f.BoolVar(&cfg.Enabled, prefix+"enabled", false, fmt.Sprintf("Enable reading TSDB index-header sections from object storage. When enabled, the configured -%s are not downloaded to local disk.", prefix+"index-sections"))
f.StringVar(&cfg.BucketIndexSections, prefix+"index-sections", SectionPostingsOffsetsTable, fmt.Sprintf("Index sections to read from object storage instead of local disk. Valid sections: %s", SectionPostingsOffsetsTable))
f.BoolVar(&cfg.WriteV2IndexHeader, prefix+"write-v2-index-header", false, "Write only symbols table to on-disk index header.")
}

func (cfg *BucketReaderConfig) Validate() error {
if !slices.Contains([]string{SectionPostingsOffsetsTable}, cfg.BucketIndexSections) {
return errInvalidIndexHeaderSection
}
if cfg.WriteV2IndexHeader && !cfg.Enabled {
return errInvalidIndexHeaderVersion
}
return nil
}
6 changes: 3 additions & 3 deletions pkg/storage/indexheader/header_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ func TestReadersComparedToIndexHeader(t *testing.T) {
t.Run(testBlock.version, func(t *testing.T) {
id := testBlock.id
indexName := filepath.Join(tmpDir, id.String(), block.IndexHeaderFilename)
require.NoError(t, WriteBinary(ctx, bkt, id, indexName))
require.NoError(t, WriteBinary(ctx, bkt, id, indexName, false))

indexFile, err := fileutil.OpenMmapFile(filepath.Join(tmpDir, id.String(), block.IndexFilename))
require.NoError(t, err)
Expand Down Expand Up @@ -529,7 +529,7 @@ func labelValuesTestCases(t test.TB) (tests map[string][]labelValuesTestCase, bl
require.NoError(t, err)

indexName := filepath.Join(tmpDir, id.String(), block.IndexHeaderFilename)
require.NoError(t, WriteBinary(ctx, bkt, id, indexName))
require.NoError(t, WriteBinary(ctx, bkt, id, indexName, false))

indexFile, err := fileutil.OpenMmapFile(filepath.Join(tmpDir, id.String(), block.IndexFilename))
require.NoError(t, err)
Expand Down Expand Up @@ -602,7 +602,7 @@ func BenchmarkBinaryWrite(t *testing.B) {

t.ResetTimer()
for i := 0; i < t.N; i++ {
require.NoError(t, WriteBinary(ctx, bkt, m.ULID, fn))
require.NoError(t, WriteBinary(ctx, bkt, m.ULID, fn, false))
}
}

Expand Down
2 changes: 1 addition & 1 deletion pkg/storage/indexheader/lazy_binary_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -204,7 +204,7 @@ func ensureIndexHeaderOnDisk(
level.Debug(logger).Log("msg", "index-header does not exist on disk; will build from bucket", "path", indexHeaderPath)

start := time.Now()
if err := WriteBinary(ctx, bkt, blockID, indexHeaderPath); err != nil {
if err := WriteBinary(ctx, bkt, blockID, indexHeaderPath, false); err != nil {
level.Error(logger).Log("msg", "failed to create index-header", "err", err)
return err
}
Expand Down
2 changes: 1 addition & 1 deletion pkg/storage/indexheader/lazy_binary_reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -694,7 +694,7 @@ func BenchmarkLazyBinaryReader_LoadReader(b *testing.B) {
require.NoError(b, err)

indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename)
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName))
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false))

diskReaderBenchFactory := func(
cachingBucket *bucketcache.CachingBucket,
Expand Down
10 changes: 5 additions & 5 deletions pkg/storage/indexheader/reader_benchmarks_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ func BenchmarkLookupSymbol(b *testing.B) {
require.NoError(b, err)

indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename)
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName))
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false))

// TODO: are these sensible values for parallelism?
for _, parallelism := range []int{1, 2, 4, 8, 20, 100} {
Expand Down Expand Up @@ -127,7 +127,7 @@ func BenchmarkLabelNames(b *testing.B) {
require.NoError(b, err)

indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename)
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName))
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false))

binaryReader, err := NewStreamBinaryReader(ctx, idIndexV2, objstore.WithNoopInstr(bkt), blockDir, Config{}, 32, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil))
require.NoError(b, err)
Expand Down Expand Up @@ -170,7 +170,7 @@ func BenchmarkLabelValuesOffsetsIndexV2(b *testing.B) {
require.NoError(b, err)

indexName := filepath.Join(dir, blockID.String(), block.IndexHeaderFilename)
require.NoError(b, WriteBinary(ctx, bkt, blockID, indexName))
require.NoError(b, WriteBinary(ctx, bkt, blockID, indexName, false))

bucketReg := prometheus.NewPedanticRegistry()

Expand Down Expand Up @@ -355,7 +355,7 @@ func BenchmarkPostingsOffset(b *testing.B) {
require.NoError(b, err)

indexName := filepath.Join(dir, idIndexV2.String(), block.IndexHeaderFilename)
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName))
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false))

b.Run(fmt.Sprintf("%vNames%vValues", nameCount, valueCount), func(b *testing.B) {
binaryReader, err := NewStreamBinaryReader(ctx, idIndexV2, objstore.WithNoopInstr(bkt), dir, Config{}, 32, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil))
Expand Down Expand Up @@ -401,7 +401,7 @@ func BenchmarkNewStreamBinaryReader(b *testing.B) {
require.NoError(b, err)

indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename)
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName))
require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false))

b.Run(fmt.Sprintf("%vNames%vValues", nameCount, valueCount), func(b *testing.B) {
for i := 0; i < b.N; i++ {
Expand Down
15 changes: 10 additions & 5 deletions pkg/storage/indexheader/sparse_header.go
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,7 @@ func BuildAndWriteSparseHeaderFromTSDBIndex(
}

allSymbolsCount, sparseSymbolsOffsets, sparsePostingsOffsets, err := buildInMemorySparseHeaderFromIndexHeader(
ctx, indexTOC, filePoolDecbufFactory, sparseSampleFactor, false, l,
ctx, sectionSource{indexTOC, filePoolDecbufFactory}, sectionSource{indexTOC, filePoolDecbufFactory}, sparseSampleFactor, false, l,
)
if err != nil {
return fmt.Errorf("cannot build sparse index-header values from full index: %w", err)
Expand All @@ -299,10 +299,15 @@ func BuildAndWriteSparseHeaderFromTSDBIndex(
return nil
}

type sectionSource struct {
toc *TOCCompat
decbufFactory streamencoding.DecbufFactory
}

func buildInMemorySparseHeaderFromIndexHeader(
ctx context.Context,
toc *TOCCompat,
decbufFactory streamencoding.DecbufFactory,
symbols sectionSource,
postingsOffsets sectionSource,
sparseSampleFactor int,
doChecksum bool,
l log.Logger,
Expand All @@ -321,13 +326,13 @@ func buildInMemorySparseHeaderFromIndexHeader(
level.Info(l).Log("msg", "creating sparse index-header from full index-header")

allSymbolsCount, sparseSymbolsOffsets, err = streamindex.SparseValuesFromSymbolsTable(
ctx, decbufFactory, int(toc.Symbols), doChecksum,
ctx, symbols.decbufFactory, int(symbols.toc.Symbols), doChecksum,
)
if err != nil {
return -1, nil, nil, err
}

sparsePostingsOffsets, err = streamindex.SparseValuesFromPostingsOffsetsTable(ctx, decbufFactory, int(toc.PostingsOffsetTable), toc.PostingsListEnd, sparseSampleFactor, doChecksum)
sparsePostingsOffsets, err = streamindex.SparseValuesFromPostingsOffsetsTable(ctx, postingsOffsets.decbufFactory, int(postingsOffsets.toc.PostingsOffsetTable), postingsOffsets.toc.PostingsListEnd, sparseSampleFactor, doChecksum)
if err != nil {
return -1, nil, nil, err
}
Expand Down
91 changes: 50 additions & 41 deletions pkg/storage/indexheader/stream_binary_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ func NewStreamBinaryReaderMetrics(reg prometheus.Registerer) *StreamBinaryReader
// 2. Reading only the Symbols table from the v1 index-header file format on disk,
// and reading the Postings Offset table directly from the full block index in the bucket.
type StreamBinaryReader struct {
indexHeaderVersion int
symbolsTOC *TOCCompat
postingsOffsetsTOC *TOCCompat

Expand Down Expand Up @@ -133,7 +134,7 @@ func NewStreamBinaryReader(
"path", localIndexHeaderPath, "err", err,
)
start := time.Now()
if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath); err != nil {
if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath, cfg.BucketReader.WriteV2IndexHeader); err != nil {
return nil, fmt.Errorf("failed to write index header: %w", err)
}
level.Info(spanLog).Log(
Expand All @@ -142,65 +143,40 @@ func NewStreamBinaryReader(
)
}

// With full index-header now on disk, initialize the local-disk-backed decbuf factory and read the TOC.
// Initialize the local-disk-backed decbuf factory and read the TOC.
// The index-header's TOC is also required to build the sparse index-header if one does not already exist.
filePoolDecbufFactory := streamencoding.NewFilePoolDecbufFactory(
localIndexHeaderPath, cfg.MaxIdleFileHandles, metrics.filePool,
)
indexHeaderTOC, _, err := TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l)
if err != nil {
// TOC read checks CRC32; assume a failure here is file corruption and attempt to recreate the index-header.
indexHeaderTOC, indexHeaderVersion, err := TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l)
// If we can't read the index header, or it's an unsupported version for the current config, we need to rebuild it
needPostingsOffsetsOnDisk := !cfg.BucketReader.Enabled
if err != nil || (indexHeaderVersion == BinaryFormatV2 && needPostingsOffsetsOnDisk) {
// TOC read checks CRC32; assume a failure here is either due to a file corruption.
level.Debug(spanLog).Log(
"msg", "failed to read table of contents from index-header on disk; will recreate from bucket block index",
"path", localIndexHeaderPath, "err", err,
"path", localIndexHeaderPath, "indexHeaderVersion", indexHeaderVersion, "err", err,
)
start := time.Now()
if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath); err != nil {
if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath, cfg.BucketReader.WriteV2IndexHeader); err != nil {
return nil, fmt.Errorf("failed to write index header: %w", err)
}
level.Info(spanLog).Log(
"msg", "created index-header on local disk from bucket block index",
"path", localIndexHeaderPath, "elapsed", time.Since(start),
)
indexHeaderTOC, _, err = TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l)
if err != nil {
// Failure after recreating index-header from bucket; assume this is unrecoverable.
return nil, fmt.Errorf("failed to read table of contents from index-header on disk after recreate from bucket block index: %w", err)
}
}

// Full index-header is now on disk.
// If we previously failed to load the sparse index-header, build it now from full header.
if !sparseHeaderLoaded {
start := time.Now()
allSymbolsCount, sparseSymbolsOffsets, sparsePostingsOffsets, err = buildInMemorySparseHeaderFromIndexHeader(
ctx, indexHeaderTOC, filePoolDecbufFactory, sparseSampleFactor, cfg.VerifyOnLoad, l,
)
if err != nil {
// Exhausted all options to load sparse index-header to memory. Not recoverable.
return nil, fmt.Errorf("cannot build sparse index-header values from full index-header: %w", err)
}

level.Info(spanLog).Log("msg", "built sparse index-header values from full index-header",
"elapsed", time.Since(start),
)

// Try to write to disk so we do not have to repeat this all again.
sparseHeaderProto := &indexheaderpb.Sparse{
Symbols: streamindex.SparseSymbolsToProto(allSymbolsCount, sparseSymbolsOffsets),
PostingsOffsetTable: streamindex.SparsePostingsOffsetsTableToProto(sparsePostingsOffsets, sparseSampleFactor),
}
if err = writeSparseHeaderProtoToDisk(localSparseHeaderPath, sparseHeaderProto, l); err != nil {
// Log an error in case there are disk issues, but we can still continue.
level.Error(spanLog).Log(
"msg", "failed to write bucket sparse index-header to disk", "err", err,
)
}
indexHeaderTOC, indexHeaderVersion, err = TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l)
}
if err != nil {
// Failure after recreating index-header from bucket; assume this is unrecoverable.
return nil, fmt.Errorf("failed to read table of contents from index-header on disk after recreate from bucket block index: %w", err)
}

// Everything is now loaded from bucket or disk.
streamBinaryReader := &StreamBinaryReader{
sparseSampleFactor: sparseSampleFactor,
indexHeaderVersion: indexHeaderVersion,
}

// Set up each of the Symbols table and Postings Offsets table readers
Expand Down Expand Up @@ -236,6 +212,39 @@ func NewStreamBinaryReader(
streamBinaryReader.postingsOffsetsTOC = indexHeaderTOC
}

// Required index-header section(s) are now on disk.
// If we previously failed to load the sparse index-header, build it now from full header.
// If the bucket reader is enabled, the postings offsets sparse index-header is built from the index header in the bucket.
if !sparseHeaderLoaded {
start := time.Now()
allSymbolsCount, sparseSymbolsOffsets, sparsePostingsOffsets, err = buildInMemorySparseHeaderFromIndexHeader(
ctx,
sectionSource{streamBinaryReader.symbolsTOC, streamBinaryReader.symbolsDecbufFactory},
sectionSource{streamBinaryReader.postingsOffsetsTOC, streamBinaryReader.postingsOffsetsDecbufFactory},
sparseSampleFactor, cfg.VerifyOnLoad, l,
)
if err != nil {
// Exhausted all options to load sparse index-header to memory. Not recoverable.
return nil, fmt.Errorf("cannot build sparse index-header values from full index-header: %w", err)
}

level.Info(spanLog).Log("msg", "built sparse index-header values from full index-header",
"elapsed", time.Since(start),
)

// Try to write to disk so we do not have to repeat this all again.
sparseHeaderProto := &indexheaderpb.Sparse{
Symbols: streamindex.SparseSymbolsToProto(allSymbolsCount, sparseSymbolsOffsets),
PostingsOffsetTable: streamindex.SparsePostingsOffsetsTableToProto(sparsePostingsOffsets, sparseSampleFactor),
}
if err = writeSparseHeaderProtoToDisk(localSparseHeaderPath, sparseHeaderProto, l); err != nil {
// Log an error in case there are disk issues, but we can still continue.
level.Error(spanLog).Log(
"msg", "failed to write bucket sparse index-header to disk", "err", err,
)
}
}

// DecbufFactory and TOC for each section are now assigned according to their configured sources.
// Finally, initialize the readers for each index-header section.
streamBinaryReader.postingsOffsetTable, err = streamindex.NewPostingsOffsetsTableReader(
Expand Down Expand Up @@ -278,7 +287,7 @@ func (r *StreamBinaryReader) IndexVersion(context.Context) (int, error) {
}

func (r *StreamBinaryReader) IndexHeaderVersion() int {
return BinaryFormatV1
return r.indexHeaderVersion
}

func (r *StreamBinaryReader) PostingsOffset(ctx context.Context, name string, value string) (rng index.Range, returnErr error) {
Expand Down
Loading
Loading