From c7e86662a0dc7ba321a5a046c2723d672d2dba37 Mon Sep 17 00:00:00 2001 From: Daniil Forshev Date: Tue, 23 Jun 2026 16:54:04 +0500 Subject: [PATCH] refactor: add ReaderProvider to fetch registry from cache once per request --- cmd/distribution/main.go | 7 ++- frac/remote.go | 71 ++++++++++++++---------- frac/sealed.go | 71 ++++++++++++++---------- frac/sealed/token/table_loader.go | 5 +- frac/sealed_loader.go | 16 +----- frac/sealed_source.go | 13 +++-- storage/index_reader.go | 89 +++++++++++++++++++------------ 7 files changed, 159 insertions(+), 113 deletions(-) diff --git a/cmd/distribution/main.go b/cmd/distribution/main.go index 527c484ff..371e7bba7 100644 --- a/cmd/distribution/main.go +++ b/cmd/distribution/main.go @@ -48,7 +48,12 @@ func getReader(path string) (storage.IndexReader, *os.File) { if err != nil { panic(err) } - return storage.NewIndexReader(readLimiter, f.Name(), f, c), f + readerProvider := storage.NewReaderProvider(readLimiter, f.Name(), f, c) + reader, err := readerProvider.GetReader() + if err != nil { + panic(err) + } + return reader, f } func readBlock(reader storage.IndexReader, blockIndex uint32) ([]byte, error) { diff --git a/frac/remote.go b/frac/remote.go index d3fc15ac9..d118d41cb 100644 --- a/frac/remote.go +++ b/frac/remote.go @@ -44,9 +44,9 @@ type Remote struct { docsReader storage.DocsReader // IsLegacy is true for fractions that use the old single .index file format. - IsLegacy bool - legacyFile storage.ImmutableFile - legacyReader storage.IndexReader + IsLegacy bool + legacyFile storage.ImmutableFile + legacyReaderProvider *storage.ReaderProvider // Per-section index files and their readers (new split format only). infoFile storage.ImmutableFile @@ -55,10 +55,10 @@ type Remote struct { idFile storage.ImmutableFile lidFile storage.ImmutableFile - tokenReader storage.IndexReader - offsetsReader storage.IndexReader - idReader storage.IndexReader - lidReader storage.IndexReader + tokenReaderProvider *storage.ReaderProvider + offsetsReaderProvider *storage.ReaderProvider + idReaderProvider *storage.ReaderProvider + lidReaderProvider *storage.ReaderProvider indexCache *IndexCache @@ -161,6 +161,14 @@ func (f *Remote) FindLIDs(ctx context.Context, ids []seq.ID) ([]seq.LID, error) return dp.FindLIDs(ids) } +func (f *Remote) mustGetReader(p *storage.ReaderProvider) storage.IndexReader { + r, err := p.GetReader() + if err != nil { + logger.Fatal("error creating IndexReader", zap.String("fraction", f.BaseFileName), zap.Error(err)) + } + return r +} + func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, error) { if err := f.init(); err != nil { logger.Error( @@ -171,14 +179,21 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e return nil, err } - tokenReader := &f.tokenReader - lidReader := &f.lidReader - idReader := &f.idReader + var ( + tokenReader storage.IndexReader + lidReader storage.IndexReader + idReader storage.IndexReader + ) if f.IsLegacy { - tokenReader = &f.legacyReader - lidReader = &f.legacyReader - idReader = &f.legacyReader + legacyReader := f.mustGetReader(f.legacyReaderProvider) + tokenReader = legacyReader + lidReader = legacyReader + idReader = legacyReader + } else { + tokenReader = f.mustGetReader(f.tokenReaderProvider) + lidReader = f.mustGetReader(f.lidReaderProvider) + idReader = f.mustGetReader(f.idReaderProvider) } return &sealedDataProvider{ @@ -190,13 +205,13 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e docsReader: &f.docsReader, blocksOffsets: f.blocksData.BlocksOffsets, lidsTable: f.blocksData.LIDsTable, - lidsLoader: lids.NewLoader(f.info.BinaryDataVer, lidReader, f.indexCache.LIDs), - tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, tokenReader, f.indexCache.Tokens), - tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, tokenReader, f.indexCache.TokenTable), + lidsLoader: lids.NewLoader(f.info.BinaryDataVer, &lidReader, f.indexCache.LIDs), + tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, &tokenReader, f.indexCache.Tokens), + tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, &tokenReader, f.indexCache.TokenTable), idsTable: &f.blocksData.IDsTable, idsProvider: seqids.NewProvider( - idReader, + &idReader, f.indexCache.MIDs, f.indexCache.RIDs, f.indexCache.Params, @@ -260,7 +275,7 @@ func (f *Remote) loadInfo() error { if err := f.openInfoLegacy(); err != nil { return err } - if f.info, err = loadInfoLegacy(f.legacyReader); err != nil { + if f.info, err = loadInfoLegacy(f.mustGetReader(f.legacyReaderProvider)); err != nil { logger.Fatal("error loading Info", zap.String("fraction", f.BaseFileName), zap.Error(err)) } return nil @@ -293,16 +308,16 @@ func (f *Remote) init() error { } if f.IsLegacy { - (&LegacyLoader{}).Load(&f.blocksData, f.info, f.legacyReader) + (&LegacyLoader{}).Load(&f.blocksData, f.info, f.mustGetReader(f.legacyReaderProvider)) f.isInited = true return nil } (&Loader{}).Load(&f.blocksData, f.info, IndexReaders{ - Token: f.tokenReader, - Offsets: f.offsetsReader, - ID: f.idReader, - LID: f.lidReader, + Token: f.mustGetReader(f.tokenReaderProvider), + Offsets: f.mustGetReader(f.offsetsReaderProvider), + ID: f.mustGetReader(f.idReaderProvider), + LID: f.mustGetReader(f.lidReaderProvider), }) f.isInited = true @@ -316,7 +331,7 @@ func (f *Remote) openInfoLegacy() error { return f.openRemoteFile(consts.IndexFileSuffix, func(file storage.ImmutableFile) { f.legacyFile = file - f.legacyReader = storage.NewIndexReader( + f.legacyReaderProvider = storage.NewReaderProvider( f.readLimiter, file.Name(), file, f.indexCache.LegacyRegistry, ) @@ -350,7 +365,7 @@ func (f *Remote) openIndex() error { consts.TokenFileSuffix, func(file storage.ImmutableFile) { f.tokenFile = file - f.tokenReader = storage.NewIndexReader( + f.tokenReaderProvider = storage.NewReaderProvider( f.readLimiter, file.Name(), file, f.indexCache.TokenRegistry, ) @@ -365,7 +380,7 @@ func (f *Remote) openIndex() error { consts.OffsetsFileSuffix, func(file storage.ImmutableFile) { f.offsetsFile = file - f.offsetsReader = storage.NewIndexReader( + f.offsetsReaderProvider = storage.NewReaderProvider( f.readLimiter, file.Name(), file, f.indexCache.OffsetsRegistry, ) @@ -380,7 +395,7 @@ func (f *Remote) openIndex() error { consts.IDFileSuffix, func(file storage.ImmutableFile) { f.idFile = file - f.idReader = storage.NewIndexReader( + f.idReaderProvider = storage.NewReaderProvider( f.readLimiter, file.Name(), file, f.indexCache.IDRegistry, ) @@ -395,7 +410,7 @@ func (f *Remote) openIndex() error { consts.LIDFileSuffix, func(file storage.ImmutableFile) { f.lidFile = file - f.lidReader = storage.NewIndexReader( + f.lidReaderProvider = storage.NewReaderProvider( f.readLimiter, file.Name(), file, f.indexCache.LIDRegistry, ) diff --git a/frac/sealed.go b/frac/sealed.go index d90c460ff..28754be91 100644 --- a/frac/sealed.go +++ b/frac/sealed.go @@ -40,9 +40,9 @@ type Sealed struct { docsReader storage.DocsReader // IsLegacy is true for fractions that use the old single .index file format. - IsLegacy bool - legacyFile *os.File - legacyReader storage.IndexReader + IsLegacy bool + legacyFile *os.File + legacyReaderProvider *storage.ReaderProvider // Per-section index files and their readers (new split format only). infoFile *os.File @@ -51,10 +51,10 @@ type Sealed struct { idFile *os.File lidFile *os.File - tokenReader storage.IndexReader - offsetsReader storage.IndexReader - idReader storage.IndexReader - lidReader storage.IndexReader + tokenReaderProvider *storage.ReaderProvider + offsetsReaderProvider *storage.ReaderProvider + idReaderProvider *storage.ReaderProvider + lidReaderProvider *storage.ReaderProvider blocksData sealed.BlocksData indexCache *IndexCache @@ -176,7 +176,7 @@ func (f *Sealed) openInfoLegacy() { } f.legacyFile = file - f.legacyReader = storage.NewIndexReader( + f.legacyReaderProvider = storage.NewReaderProvider( f.readLimiter, file.Name(), file, f.indexCache.LegacyRegistry, ) @@ -216,7 +216,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open token file", zap.String("file", name), zap.Error(err)) } f.tokenFile = file - f.tokenReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.TokenRegistry) + f.tokenReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.TokenRegistry) } if f.offsetsFile == nil { @@ -226,7 +226,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open offsets file", zap.String("file", name), zap.Error(err)) } f.offsetsFile = file - f.offsetsReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.OffsetsRegistry) + f.offsetsReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.OffsetsRegistry) } if f.idFile == nil { @@ -236,7 +236,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open id file", zap.String("file", name), zap.Error(err)) } f.idFile = file - f.idReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.IDRegistry) + f.idReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.IDRegistry) } if f.lidFile == nil { @@ -246,7 +246,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open lid file", zap.String("file", name), zap.Error(err)) } f.lidFile = file - f.lidReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.LIDRegistry) + f.lidReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.LIDRegistry) } } @@ -279,12 +279,20 @@ func (f *Sealed) openDocs() { f.docsReader = storage.NewDocsReader(f.readLimiter, f.docsFile, f.docsCache) } +func (f *Sealed) mustGetReader(p *storage.ReaderProvider) storage.IndexReader { + r, err := p.GetReader() + if err != nil { + logger.Fatal("error creating IndexReader", zap.String("fraction", f.BaseFileName), zap.Error(err)) + } + return r +} + func (f *Sealed) loadInfo() { var err error if f.IsLegacy { f.openInfoLegacy() - if f.info, err = loadInfoLegacy(f.legacyReader); err != nil { + if f.info, err = loadInfoLegacy(f.mustGetReader(f.legacyReaderProvider)); err != nil { logger.Fatal("error loading Info", zap.String("fraction", f.BaseFileName), zap.Error(err)) } return @@ -308,16 +316,16 @@ func (f *Sealed) init(full bool) { } if f.IsLegacy { - (&LegacyLoader{}).Load(&f.blocksData, f.info, f.legacyReader) + (&LegacyLoader{}).Load(&f.blocksData, f.info, f.mustGetReader(f.legacyReaderProvider)) f.isInited = true return } (&Loader{}).Load(&f.blocksData, f.info, IndexReaders{ - Token: f.tokenReader, - Offsets: f.offsetsReader, - ID: f.idReader, - LID: f.lidReader, + Token: f.mustGetReader(f.tokenReaderProvider), + Offsets: f.mustGetReader(f.offsetsReaderProvider), + ID: f.mustGetReader(f.idReaderProvider), + LID: f.mustGetReader(f.lidReaderProvider), }) f.isInited = true @@ -495,14 +503,21 @@ func (f *Sealed) FindLIDs(ctx context.Context, ids []seq.ID) ([]seq.LID, error) func (f *Sealed) createDataProvider(ctx context.Context) *sealedDataProvider { f.init(true) - tokenReader := &f.tokenReader - lidReader := &f.lidReader - idReader := &f.idReader + var ( + tokenReader storage.IndexReader + lidReader storage.IndexReader + idReader storage.IndexReader + ) if f.IsLegacy { - tokenReader = &f.legacyReader - lidReader = &f.legacyReader - idReader = &f.legacyReader + legacyReader := f.mustGetReader(f.legacyReaderProvider) + tokenReader = legacyReader + lidReader = legacyReader + idReader = legacyReader + } else { + tokenReader = f.mustGetReader(f.tokenReaderProvider) + lidReader = f.mustGetReader(f.lidReaderProvider) + idReader = f.mustGetReader(f.idReaderProvider) } return &sealedDataProvider{ @@ -514,13 +529,13 @@ func (f *Sealed) createDataProvider(ctx context.Context) *sealedDataProvider { docsReader: &f.docsReader, blocksOffsets: f.blocksData.BlocksOffsets, lidsTable: f.blocksData.LIDsTable, - lidsLoader: lids.NewLoader(f.info.BinaryDataVer, lidReader, f.indexCache.LIDs), - tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, tokenReader, f.indexCache.Tokens), - tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, tokenReader, f.indexCache.TokenTable), + lidsLoader: lids.NewLoader(f.info.BinaryDataVer, &lidReader, f.indexCache.LIDs), + tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, &tokenReader, f.indexCache.Tokens), + tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, &tokenReader, f.indexCache.TokenTable), idsTable: &f.blocksData.IDsTable, idsProvider: seqids.NewProvider( - idReader, + &idReader, f.indexCache.MIDs, f.indexCache.RIDs, f.indexCache.Params, diff --git a/frac/sealed/token/table_loader.go b/frac/sealed/token/table_loader.go index beb61d485..f4bf56c76 100644 --- a/frac/sealed/token/table_loader.go +++ b/frac/sealed/token/table_loader.go @@ -141,10 +141,7 @@ func (l *TableLoader) loadBlocksLegacy() ([]TableBlock, error) { } func (l *TableLoader) loadBlocks() ([]TableBlock, error) { - blocksCount, err := l.reader.BlocksCount() - if err != nil { - return nil, err - } + blocksCount := l.reader.BlocksCount() var blocks []TableBlock for blockIndex := l.tableIndex; blockIndex < uint32(blocksCount); blockIndex++ { diff --git a/frac/sealed_loader.go b/frac/sealed_loader.go index 10d95a2d3..17782072a 100644 --- a/frac/sealed_loader.go +++ b/frac/sealed_loader.go @@ -231,13 +231,7 @@ func (l *Loader) loadIDsTable(r storage.IndexReader, info *common.Info) seqids.T IDsTotal: info.DocsTotal + 1, // Increment by one for [seq.SystemID] } - blocksCount, err := r.BlocksCount() - if err != nil { - logger.Fatal( - "cannot get block count", - zap.Error(err), - ) - } + blocksCount := r.BlocksCount() for blockIdx := 0; blockIdx < blocksCount; blockIdx += 3 { header, err := r.GetBlockHeader(uint32(blockIdx)) @@ -269,13 +263,7 @@ func (l *Loader) loadLIDsTable(r storage.IndexReader) (*lids.Table, error) { isContinued []bool ) - blocksCount, err := r.BlocksCount() - if err != nil { - logger.Fatal( - "cannot get block count", - zap.Error(err), - ) - } + blocksCount := r.BlocksCount() for blockIdx := 0; blockIdx < blocksCount; blockIdx++ { header, err := r.GetBlockHeader(uint32(blockIdx)) diff --git a/frac/sealed_source.go b/frac/sealed_source.go index cea4f0904..31f7e9f2e 100644 --- a/frac/sealed_source.go +++ b/frac/sealed_source.go @@ -30,19 +30,24 @@ type SealedSource struct { func NewSealedSource(f *Sealed) *SealedSource { f.init(true) + + idReader := f.mustGetReader(f.idReaderProvider) + lidReader := f.mustGetReader(f.lidReaderProvider) + tokenReader := f.mustGetReader(f.tokenReaderProvider) + return &SealedSource{ f: f, idsProvider: seqids.NewProvider( - &f.idReader, + &idReader, f.indexCache.MIDs, f.indexCache.RIDs, f.indexCache.Params, &f.blocksData.IDsTable, f.info.BinaryDataVer, ), - lidsLoader: lids.NewLoader(f.Info().BinaryDataVer, &f.lidReader, f.indexCache.LIDs), - tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, &f.tokenReader, f.indexCache.Tokens), - tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, &f.tokenReader, f.indexCache.TokenTable), + lidsLoader: lids.NewLoader(f.Info().BinaryDataVer, &lidReader, f.indexCache.LIDs), + tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, &tokenReader, f.indexCache.Tokens), + tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, &tokenReader, f.indexCache.TokenTable), } } diff --git a/storage/index_reader.go b/storage/index_reader.go index c7dee18c9..2cbe71fee 100644 --- a/storage/index_reader.go +++ b/storage/index_reader.go @@ -10,7 +10,7 @@ import ( "github.com/ozontech/seq-db/util" ) -type IndexReader struct { +type ReaderProvider struct { limiter *ReadLimiter reader io.ReaderAt @@ -19,47 +19,41 @@ type IndexReader struct { cache *cache.Cache[[]byte] } -func NewIndexReader( - limiter *ReadLimiter, readerName string, - reader io.ReaderAt, registryCache *cache.Cache[[]byte], -) IndexReader { - return IndexReader{ +func NewReaderProvider( + limiter *ReadLimiter, + readerName string, + reader io.ReaderAt, + registryCache *cache.Cache[[]byte], +) *ReaderProvider { + return &ReaderProvider{ limiter: limiter, - reader: reader, readerName: readerName, + reader: reader, cache: registryCache, } } -func (r *IndexReader) GetBlockHeader(index uint32) (IndexBlockHeader, error) { +func (r *ReaderProvider) GetReader() (IndexReader, error) { registry, err := r.cache.GetWithError(1, func() ([]byte, int, error) { data, err := r.readRegistry() return data, cap(data), err }) if err != nil { - return nil, err - } - - if (uint64(index)+1)*IndexBlockHeaderSize > uint64(len(registry)) { - return nil, fmt.Errorf( - "too large index block in file %s, with index %d, registry size %d", - r.readerName, index, len(registry), - ) + return IndexReader{}, err } - pos := index * IndexBlockHeaderSize - return registry[pos : pos+IndexBlockHeaderSize], nil + return NewIndexReader(r.limiter, r.readerName, r.reader, registry), nil } -func (r *IndexReader) readRegistry() ([]byte, error) { +func (r *ReaderProvider) readRegistry() ([]byte, error) { numBuf := make([]byte, 16) n, err := r.limiter.ReadAt(r.reader, numBuf, 0) if err != nil { - return nil, fmt.Errorf("can't read disk registry, %s", err.Error()) + return nil, fmt.Errorf("can't read disk registry from file %s: %s", r.readerName, err.Error()) } if n == 0 { - return nil, fmt.Errorf("can't read disk registry, n=0") + return nil, fmt.Errorf("can't read disk registry from file %s, n=0", r.readerName) } pos := binary.LittleEndian.Uint64(numBuf) @@ -68,20 +62,55 @@ func (r *IndexReader) readRegistry() ([]byte, error) { n, err = r.limiter.ReadAt(r.reader, buf, int64(pos)) if err != nil && err != io.EOF { - return nil, fmt.Errorf("can't read disk registry, %s", err.Error()) + return nil, fmt.Errorf("can't read disk registry from file %s: %s", r.readerName, err.Error()) } if uint64(n) != l { - return nil, fmt.Errorf("can't read disk registry, read=%d, requested=%d", n, l) + return nil, fmt.Errorf("can't read disk registry, read=%d, requested=%d in file %s", n, l, r.readerName) } if len(buf)%IndexBlockHeaderSize != 0 { - return nil, fmt.Errorf("wrong registry format") + return nil, fmt.Errorf("wrong registry format in file %s", r.readerName) } return buf, nil } +type IndexReader struct { + limiter *ReadLimiter + + reader io.ReaderAt + readerName string + + registry []byte +} + +func NewIndexReader( + limiter *ReadLimiter, + readerName string, + reader io.ReaderAt, + registry []byte, +) IndexReader { + return IndexReader{ + limiter: limiter, + reader: reader, + readerName: readerName, + registry: registry, + } +} + +func (r *IndexReader) GetBlockHeader(index uint32) (IndexBlockHeader, error) { + if (uint64(index)+1)*IndexBlockHeaderSize > uint64(len(r.registry)) { + return nil, fmt.Errorf( + "too large index block in file %s, with index %d, registry size %d", + r.readerName, index, len(r.registry), + ) + } + + pos := index * IndexBlockHeaderSize + return r.registry[pos : pos+IndexBlockHeaderSize], nil +} + func (r *IndexReader) ReadIndexBlock(blockIndex uint32, dst []byte) ([]byte, uint64, error) { header, err := r.GetBlockHeader(blockIndex) if err != nil { @@ -108,14 +137,6 @@ func (r *IndexReader) ReadIndexBlock(blockIndex uint32, dst []byte) ([]byte, uin return dst, uint64(n), err } -func (r *IndexReader) BlocksCount() (int, error) { - registry, err := r.cache.GetWithError(1, func() ([]byte, int, error) { - data, err := r.readRegistry() - return data, cap(data), err - }) - if err != nil { - return 0, err - } - - return len(registry) / IndexBlockHeaderSize, nil +func (r *IndexReader) BlocksCount() int { + return len(r.registry) / IndexBlockHeaderSize }