mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 03:47:51 +00:00
remap storage on reopen, instead of remarshalling it
When we do a snapshot, we may end up with containers which are mmapped to the old file, and containers which have allocated storage identical to the contents of the new file. It would be nicer if they were mapped to it. But unmarshalling the entire file is expensive. Instead, we remap it. (Or, if we couldn't mmap it, just make sure the old stuff is no longer using the old storage space before we munmap it.)
This commit is contained in:
parent
565288f6c2
commit
b369dace69
4 changed files with 221 additions and 39 deletions
174
fragment.go
174
fragment.go
|
|
@ -177,7 +177,7 @@ func (f *fragment) Open() error {
|
|||
if err := func() error {
|
||||
// Initialize storage in a function so we can close if anything goes wrong.
|
||||
f.Logger.Debugf("open storage for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
|
||||
if err := f.openStorage(); err != nil {
|
||||
if err := f.openStorage(true); err != nil {
|
||||
return errors.Wrap(err, "opening storage")
|
||||
}
|
||||
|
||||
|
|
@ -215,12 +215,43 @@ func (f *fragment) reopen() (mustClose bool, err error) {
|
|||
return mustClose, nil
|
||||
}
|
||||
|
||||
// openStorage opens the storage bitmap.
|
||||
func (f *fragment) openStorage() error {
|
||||
// openStorage opens the storage bitmap. Usually you also want to read in
|
||||
// the storage, but in the case where we just wrote that file, such as
|
||||
// unprotectedWriteToFragment, we could also just... not. If we didn't
|
||||
// have existing storage, we probably need to unmarshal the data. If the
|
||||
// file we're asked to open is empty, we probably don't.
|
||||
//
|
||||
// If we already had mapped storage previously, we want to unmap that, and
|
||||
// possibly remap it from the file, but we don't need a full unmarshal, just
|
||||
// an update of mapped pointers.
|
||||
//
|
||||
// unmarshalData is somewhat overloaded. it tells us whether or not we
|
||||
// need to actually create a bitmap from the data (if the data exists to
|
||||
// do this from).
|
||||
//
|
||||
// usually unmarshalData is only set to false when we're in the middle of
|
||||
// a snapshot, and unprotectedWriteToFragment just wrote the in-memory data
|
||||
// out.
|
||||
//
|
||||
// If we have existing storage data, and we successfully get new data,
|
||||
// we will unmap the existing storage data.
|
||||
//
|
||||
// This function's design is probably a problem -- it is trying to handle
|
||||
// both cases where there was existing data before, and cases where we
|
||||
// just wrote the data.
|
||||
func (f *fragment) openStorage(unmarshalData bool) error {
|
||||
oldStorageData := f.storageData
|
||||
// there's a few places where we might encounter an error, but need
|
||||
// to continue past it through other error checks, before returning it.
|
||||
var lastError error
|
||||
|
||||
// Create a roaring bitmap to serve as storage for the shard.
|
||||
if f.storage == nil {
|
||||
f.storage = roaring.NewFileBitmap()
|
||||
f.storage.Flags = f.flags
|
||||
// if we didn't actually have storage, we *do* need to
|
||||
// unmarshal this data in order to have any.
|
||||
unmarshalData = true
|
||||
}
|
||||
// Open the data file to be mmap'd and used as an ops log.
|
||||
file, mustClose, err := syswrap.OpenFile(f.path, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666)
|
||||
|
|
@ -237,13 +268,24 @@ func (f *fragment) openStorage() error {
|
|||
return fmt.Errorf("flock: %s", err)
|
||||
}
|
||||
|
||||
// data is the data we would unmarshal from, if we're unmarshalling; it might
|
||||
// be obtained by calling ReadAll on a file.
|
||||
//
|
||||
// newStorageData is the data we should map things to. it is set only if
|
||||
// mmapped; if we didn't mmap (say, we couldn't), we won't want to unmap
|
||||
// the ioutil byte slice. (Theoretically, we shouldn't be using the mapped
|
||||
// flag in that case...)
|
||||
var data []byte
|
||||
var newStorageData []byte
|
||||
|
||||
// If the file is empty then initialize it with an empty bitmap.
|
||||
fi, err := f.file.Stat()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "statting file before")
|
||||
} else if fi.Size() == 0 {
|
||||
bi := bufio.NewWriter(f.file)
|
||||
if _, err := f.storage.WriteTo(bi); err != nil {
|
||||
var err error
|
||||
if _, err = f.storage.WriteTo(bi); err != nil {
|
||||
return fmt.Errorf("init storage file: %s", err)
|
||||
}
|
||||
bi.Flush()
|
||||
|
|
@ -251,39 +293,94 @@ func (f *fragment) openStorage() error {
|
|||
if err != nil {
|
||||
return errors.Wrap(err, "statting file after")
|
||||
}
|
||||
// there's nothing here, we're not going to try to unmarshal it.
|
||||
unmarshalData = false
|
||||
f.rowCache = &simpleCache{make(map[uint64]*Row)}
|
||||
} else {
|
||||
// Mmap the underlying file so it can be zero copied.
|
||||
data, err := syswrap.Mmap(int(f.file.Fd()), 0, int(fi.Size()), syscall.PROT_READ, syscall.MAP_SHARED)
|
||||
data, err = syswrap.Mmap(int(f.file.Fd()), 0, int(fi.Size()), syscall.PROT_READ, syscall.MAP_SHARED)
|
||||
if err == syswrap.ErrMaxMapCountReached {
|
||||
f.Logger.Debugf("maximum number of maps reached, reading file instead")
|
||||
data, err = ioutil.ReadAll(file)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "failure file readall")
|
||||
if unmarshalData {
|
||||
data, err = ioutil.ReadAll(file)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "failure file readall")
|
||||
}
|
||||
}
|
||||
} else if err != nil {
|
||||
return errors.Wrap(err, "mmap failed")
|
||||
} else {
|
||||
f.storageData = data
|
||||
// Advise the kernel that the mmap is accessed randomly.
|
||||
if err := madvise(f.storageData, syscall.MADV_RANDOM); err != nil {
|
||||
return fmt.Errorf("madvise: %s", err)
|
||||
}
|
||||
newStorageData = data
|
||||
}
|
||||
|
||||
if err := f.storage.UnmarshalBinary(data); err != nil {
|
||||
return fmt.Errorf("unmarshal storage: file=%s, err=%s", f.file.Name(), err)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
f.opN = f.storage.Info().OpN
|
||||
if unmarshalData {
|
||||
f.storageData = newStorageData
|
||||
// We're about to either re-read the bitmap, or fail to do so
|
||||
// and unconditionally unmap the existing stuff. Either way, we
|
||||
// want to unmap the old storage data after we're done here, but
|
||||
// we can't unmap it yet because it's still live until sometime
|
||||
// later, but we can't unmap it later, because we could return
|
||||
// early... this is what defer is for.
|
||||
if oldStorageData != nil {
|
||||
defer func() {
|
||||
unmapErr := syswrap.Munmap(oldStorageData)
|
||||
if unmapErr != nil {
|
||||
f.Logger.Printf("unmap of old storage failed: %s", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
// so we have a problem here: if this fails, it's unclear whether
|
||||
// *either* or *both* of old and new storage data might be in use.
|
||||
// So we call the thing that should unconditionally unmap both of them...
|
||||
if err := f.storage.UnmarshalBinary(data); err != nil {
|
||||
f.storage.RemapRoaringStorage(nil)
|
||||
return fmt.Errorf("unmarshal storage: file=%s, err=%s", f.file.Name(), err)
|
||||
}
|
||||
f.rowCache = &simpleCache{make(map[uint64]*Row)}
|
||||
f.opN = f.storage.OpN()
|
||||
// slightly incorrect, but we can't measure it so...
|
||||
} else {
|
||||
// we're moving to new storage, so instead of using the OpN
|
||||
// derived from reading that storage, we notify the bitmap that
|
||||
// OpN is now effectively zero.
|
||||
f.opN = 0
|
||||
f.storage.SetOpN(0)
|
||||
// if oldStorageData is nil, this just tries to unmap any bits that
|
||||
// are currently mapped. otherwise, it will point them at this
|
||||
// storage (if the containers match).
|
||||
var mappedAny bool
|
||||
mappedAny, lastError = f.storage.RemapRoaringStorage(newStorageData)
|
||||
if oldStorageData != nil {
|
||||
unmapErr := syswrap.Munmap(oldStorageData)
|
||||
if unmapErr != nil {
|
||||
f.Logger.Printf("unmap of old storage failed: %s", err)
|
||||
}
|
||||
}
|
||||
if mappedAny {
|
||||
// Advise the kernel that the mmap is accessed randomly.
|
||||
if err := madvise(newStorageData, syscall.MADV_RANDOM); err != nil {
|
||||
lastError = fmt.Errorf("madvise: %s", err)
|
||||
}
|
||||
} else {
|
||||
// if we did map data, but for some reason none of it got used
|
||||
// as backing store, we can unmap it, and set the slice to nil,
|
||||
// so we don't keep the now-invalid slice in f.storageData.
|
||||
if newStorageData != nil {
|
||||
unmapErr := syswrap.Munmap(newStorageData)
|
||||
if unmapErr != nil {
|
||||
lastError = fmt.Errorf("unmapping unused storage data: %s", err)
|
||||
}
|
||||
newStorageData = nil
|
||||
}
|
||||
}
|
||||
f.storageData = newStorageData
|
||||
}
|
||||
|
||||
// Attach the file to the bitmap to act as a write-ahead log.
|
||||
f.storage.OpWriter = f.file
|
||||
f.rowCache = &simpleCache{make(map[uint64]*Row)}
|
||||
|
||||
return nil
|
||||
|
||||
return lastError
|
||||
}
|
||||
|
||||
// openCache initializes the cache from row ids persisted to disk.
|
||||
|
|
@ -343,7 +440,7 @@ func (f *fragment) close() error {
|
|||
}
|
||||
|
||||
// Close underlying storage.
|
||||
if err := f.closeStorage(); err != nil {
|
||||
if err := f.closeStorage(true); err != nil {
|
||||
f.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path)
|
||||
return errors.Wrap(err, "closing storage")
|
||||
}
|
||||
|
|
@ -374,13 +471,19 @@ func (f *fragment) safeClose() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (f *fragment) closeStorage() error {
|
||||
// closeStorage attempts to close storage, including unmapping the old
|
||||
// storage if includeMap is true. This would normally make sense if you're
|
||||
// expecting to be done using the fragment, or to reload it. But it's also
|
||||
// okay to just leave stuff mmapped; you don't have to keep the file
|
||||
// descriptor open. So in some cases, we'll just leave the old mmapping
|
||||
// in place, rather than regenerating everything from the new file.
|
||||
func (f *fragment) closeStorage(includeMap bool) error {
|
||||
// Clear the storage bitmap so it doesn't access the closed mmap.
|
||||
|
||||
//f.storage = roaring.NewBitmap()
|
||||
|
||||
// Unmap the file.
|
||||
if f.storageData != nil {
|
||||
if includeMap && f.storageData != nil {
|
||||
if err := syswrap.Munmap(f.storageData); err != nil {
|
||||
return fmt.Errorf("munmap: %s", err)
|
||||
}
|
||||
|
|
@ -1994,8 +2097,8 @@ func (f *fragment) importValueSmallWrite(columnIDs []uint64, values []int64, bit
|
|||
}
|
||||
return nil
|
||||
}(); err != nil {
|
||||
_ = f.closeStorage()
|
||||
_ = f.openStorage()
|
||||
_ = f.closeStorage(true)
|
||||
_ = f.openStorage(true)
|
||||
return err
|
||||
}
|
||||
rowSet := make(map[uint64]struct{}, bitDepth+1)
|
||||
|
|
@ -2033,8 +2136,8 @@ func (f *fragment) importValue(columnIDs []uint64, values []int64, bitDepth uint
|
|||
}
|
||||
return nil
|
||||
}(); err != nil {
|
||||
_ = f.closeStorage()
|
||||
_ = f.openStorage()
|
||||
_ = f.closeStorage(true)
|
||||
_ = f.openStorage(true)
|
||||
return err
|
||||
}
|
||||
|
||||
|
|
@ -2118,7 +2221,6 @@ func (f *fragment) snapshot() error {
|
|||
// unprotectedWriteToFragment writes the fragment f with bm as the data. It is unprotected, and
|
||||
// f.mu must be locked when calling it.
|
||||
func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err error) { // nolint: interfacer
|
||||
|
||||
completeMessage := fmt.Sprintf("fragment: snapshot complete %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
|
||||
start := time.Now()
|
||||
defer track(start, completeMessage, f.stats, f.Logger)
|
||||
|
|
@ -2142,7 +2244,7 @@ func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err e
|
|||
}
|
||||
|
||||
// Close current storage.
|
||||
if err := f.closeStorage(); err != nil {
|
||||
if err := f.closeStorage(false); err != nil {
|
||||
return n, fmt.Errorf("close storage: %s", err)
|
||||
}
|
||||
|
||||
|
|
@ -2151,8 +2253,12 @@ func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err e
|
|||
return n, fmt.Errorf("rename snapshot: %s", err)
|
||||
}
|
||||
|
||||
// if we reloaded from the file, we'd end up with this bitmap
|
||||
// as our storage. so... let's use this bitmap. as our storage.
|
||||
f.storage = bm
|
||||
|
||||
// Reopen storage.
|
||||
if err := f.openStorage(); err != nil {
|
||||
if err := f.openStorage(false); err != nil {
|
||||
return n, fmt.Errorf("open storage: %s", err)
|
||||
}
|
||||
|
||||
|
|
@ -2341,7 +2447,7 @@ func (f *fragment) readStorageFromArchive(r io.Reader) error {
|
|||
}
|
||||
|
||||
// Close current storage.
|
||||
if err := f.closeStorage(); err != nil {
|
||||
if err := f.closeStorage(true); err != nil {
|
||||
return errors.Wrap(err, "closing")
|
||||
}
|
||||
|
||||
|
|
@ -2351,7 +2457,7 @@ func (f *fragment) readStorageFromArchive(r io.Reader) error {
|
|||
}
|
||||
|
||||
// Reopen storage.
|
||||
if err := f.openStorage(); err != nil {
|
||||
if err := f.openStorage(true); err != nil {
|
||||
return errors.Wrap(err, "opening")
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -117,7 +117,7 @@ func TestFragment_RowcacheMap(t *testing.T) {
|
|||
_, _ = f.setBit(0, uint64(i*32))
|
||||
}
|
||||
// force snapshot so we get a mmapped row...
|
||||
_ = f.snapshot()
|
||||
_ = f.Snapshot()
|
||||
row := f.row(0)
|
||||
segment := row.Segments()[0]
|
||||
bitmap := segment.data
|
||||
|
|
@ -3227,7 +3227,7 @@ func TestImportClearRestart(t *testing.T) {
|
|||
f2.MaxOpN = maxOpN
|
||||
f2.CacheType = f.CacheType
|
||||
|
||||
err = f.closeStorage()
|
||||
err = f.closeStorage(true)
|
||||
if err != nil {
|
||||
t.Fatalf("closing storage: %v", err)
|
||||
}
|
||||
|
|
@ -3261,7 +3261,7 @@ func TestImportClearRestart(t *testing.T) {
|
|||
f3.MaxOpN = maxOpN
|
||||
f3.CacheType = f.CacheType
|
||||
|
||||
err = f2.closeStorage()
|
||||
err = f2.closeStorage(true)
|
||||
if err != nil {
|
||||
t.Fatalf("f2 closing storage: %v", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1240,6 +1240,82 @@ func (r *roaringIterator) Current() (key uint64, cType byte, n int, length int,
|
|||
return r.currentKey, r.currentType, r.currentN, r.currentLen, r.currentPointer, r.lastErr
|
||||
}
|
||||
|
||||
// RemapRoaringStorage tries to update all containers to refer to
|
||||
// the roaring bitmap in the provided []byte. If any containers are
|
||||
// marked as mapped, but do not match the provided storage, they will
|
||||
// be unmapped. The boolean return indicates whether or not any
|
||||
// containers were mapped to the given storage.
|
||||
//
|
||||
// Regardless, after this function runs, no containers have
|
||||
// mapped storage which does not refer to data; either they got mapped
|
||||
// to the new storage, or storage was allocated for them.
|
||||
//
|
||||
// Data should be in the Pilosa roaring format.
|
||||
func (b *Bitmap) RemapRoaringStorage(data []byte) (mappedAny bool, returnErr error) {
|
||||
if b.Containers == nil {
|
||||
return false, nil
|
||||
}
|
||||
var itr *roaringIterator
|
||||
var err error
|
||||
var itrKey uint64
|
||||
var itrCType byte
|
||||
var itrN int
|
||||
var itrPointer *uint16
|
||||
var itrErr error
|
||||
|
||||
if data != nil {
|
||||
itr, err = newRoaringIterator(data)
|
||||
}
|
||||
// don't return early: we still have to do the unmapping
|
||||
if err != nil {
|
||||
returnErr = err
|
||||
}
|
||||
|
||||
if itr != nil {
|
||||
itrKey, itrCType, itrN, _, itrPointer, itrErr = itr.Next()
|
||||
}
|
||||
if itrErr != nil {
|
||||
// iterator errored out, so we won't check it in the loop below
|
||||
itr = nil
|
||||
}
|
||||
|
||||
b.Containers.UpdateEvery(func(key uint64, oldC *Container, existed bool) (newC *Container, write bool) {
|
||||
if itr != nil {
|
||||
for itrKey < key && itrErr == nil {
|
||||
itrKey, itrCType, itrN, _, itrPointer, itrErr = itr.Next()
|
||||
}
|
||||
if itrErr != nil {
|
||||
itr = nil
|
||||
}
|
||||
// container might be similar enough that we should trust it:
|
||||
if itrKey == key && itrCType == oldC.typ() && itrN == int(oldC.N()) {
|
||||
if oldC.frozen() {
|
||||
// we don't use Clone, because that would copy the
|
||||
// storage, and we don't need that.
|
||||
var halfCopy Container
|
||||
halfCopy = *oldC
|
||||
halfCopy.flags &^= flagFrozen
|
||||
newC = &halfCopy
|
||||
} else {
|
||||
newC = oldC
|
||||
}
|
||||
mappedAny = true
|
||||
newC.pointer = itrPointer
|
||||
newC.flags |= flagMapped
|
||||
return newC, true
|
||||
}
|
||||
}
|
||||
// if the container isn't mapped, we don't need to do anything
|
||||
if !oldC.Mapped() {
|
||||
return oldC, false
|
||||
}
|
||||
// forcibly unmap it, so the old mapping can be unmapped safely.
|
||||
newC = oldC.unmapOrClone()
|
||||
return newC, true
|
||||
})
|
||||
return mappedAny, returnErr
|
||||
}
|
||||
|
||||
// ImportRoaringBits sets-or-clears bits based on a provided Roaring bitmap.
|
||||
// This should be equivalent to unmarshalling the bitmap, then executing
|
||||
// either `b = Union(b, newB)` or `b = Difference(b, newB)`, but with lower
|
||||
|
|
|
|||
4
view.go
4
view.go
|
|
@ -441,11 +441,11 @@ func upgradeViewBSIv2(v *view, bitDepth uint) (ok bool, _ error) {
|
|||
|
||||
if tmpPath, err := upgradeRoaringBSIv2(frag, bitDepth); err != nil {
|
||||
return ok, errors.Wrap(err, "upgrading bsi v2")
|
||||
} else if err := frag.closeStorage(); err != nil {
|
||||
} else if err := frag.closeStorage(true); err != nil {
|
||||
return ok, errors.Wrap(err, "closing after bsi v2 upgrade")
|
||||
} else if err := os.Rename(tmpPath, frag.path); err != nil {
|
||||
return ok, errors.Wrap(err, "renaming after bsi v2 upgrade")
|
||||
} else if err := frag.openStorage(); err != nil {
|
||||
} else if err := frag.openStorage(true); err != nil {
|
||||
return ok, errors.Wrap(err, "re-opening after bsi v2 upgrade")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue