mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
Merge pull request #1809 from molecula/fb992
[FB-992] Implement RBF Async Checkpoint
This commit is contained in:
commit
06204bcf7b
13 changed files with 576 additions and 134 deletions
|
|
@ -85,7 +85,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) {
|
|||
flags.BoolVar(&srv.Config.Storage.FsyncEnabled, "storage.fsync", true, "enable fsync fully safe flush-to-disk")
|
||||
|
||||
// RowcacheOn
|
||||
flags.BoolVar((&srv.Config.RowcacheOn), "rowcache-on", srv.Config.RowcacheOn, "turn on the rowcache for all backends (may speed some queries)")
|
||||
flags.BoolVar((&srv.Config.RowcacheOn), "rowcache-on", srv.Config.RowcacheOn, "Do not use, permanently disabled. Flag exists for backwards compatibility and will be removed.")
|
||||
|
||||
// RBF specific flags. See pilosa/rbf/cfg/cfg.go for definitions.
|
||||
srv.Config.RBFConfig.DefineFlags(flags)
|
||||
|
|
|
|||
89
executor.go
89
executor.go
|
|
@ -51,6 +51,9 @@ type executor struct {
|
|||
Node *topology.Node
|
||||
Cluster *cluster
|
||||
|
||||
// how many jobs the work queue has seen
|
||||
workCounter uint64
|
||||
|
||||
// Client used for remote requests.
|
||||
client InternalQueryClient
|
||||
|
||||
|
|
@ -61,6 +64,7 @@ type executor struct {
|
|||
workMu sync.RWMutex
|
||||
workersWG sync.WaitGroup
|
||||
workerPoolSize int
|
||||
currentWorkers int64
|
||||
work chan job
|
||||
|
||||
// Maximum per-request memory usage (Extract() only)
|
||||
|
|
@ -128,15 +132,61 @@ func newExecutor(opts ...executorOption) *executor {
|
|||
e.work = make(chan job, e.workerPoolSize)
|
||||
_ = testhook.Opened(NewAuditor(), e, nil)
|
||||
for i := 0; i < e.workerPoolSize; i++ {
|
||||
e.workersWG.Add(1)
|
||||
go func() {
|
||||
defer e.workersWG.Done()
|
||||
worker(e.work)
|
||||
}()
|
||||
e.addWorker()
|
||||
}
|
||||
go func() {
|
||||
// background task: every so often, check to see whether we have
|
||||
// work in the queue but none has been taken for a while. if so, we
|
||||
// need more workers.
|
||||
prev := atomic.LoadUint64(&e.workCounter)
|
||||
periodic := time.NewTicker(50 * time.Millisecond)
|
||||
defer periodic.Stop()
|
||||
running := true
|
||||
idle := 0
|
||||
for running {
|
||||
<-periodic.C
|
||||
func() {
|
||||
e.workMu.RLock()
|
||||
defer e.workMu.RUnlock()
|
||||
if e.shutdown {
|
||||
running = false
|
||||
return
|
||||
}
|
||||
if len(e.work) == 0 {
|
||||
idle++
|
||||
if idle > 10 && atomic.LoadInt64(&e.currentWorkers) > int64(e.workerPoolSize*2) {
|
||||
select {
|
||||
case e.work <- job{idleHands: true}:
|
||||
// we closed an excess worker
|
||||
default:
|
||||
// somehow between our test above and now the work
|
||||
// queue FILLED UP and we stoically accept this
|
||||
}
|
||||
idle = 0
|
||||
}
|
||||
return
|
||||
}
|
||||
next := atomic.LoadUint64(&e.workCounter)
|
||||
if next == prev {
|
||||
e.addWorker()
|
||||
}
|
||||
prev = next
|
||||
}()
|
||||
}
|
||||
}()
|
||||
return e
|
||||
}
|
||||
|
||||
func (e *executor) addWorker() {
|
||||
e.workersWG.Add(1)
|
||||
atomic.AddInt64(&e.currentWorkers, 1)
|
||||
go func() {
|
||||
defer e.workersWG.Done()
|
||||
e.worker(e.work)
|
||||
atomic.AddInt64(&e.currentWorkers, -1)
|
||||
}()
|
||||
}
|
||||
|
||||
func (e *executor) Close() error {
|
||||
e.workMu.Lock()
|
||||
defer e.workMu.Unlock()
|
||||
|
|
@ -4493,7 +4543,11 @@ func (e *executor) executeRowShard(ctx context.Context, qcx *Qcx, index string,
|
|||
return nil, err
|
||||
}
|
||||
defer finisher(&err0)
|
||||
return frag.row(tx, rowID)
|
||||
row, err := frag.row(tx, rowID)
|
||||
if qcx.write && err == nil {
|
||||
row = row.Clone()
|
||||
}
|
||||
return row, err
|
||||
}
|
||||
|
||||
// If no quantum exists then return an empty bitmap.
|
||||
|
|
@ -4532,15 +4586,21 @@ func (e *executor) executeRowShard(ctx context.Context, qcx *Qcx, index string,
|
|||
if len(rows) == 0 {
|
||||
return &Row{}, nil
|
||||
} else if len(rows) == 1 {
|
||||
if qcx.write {
|
||||
return rows[0].Clone(), nil
|
||||
}
|
||||
return rows[0], nil
|
||||
}
|
||||
row := rows[0].Union(rows[1:]...)
|
||||
if qcx.write {
|
||||
row = row.Clone()
|
||||
}
|
||||
return row, nil
|
||||
|
||||
}
|
||||
|
||||
// executeRowBSIGroupShard executes a range(bsiGroup) call for a local shard.
|
||||
func (e *executor) executeRowBSIGroupShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (_ *Row, err0 error) {
|
||||
func (e *executor) executeRowBSIGroupShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (cloneable *Row, err0 error) {
|
||||
span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeRowBSIGroupShard")
|
||||
defer span.Finish()
|
||||
|
||||
|
|
@ -4572,6 +4632,11 @@ func (e *executor) executeRowBSIGroupShard(ctx context.Context, qcx *Qcx, index
|
|||
return nil, err
|
||||
}
|
||||
defer finisher(&err0)
|
||||
defer func() {
|
||||
if qcx.write && cloneable != nil {
|
||||
cloneable = cloneable.Clone()
|
||||
}
|
||||
}()
|
||||
|
||||
// EQ null _exists - frag.NotNull()
|
||||
// NEQ null frag.NotNull()
|
||||
|
|
@ -4822,6 +4887,9 @@ func (e *executor) executeNotShard(ctx context.Context, qcx *Qcx, index string,
|
|||
if existenceRow, err = existenceFrag.row(tx, 0); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if qcx.write {
|
||||
existenceRow = existenceRow.Clone()
|
||||
}
|
||||
}
|
||||
// the finishers returned by a write tx, which we might be in if there's
|
||||
// a higher-level write in this call OR ANY OTHER CALL, are safe to
|
||||
|
|
@ -5915,10 +5983,15 @@ type job struct {
|
|||
ctx context.Context
|
||||
memoryAvailable *int64 // shared, atomic value
|
||||
resultChan chan mapResponse
|
||||
idleHands bool
|
||||
}
|
||||
|
||||
func worker(work chan job) {
|
||||
func (e *executor) worker(work chan job) {
|
||||
for j := range work {
|
||||
atomic.AddUint64(&e.workCounter, 1)
|
||||
if j.idleHands {
|
||||
return
|
||||
}
|
||||
// Skip out early if the context is done, but still send
|
||||
// an ack so mapperLocal can be sure we aren't about to
|
||||
// work on something it sent us.
|
||||
|
|
|
|||
75
fragment.go
75
fragment.go
|
|
@ -593,8 +593,8 @@ func (f *fragment) mutexCheck(tx Tx, details bool, limit int) (map[uint64][]uint
|
|||
|
||||
// row returns a row by ID.
|
||||
func (f *fragment) row(tx Tx, rowID uint64) (*Row, error) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.mu.RLock()
|
||||
defer f.mu.RUnlock()
|
||||
return f.unprotectedRow(tx, rowID)
|
||||
}
|
||||
|
||||
|
|
@ -937,9 +937,12 @@ func (f *fragment) unprotectedClearRow(tx Tx, rowID uint64) (changed bool, err e
|
|||
return changed, nil
|
||||
}
|
||||
|
||||
// unprotectedClearBlock clears all rows for a given block.
|
||||
// clearBlock clears all rows for a given block.
|
||||
// This updates both the on-disk storage and the in-cache bitmap.
|
||||
func (f *fragment) unprotectedClearBlock(tx Tx, block int) (changed bool, err error) {
|
||||
func (f *fragment) clearBlock(tx Tx, block int) (changed bool, err error) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
firstRow := uint64(block * HashBlockSize)
|
||||
var wp *io.Writer
|
||||
if f.storage != nil {
|
||||
|
|
@ -2708,20 +2711,24 @@ func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDep
|
|||
func (f *fragment) importRoaring(ctx context.Context, tx Tx, data []byte, clear bool) error {
|
||||
span, ctx := tracing.StartSpanFromContext(ctx, "fragment.importRoaring")
|
||||
defer span.Finish()
|
||||
span, ctx = tracing.StartSpanFromContext(ctx, "importRoaring.AcquireFragmentLock")
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
span.Finish()
|
||||
|
||||
return f.unprotectedImportRoaring(ctx, tx, data, clear)
|
||||
rowSet, updateCache, err := f.doImportRoaring(ctx, tx, data, clear)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "doImportRoaring")
|
||||
}
|
||||
if updateCache {
|
||||
return f.updateCachePostImport(ctx, rowSet)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fragment) unprotectedImportRoaring(ctx context.Context, tx Tx, data []byte, clear bool) error {
|
||||
func (f *fragment) doImportRoaring(ctx context.Context, tx Tx, data []byte, clear bool) (map[uint64]int, bool, error) {
|
||||
f.mu.RLock()
|
||||
defer f.mu.RUnlock()
|
||||
rowSize := uint64(1 << shardVsContainerExponent)
|
||||
span, ctx := tracing.StartSpanFromContext(ctx, "importRoaring.ImportRoaringBits")
|
||||
defer span.Finish()
|
||||
|
||||
useRowCache := storage.RowCacheEnabled()
|
||||
var changed int
|
||||
var rowSet map[uint64]int
|
||||
var wp *io.Writer
|
||||
if f.storage != nil {
|
||||
|
|
@ -2734,37 +2741,37 @@ func (f *fragment) unprotectedImportRoaring(ctx context.Context, tx Tx, data []b
|
|||
return err
|
||||
}
|
||||
|
||||
changed, rowSet, err = tx.ImportRoaringBits(f.index(), f.field(), f.view(), f.shard, rit, clear, true, rowSize)
|
||||
_, rowSet, err = tx.ImportRoaringBits(f.index(), f.field(), f.view(), f.shard, rit, clear, true, rowSize)
|
||||
return err
|
||||
})
|
||||
|
||||
span.Finish()
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, false, err
|
||||
}
|
||||
|
||||
updateCache := f.CacheType != CacheTypeNone
|
||||
return rowSet, updateCache, err
|
||||
}
|
||||
|
||||
func (f *fragment) updateCachePostImport(ctx context.Context, rowSet map[uint64]int) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
anyChanged := false
|
||||
|
||||
for rowID, changes := range rowSet {
|
||||
if changes == 0 {
|
||||
continue
|
||||
}
|
||||
if useRowCache && f.rowCache != nil {
|
||||
f.rowCache.Add(rowID, nil)
|
||||
}
|
||||
if updateCache {
|
||||
anyChanged = true
|
||||
if changes < 0 {
|
||||
absChanges := uint64(-1 * changes)
|
||||
if absChanges <= f.cache.Get(rowID) {
|
||||
f.cache.BulkAdd(rowID, f.cache.Get(rowID)-absChanges)
|
||||
} else {
|
||||
f.cache.BulkAdd(rowID, 0)
|
||||
}
|
||||
anyChanged = true
|
||||
if changes < 0 {
|
||||
absChanges := uint64(-1 * changes)
|
||||
if absChanges <= f.cache.Get(rowID) {
|
||||
f.cache.BulkAdd(rowID, f.cache.Get(rowID)-absChanges)
|
||||
} else {
|
||||
f.cache.BulkAdd(rowID, f.cache.Get(rowID)+uint64(changes))
|
||||
f.cache.BulkAdd(rowID, 0)
|
||||
}
|
||||
} else {
|
||||
f.cache.BulkAdd(rowID, f.cache.Get(rowID)+uint64(changes))
|
||||
}
|
||||
}
|
||||
// we only set this if we need to update the cache
|
||||
|
|
@ -2772,26 +2779,18 @@ func (f *fragment) unprotectedImportRoaring(ctx context.Context, tx Tx, data []b
|
|||
f.cache.Invalidate()
|
||||
}
|
||||
|
||||
span, _ = tracing.StartSpanFromContext(ctx, "importRoaring.incrementOpN")
|
||||
|
||||
f.incrementOpN(changed)
|
||||
|
||||
span.Finish()
|
||||
return nil
|
||||
}
|
||||
|
||||
// importRoaringOverwrite overwrites the specified block with the provided data.
|
||||
func (f *fragment) importRoaringOverwrite(ctx context.Context, tx Tx, data []byte, block int) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
// Clear the existing data from fragment block.
|
||||
if _, err := f.unprotectedClearBlock(tx, block); err != nil {
|
||||
if _, err := f.clearBlock(tx, block); err != nil {
|
||||
return errors.Wrapf(err, "clearing block: %d", block)
|
||||
}
|
||||
|
||||
// Union the new block data with the fragment data.
|
||||
return f.unprotectedImportRoaring(ctx, tx, data, false)
|
||||
return f.importRoaring(ctx, tx, data, false)
|
||||
}
|
||||
|
||||
// incrementOpN increase the operation count by one.
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@
|
|||
package cfg
|
||||
|
||||
import (
|
||||
"github.com/molecula/featurebase/v2/logger"
|
||||
"github.com/spf13/pflag"
|
||||
)
|
||||
|
||||
|
|
@ -35,6 +36,11 @@ type Config struct {
|
|||
// CursorCacheSize is the number of copies of Cursor{} to keep in our
|
||||
// readyCursorCh arena to avoid GC pressure.
|
||||
CursorCacheSize int64 `toml:"cursor-cache-size"`
|
||||
|
||||
// Logger specifies a logger for asynchronous errors, such as
|
||||
// background checkpoints. It cannot be set from toml. The default is
|
||||
// to use stderr.
|
||||
Logger logger.Logger `toml:"-"`
|
||||
}
|
||||
|
||||
func NewDefaultConfig() *Config {
|
||||
|
|
|
|||
|
|
@ -973,8 +973,8 @@ func TestCursor_SplitBranchCells(t *testing.T) {
|
|||
}
|
||||
//
|
||||
c, _ := tx.Cursor("x") //added just for dot code coverage
|
||||
c.Dump("ignore for coverage")
|
||||
|
||||
c.Dump("test.dump")
|
||||
os.Remove("test.dump")
|
||||
}
|
||||
|
||||
func TestCursor_RemoveCells(t *testing.T) {
|
||||
|
|
|
|||
367
rbf/db.go
367
rbf/db.go
|
|
@ -11,6 +11,7 @@ import (
|
|||
"syscall"
|
||||
|
||||
"github.com/benbjohnson/immutable"
|
||||
"github.com/molecula/featurebase/v2/logger"
|
||||
rbfcfg "github.com/molecula/featurebase/v2/rbf/cfg"
|
||||
"github.com/molecula/featurebase/v2/syswrap"
|
||||
)
|
||||
|
|
@ -27,6 +28,16 @@ var cursorSyncPool = &sync.Pool{
|
|||
},
|
||||
}
|
||||
|
||||
// txWaiter is a representation of "i need to wait for txs to complete".
|
||||
// it is created with a function, and will run that function, with the db
|
||||
// lock held, at some point after every Tx that was open when it was created
|
||||
// has closed. WARNING: A txWaiter may hold db.rwmu.
|
||||
type txWaiter struct {
|
||||
ready chan struct{}
|
||||
waitingOn map[*Tx]struct{}
|
||||
callback func()
|
||||
}
|
||||
|
||||
// DB options like MaxSize, FsyncEnabled, DoAllocZero
|
||||
// can be set before calling DB.Open().
|
||||
type DB struct {
|
||||
|
|
@ -38,15 +49,21 @@ type DB struct {
|
|||
pageMap *PageMap // pgno-to-WALID mapping
|
||||
txs map[*Tx]struct{} // active transactions
|
||||
opened bool // true if open
|
||||
logger logger.Logger // for diagnostics from async things
|
||||
|
||||
wal []byte // wal mmap
|
||||
walFile *os.File // wal file descriptor
|
||||
walPageN int // wal page count
|
||||
wal []byte // wal mmap
|
||||
walFile *os.File // wal file descriptor
|
||||
walPageN int // wal page count
|
||||
baseWALID int64 // WAL ID of first page
|
||||
|
||||
mu sync.RWMutex // general mutex
|
||||
rwmu sync.Mutex // mutex for restricting single writer
|
||||
haltCond *sync.Cond // condition for resuming txs after checkpoint
|
||||
|
||||
txWaiters []*txWaiter // things waiting for Txs to close
|
||||
|
||||
isDead error // this database died in an unrecoverable way, error out opens
|
||||
|
||||
// Path represents the path to the database file.
|
||||
Path string
|
||||
}
|
||||
|
|
@ -62,6 +79,11 @@ func NewDB(path string, cfg *rbfcfg.Config) *DB {
|
|||
txs: make(map[*Tx]struct{}),
|
||||
pageMap: NewPageMap(),
|
||||
Path: path,
|
||||
logger: cfg.Logger,
|
||||
}
|
||||
if db.logger == nil {
|
||||
// default to writing to stdout if not told otherwise
|
||||
db.logger = logger.NewStandardLogger(os.Stderr)
|
||||
}
|
||||
db.haltCond = sync.NewCond(&db.mu)
|
||||
|
||||
|
|
@ -133,8 +155,12 @@ func (db *DB) Open() (err error) {
|
|||
// Open write-ahead log & checkpoint to the end since no transactions are open.
|
||||
if err := db.openWAL(); err != nil {
|
||||
return fmt.Errorf("wal open: %w", err)
|
||||
} else if err := db.checkpoint(); err != nil {
|
||||
return fmt.Errorf("checkpoint: %w", err)
|
||||
} else {
|
||||
// checkpoint wants to hold the rwmu lock.
|
||||
db.rwmu.Lock()
|
||||
if err := db.checkpoint(); err != nil {
|
||||
return fmt.Errorf("startup checkpoint: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
@ -158,10 +184,12 @@ func (db *DB) openWAL() (err error) {
|
|||
|
||||
// Determine the number of whole pages in the WAL.
|
||||
var pageN int
|
||||
var fileSize int64
|
||||
if fi, err := db.walFile.Stat(); err != nil {
|
||||
return fmt.Errorf("wal stat: %w", err)
|
||||
} else {
|
||||
pageN = int(fi.Size() / PageSize)
|
||||
fileSize = fi.Size()
|
||||
pageN = int(fileSize / PageSize)
|
||||
}
|
||||
|
||||
// Read backwards through the WAL to find the last valid meta page.
|
||||
|
|
@ -169,28 +197,88 @@ func (db *DB) openWAL() (err error) {
|
|||
if page, err := db.readWALPageAt(pageN - 1); err != nil {
|
||||
return err
|
||||
} else if IsMetaPage(page) {
|
||||
// We now face a challenge. Probably this is a meta page.
|
||||
// But consider a sequence of pages written which gets
|
||||
// interrupted right before the meta page is written.
|
||||
// If the last page is a bitmap page, it could LOOK LIKE a meta
|
||||
// page. So we have to check the page before it. If that page
|
||||
// is a bitmap header, then actually this is a bitmap page, right?
|
||||
// If that page doesn't exist, of course, we're fine, except
|
||||
// for the philosophical question of why we wrote a meta page
|
||||
// when no pages had changed.
|
||||
if pageN > 1 {
|
||||
if page, err = db.readWALPageAt(pageN - 2); err != nil {
|
||||
return err
|
||||
}
|
||||
if IsBitmapHeader(page) {
|
||||
// But wait!
|
||||
// What if this *is* a meta page, and the page before it is
|
||||
// actually a *bitmap page* that looks like a bitmap header? And
|
||||
// so on.
|
||||
//
|
||||
// Rather than try to resolve this, in this insanely unlikely
|
||||
// situation, we read from the beginning which allows us to
|
||||
// always know what we're seeing, because every bitmap page
|
||||
// comes *after* a bitmap header page, and thus, we know when
|
||||
// we might be seeing one.
|
||||
pageN, err = db.methodicalWALPageN(pageN)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
// Truncate WAL to the last valid meta page.
|
||||
if err := db.walFile.Truncate(int64(pageN * PageSize)); err != nil {
|
||||
return fmt.Errorf("wal truncate: %w", err)
|
||||
} else if _, err := db.walFile.Seek(int64(pageN*PageSize), io.SeekStart); err != nil {
|
||||
if fileSize != int64(pageN*PageSize) {
|
||||
if err := db.walFile.Truncate(int64(pageN * PageSize)); err != nil {
|
||||
return fmt.Errorf("wal truncate: %w", err)
|
||||
}
|
||||
}
|
||||
if _, err := db.walFile.Seek(int64(pageN*PageSize), io.SeekStart); err != nil {
|
||||
return fmt.Errorf("wal seek: %w", err)
|
||||
}
|
||||
db.walPageN = pageN
|
||||
db.baseWALID = readMetaWALID(db.data)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// checkpoint moves all WAL pages to the main DB file.
|
||||
// Must be called by a write transaction while under db.mu lock.
|
||||
func (db *DB) checkpoint() error {
|
||||
// methodicalWALPageN tries to determine the last meta page in a very reliable
|
||||
// but slow way. This handles the theoretical but hard to imagine creating
|
||||
// edge case where we have a bitmap page which happens to look like a meta
|
||||
// page, and the write got interrupted before the meta page got written.
|
||||
func (db *DB) methodicalWALPageN(pageN int) (lastMeta int, err error) {
|
||||
for i := 0; i < pageN; i++ {
|
||||
var page []byte
|
||||
if page, err = db.readWALPageAt(i); err != nil {
|
||||
return -1, err
|
||||
}
|
||||
switch {
|
||||
case IsMetaPage(page):
|
||||
lastMeta = i
|
||||
case IsBitmapHeader(page):
|
||||
// skip the bitmap page, which we can't usefully evaluate
|
||||
i++
|
||||
}
|
||||
}
|
||||
return lastMeta, nil
|
||||
}
|
||||
|
||||
// checkpoint moves all WAL pages to the main DB file. Must be called
|
||||
// while holding both db.mu and db.rwmu. Should release db.rwmu, but not
|
||||
// db.mu.
|
||||
func (db *DB) checkpoint() (err error) {
|
||||
// if we don't spin off a possible async waiter, we should release the
|
||||
// write lock, if we do, that will release it.
|
||||
releaseLock := true
|
||||
defer func() {
|
||||
if releaseLock {
|
||||
db.rwmu.Unlock()
|
||||
}
|
||||
}()
|
||||
if !db.opened {
|
||||
return nil
|
||||
} else if len(db.txs) > 0 {
|
||||
return nil // skip if transactions open
|
||||
}
|
||||
|
||||
// Check if there are any WAL pages, if not do nothing as
|
||||
|
|
@ -199,48 +287,112 @@ func (db *DB) checkpoint() error {
|
|||
if db.walPageN == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
for i := 0; i < db.walPageN; i++ {
|
||||
page, err := db.readWALPageAt(i)
|
||||
if err != nil {
|
||||
return err
|
||||
// wake up things waiting on haltCond when we're done, even if we fail.
|
||||
// Otherwise, we deadlock with them all stuck waiting on that forever.
|
||||
defer func() {
|
||||
if err != nil && db.isDead == nil {
|
||||
db.isDead = err
|
||||
}
|
||||
db.haltCond.Broadcast()
|
||||
}()
|
||||
|
||||
// Determine page number. Meta pages are always on zero & bitmap
|
||||
// headers specify the page number of the next page in the WAL.
|
||||
// All other pages have their page number in the page data.
|
||||
var pgno uint32
|
||||
if IsBitmapHeader(page) {
|
||||
pgno = readPageNo(page)
|
||||
if page, err = db.readWALPageAt(i + 1); err != nil {
|
||||
return err
|
||||
// Copy the pages from the WAL back to the database outside of the lock.
|
||||
if err := func() error {
|
||||
db.mu.Unlock() // This is intentionally reversed so run w/o lock
|
||||
defer db.mu.Lock()
|
||||
|
||||
var page []byte
|
||||
// We might have either a *PageMap or just the file. If we have the file,
|
||||
// building the PageMap is fairly expensive because it's fancy and immutable.
|
||||
// If we have the PageMap *or* some other map, that's two different things
|
||||
// to iterate. If we have the PageMap, building a map from it is relatively
|
||||
// cheap, so we'll do it that way.
|
||||
pages := make(map[uint32]int)
|
||||
|
||||
if db.pageMap.size == 0 {
|
||||
// you'd think we're done, but actually this PROBABLY means that
|
||||
// this is initial startup, and we haven't read the file yet. We scan
|
||||
// the file for pages, because it turns out most of them probably
|
||||
// got overwritten.
|
||||
for i := 0; i < db.walPageN; i++ {
|
||||
page, err = db.readWALPageAt(i)
|
||||
if err != nil {
|
||||
return fmt.Errorf("reading WAL page %d: %w", i, err)
|
||||
}
|
||||
|
||||
// Determine page number. Meta pages are always on zero & bitmap
|
||||
// headers specify the page number of the next page in the WAL.
|
||||
// All other pages have their page number in the page data.
|
||||
var pgno uint32
|
||||
if IsBitmapHeader(page) {
|
||||
pgno = readPageNo(page)
|
||||
if i+1 < db.walPageN {
|
||||
if page, err = db.readWALPageAt(i + 1); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
return fmt.Errorf("last page of WAL file (%d) is bitmap header", i)
|
||||
}
|
||||
i++ // bitmaps in WAL are two pages
|
||||
} else if !IsMetaPage(page) {
|
||||
pgno = readPageNo(page)
|
||||
}
|
||||
// record where in the file we have this page
|
||||
pages[pgno] = i
|
||||
}
|
||||
} else {
|
||||
itr := db.pageMap.Iterator()
|
||||
itr.First()
|
||||
for k, v, ok := itr.Next(); ok; k, v, ok = itr.Next() {
|
||||
pages[k] = int(v - db.baseWALID - 1)
|
||||
}
|
||||
i++ // bitmaps in WAL are two pages
|
||||
} else if !IsMetaPage(page) {
|
||||
pgno = readPageNo(page)
|
||||
}
|
||||
|
||||
// Write data to the data file.
|
||||
if err := db.writeDBPage(pgno, page); err != nil {
|
||||
return err
|
||||
// fmt.Printf("checkpoint: walPageN %d, PageMap size %d\n", db.walPageN, db.pageMap.size)
|
||||
for pgno, walID := range pages {
|
||||
page, err = db.readWALPageAt(walID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("reading page %d [page number %d]: %v", walID, pgno, err)
|
||||
}
|
||||
|
||||
// Write data to the data file.
|
||||
if err = db.writeDBPage(pgno, page); err != nil {
|
||||
return fmt.Errorf("writing page %d: %v", pgno, err)
|
||||
}
|
||||
}
|
||||
|
||||
// Ensure database file is synced and then truncate the WAL file.
|
||||
if err = db.fsync(db.file); err != nil {
|
||||
return fmt.Errorf("db file sync: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Ensure database file is synced and then truncate the WAL file.
|
||||
if err := db.fsync(db.file); err != nil {
|
||||
return fmt.Errorf("db file sync: %w", err)
|
||||
} else if err := db.walFile.Truncate(0); err != nil {
|
||||
return fmt.Errorf("truncate wal file: %w", err)
|
||||
} else if err := db.fsync(db.walFile); err != nil {
|
||||
return fmt.Errorf("wal file sync: %w", err)
|
||||
} else if _, err := db.walFile.Seek(0, io.SeekStart); err != nil {
|
||||
return fmt.Errorf("seek wal file: %w", err)
|
||||
}
|
||||
// now we've updated the file. There are existing transactions that are still
|
||||
// using the WAL, though. So we wait for them to terminate before we unlock
|
||||
// the rwmu and update the metadata about the WAL.
|
||||
releaseLock = false
|
||||
db.walPageN = 0
|
||||
db.pageMap = NewPageMap()
|
||||
|
||||
// Notify halted transactions that the WAL has been checkpointed.
|
||||
db.haltCond.Broadcast()
|
||||
db.afterCurrentTx(func() {
|
||||
defer db.rwmu.Unlock()
|
||||
db.baseWALID = readMetaWALID(db.data)
|
||||
db.mu.Unlock()
|
||||
defer db.mu.Lock()
|
||||
|
||||
if err = db.walFile.Truncate(0); err != nil {
|
||||
db.logger.Errorf("truncate wal file: %w", err)
|
||||
} else if err = db.fsync(db.walFile); err != nil {
|
||||
db.logger.Errorf("wal file sync: %w", err)
|
||||
} else if _, err = db.walFile.Seek(0, io.SeekStart); err != nil {
|
||||
db.logger.Errorf("seek wal file: %w", err)
|
||||
}
|
||||
|
||||
})
|
||||
|
||||
return nil
|
||||
}
|
||||
|
|
@ -450,10 +602,25 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) {
|
|||
cleanup()
|
||||
return nil, ErrClosed
|
||||
}
|
||||
if db.isDead != nil {
|
||||
err := db.isDead
|
||||
cleanup()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Wait for WAL size to be below threshold.
|
||||
for int64(db.walPageN*PageSize) > db.cfg.MaxWALCheckpointSize {
|
||||
db.haltCond.Wait()
|
||||
// Wait for WAL size to be below threshold, if we're going to write.
|
||||
// Reads don't care.
|
||||
if writable {
|
||||
for int64(db.walPageN*PageSize) > db.cfg.MaxWALCheckpointSize {
|
||||
if db.isDead != nil {
|
||||
err := db.isDead
|
||||
cleanup()
|
||||
return nil, err
|
||||
}
|
||||
// This implicitly releases db.mu.Lock and comes back with it
|
||||
// held again.
|
||||
db.haltCond.Wait()
|
||||
}
|
||||
}
|
||||
|
||||
tx := &Tx{
|
||||
|
|
@ -502,26 +669,95 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) {
|
|||
return tx, nil
|
||||
}
|
||||
|
||||
// removeTx removes an active transaction from the database.
|
||||
func (db *DB) removeTx(tx *Tx) error {
|
||||
// Release writer lock if tx is writable.
|
||||
if tx.writable {
|
||||
tx.db.rwmu.Unlock()
|
||||
// afterCurrentTx produces runs the provided callback, with the db lock
|
||||
// held, after all current Tx terminate. It should be called with the db
|
||||
// lock held.
|
||||
func (db *DB) afterCurrentTx(callback func()) {
|
||||
if len(db.txs) == 0 {
|
||||
callback()
|
||||
return
|
||||
}
|
||||
txw := &txWaiter{}
|
||||
txw.ready = make(chan struct{})
|
||||
txw.callback = callback
|
||||
txw.waitingOn = make(map[*Tx]struct{}, len(db.txs))
|
||||
for k := range db.txs {
|
||||
txw.waitingOn[k] = struct{}{}
|
||||
}
|
||||
db.txWaiters = append(db.txWaiters, txw)
|
||||
go func() {
|
||||
<-txw.ready
|
||||
// fmt.Printf("afterCurrentTx: locking db\n")
|
||||
db.mu.Lock()
|
||||
defer db.mu.Unlock()
|
||||
// fmt.Printf("afterCurrentTx: running callback\n")
|
||||
txw.callback()
|
||||
}()
|
||||
return
|
||||
}
|
||||
|
||||
// removeTx removes an active transaction from the database. it obtains
|
||||
// the db lock, and currently drops it, but will later possibly be leaving
|
||||
// it retained by an asynchronous op that wants to happen before we start
|
||||
// running new tx.
|
||||
func (db *DB) removeTx(tx *Tx) error {
|
||||
// We might want to trigger a checkpoint. Only for writable
|
||||
// transactions, and only when either there's nothing else open or we
|
||||
// really need to.
|
||||
checkpoint := false
|
||||
if tx.writable {
|
||||
walSize := db.walSize()
|
||||
if walSize > db.cfg.MinWALCheckpointSize {
|
||||
// Might be a good time for a checkpoint. We'll do a checkpoint
|
||||
// if we're the only transaction, or if we have to.
|
||||
if len(db.txs) == 1 || walSize > db.cfg.MaxWALCheckpointSize {
|
||||
checkpoint = true
|
||||
}
|
||||
}
|
||||
// During checkpointing, we'll be preventing writes, but allowing reads.
|
||||
if !checkpoint {
|
||||
tx.db.rwmu.Unlock()
|
||||
}
|
||||
}
|
||||
// remove ourselves from the list of transactions the db is keeping.
|
||||
delete(tx.db.txs, tx)
|
||||
for i := 0; i < len(tx.db.txWaiters); i++ {
|
||||
txw := tx.db.txWaiters[i]
|
||||
// in practice this probably never matters, but theoretically the
|
||||
// goroutine that's waiting on the condition variable may
|
||||
// not have performed its first test on len(txw.waitingOn) yet.
|
||||
delete(txw.waitingOn, tx)
|
||||
// let it know we're done. we've still got db.mu.lock, so it won't
|
||||
// happen just yet, but it'll be able to continue.
|
||||
if len(txw.waitingOn) == 0 {
|
||||
// remove us from the db's list
|
||||
copy(db.txWaiters[i:], db.txWaiters[i+1:])
|
||||
db.txWaiters = db.txWaiters[:len(db.txWaiters)-1]
|
||||
close(txw.ready)
|
||||
// decrement i so we don't skip an entry we just copied in to [i]
|
||||
i--
|
||||
}
|
||||
}
|
||||
|
||||
// Disassociate from db.
|
||||
tx.db = nil
|
||||
|
||||
// Write pages from WAL to DB.
|
||||
// TODO(bbj): Move this to an async goroutine.
|
||||
if len(db.txs) == 0 && db.walSize() > db.cfg.MinWALCheckpointSize {
|
||||
if err := db.checkpoint(); err != nil {
|
||||
return fmt.Errorf("checkpoint: %w", err)
|
||||
}
|
||||
if checkpoint {
|
||||
// We need to run a checkpoint. This can be semi-asynchronous.
|
||||
// It needs to wait until every existing transaction has finished,
|
||||
// because every existing transaction could want to look up pages
|
||||
// which are in the database before our operations, but which should
|
||||
// now be in the WAL. We want them to use the WAL instead.
|
||||
// fmt.Printf("possibly-async checkpoint...\n")
|
||||
db.afterCurrentTx(func() {
|
||||
// We still hold db.rwmu here. checkpoint unlocks it when it's
|
||||
// ready.
|
||||
// fmt.Printf("checkpoint starting\n")
|
||||
if err := db.checkpoint(); err != nil {
|
||||
db.logger.Errorf("async checkpoint: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -547,14 +783,9 @@ func (db *DB) readDBPage(pgno uint32) ([]byte, error) {
|
|||
return db.data[offset : offset+PageSize], nil
|
||||
}
|
||||
|
||||
// baseWALID returns the WAL ID stored in the database file meta page.
|
||||
func (db *DB) baseWALID() int64 {
|
||||
return readMetaWALID(db.data)
|
||||
}
|
||||
|
||||
// readWALPageByID reads a WAL page by WAL ID.
|
||||
func (db *DB) readWALPageByID(id int64) ([]byte, error) {
|
||||
return db.readWALPageAt(int(id - db.baseWALID() - 1))
|
||||
return db.readWALPageAt(int(id - db.baseWALID - 1))
|
||||
}
|
||||
|
||||
// readWALPageAt reads the i-th page in the WAL file.
|
||||
|
|
|
|||
120
rbf/db_test.go
120
rbf/db_test.go
|
|
@ -13,6 +13,7 @@ import (
|
|||
|
||||
_ "net/http/pprof"
|
||||
|
||||
"github.com/felixge/fgprof"
|
||||
"github.com/molecula/featurebase/v2/rbf"
|
||||
rbfcfg "github.com/molecula/featurebase/v2/rbf/cfg"
|
||||
"golang.org/x/sync/errgroup"
|
||||
|
|
@ -290,7 +291,8 @@ func TestDB_MultiTx(t *testing.T) {
|
|||
|
||||
time.Sleep(time.Duration(rand.Intn(100)) * time.Millisecond)
|
||||
|
||||
for i := 0; i < rand.Intn(1000); i++ {
|
||||
n := rand.Intn(500) + 500
|
||||
for i := 0; i < n; i++ {
|
||||
v := rand.Intn(1 << 20)
|
||||
if _, err := tx.Contains("x", uint64(v)); err != nil {
|
||||
return err
|
||||
|
|
@ -315,7 +317,8 @@ func TestDB_MultiTx(t *testing.T) {
|
|||
}
|
||||
defer tx.Rollback()
|
||||
|
||||
for j := 0; j < rand.Intn(100); j++ {
|
||||
n := rand.Intn(90) + 10
|
||||
for j := 0; j < n; j++ {
|
||||
v := rand.Intn(1 << 20)
|
||||
if _, err := tx.Add("x", uint64(v)); err != nil {
|
||||
t.Fatal(err)
|
||||
|
|
@ -336,6 +339,119 @@ func TestDB_MultiTx(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
// premake pool of random values
|
||||
const randPool = (1 << 18)
|
||||
|
||||
// benchmarkOneCheckpoint
|
||||
func benchmarkOneCheckpoint(b *testing.B, randInts []int) {
|
||||
cfg := rbfcfg.NewDefaultConfig()
|
||||
// extremely low to force checkpointing
|
||||
cfg.MinWALCheckpointSize = rbf.PageSize * 16
|
||||
cfg.MaxWALCheckpointSize = rbf.PageSize * 64
|
||||
var _ rbfcfg.Config
|
||||
db := MustOpenDB(b, cfg)
|
||||
defer MustCloseDB(b, db)
|
||||
|
||||
// Run multiple readers in separate goroutines.
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
g, ctx := errgroup.WithContext(ctx)
|
||||
for i := 0; i < 8; i++ {
|
||||
i := i
|
||||
g.Go(func() error {
|
||||
for {
|
||||
if ctx.Err() != nil {
|
||||
return nil // cancelled, return no error
|
||||
} else if err := func() error {
|
||||
tx, err := db.Begin(false)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
|
||||
time.Sleep(time.Duration(rand.Intn(int(3 * time.Millisecond))))
|
||||
|
||||
times := rand.Intn(1000) + 1
|
||||
for j := 0; j < times; j++ {
|
||||
v := randInts[((i<<10)+j)%(randPool-1)]
|
||||
if _, err := tx.Contains("x", uint64(v)); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}(); err != nil {
|
||||
return err
|
||||
}
|
||||
// time.Sleep(time.Duration(rand.Intn(int(3 * time.Millisecond))))
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// Continuously set/clear bits while readers are executing.
|
||||
next := 0
|
||||
for i := 0; i < 1000; i++ {
|
||||
func() {
|
||||
tx, err := db.Begin(true)
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
defer tx.Rollback()
|
||||
|
||||
times := rand.Intn(100)
|
||||
for j := 0; j < times; j++ {
|
||||
v := randInts[next]
|
||||
next = (next + 1) % (randPool - 1)
|
||||
if j&7 == 0 {
|
||||
// some removes but they're less frequent
|
||||
if _, err := tx.Remove("x", uint64(v)); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
} else {
|
||||
if _, err := tx.Add("x", uint64(v)); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// Stop readers & wait.
|
||||
cancel()
|
||||
if err := g.Wait(); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkDbCheckpoint(b *testing.B) {
|
||||
out, err := os.Create("cp.out")
|
||||
if err != nil {
|
||||
b.Fatalf("creating log file: %v", err)
|
||||
}
|
||||
done := fgprof.Start(out, fgprof.FormatPprof)
|
||||
b.StopTimer()
|
||||
// premake these because otherwise it's >5% of CPU in the reads
|
||||
randInts := make([]int, randPool)
|
||||
for i := range randInts {
|
||||
v1, v2 := rand.Intn(1<<24), rand.Intn(1<<24)
|
||||
// minimum gives us a skewed distribution which makes lower values more
|
||||
// likely than higher values, so we get a mix of container types
|
||||
if v1 < v2 {
|
||||
randInts[i] = v1
|
||||
} else {
|
||||
randInts[i] = v2
|
||||
}
|
||||
}
|
||||
b.StartTimer()
|
||||
for i := 0; i < b.N; i++ {
|
||||
benchmarkOneCheckpoint(b, randInts)
|
||||
}
|
||||
b.StopTimer()
|
||||
done()
|
||||
}
|
||||
|
||||
// better diagnosis of deadlocks/hung situations versus just really slow "Quick" tests.
|
||||
func TestMain(m *testing.M) {
|
||||
l, err := net.Listen("tcp", ":0")
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ import (
|
|||
"sort"
|
||||
"testing"
|
||||
|
||||
"github.com/molecula/featurebase/v2/logger"
|
||||
"github.com/molecula/featurebase/v2/rbf"
|
||||
rbfcfg "github.com/molecula/featurebase/v2/rbf/cfg"
|
||||
"github.com/molecula/featurebase/v2/testhook"
|
||||
|
|
@ -65,6 +66,13 @@ func NewDB(tb testing.TB, cfg ...*rbfcfg.Config) *rbf.DB {
|
|||
// MustOpenDB returns a db opened on a temporary file. On error, fail test.
|
||||
func MustOpenDB(tb testing.TB, cfg ...*rbfcfg.Config) *rbf.DB {
|
||||
tb.Helper()
|
||||
if len(cfg) == 0 || cfg[0] == nil {
|
||||
newconf := rbfcfg.NewDefaultConfig()
|
||||
newconf.Logger = logger.NewLogfLogger(tb)
|
||||
cfg = []*rbfcfg.Config{newconf}
|
||||
} else if cfg[0].Logger == nil {
|
||||
cfg[0].Logger = logger.NewLogfLogger(tb)
|
||||
}
|
||||
db := NewDB(tb, cfg...)
|
||||
if err := db.Open(); err != nil {
|
||||
tb.Fatal(err)
|
||||
|
|
|
|||
15
rbf/tx.go
15
rbf/tx.go
|
|
@ -109,20 +109,25 @@ func (tx *Tx) Commit() error {
|
|||
// future plan: after checkpoint is moved to background
|
||||
// or not every removeTx, then we can move the
|
||||
// tx.db.rootRecords = tx.rootRecords into removeTx().
|
||||
|
||||
//
|
||||
// ... or maybe not: let's do that part here, and then removeTx
|
||||
// may or may not start a checkpoint, possibly asynchronously.
|
||||
//
|
||||
// avoid race detector firing on a write race here
|
||||
// vs the read of rootRecords at db.Begin()
|
||||
// vs the read of rootRecords at db.Begin(), then release
|
||||
// the lock, because we need removeTx to grab the lock to
|
||||
// work, but if it wants to checkpoint, it wants to be able to return
|
||||
// to us here and still be holding the lock.
|
||||
tx.db.mu.Lock()
|
||||
defer tx.db.mu.Unlock()
|
||||
tx.db.rootRecords = tx.rootRecords
|
||||
tx.db.pageMap = tx.pageMap
|
||||
tx.db.walPageN = tx.walPageN
|
||||
return tx.db.removeTx(tx)
|
||||
tx.db.mu.Unlock()
|
||||
}
|
||||
|
||||
// Disconnect transaction from DB.
|
||||
tx.db.mu.Lock()
|
||||
defer tx.db.mu.Unlock()
|
||||
// Disconnect transaction from DB.
|
||||
return tx.db.removeTx(tx)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -476,7 +476,6 @@ func NewServer(opts ...ServerOption) (*Server, error) {
|
|||
}
|
||||
s.holder = NewHolder(path, s.holderConfig)
|
||||
s.holder.Stats.SetLogger(s.logger)
|
||||
s.holder.Logger.Infof("RowCacheOn: %v", s.holderConfig.RowcacheOn)
|
||||
cwd, err := os.Getwd()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
|
|||
|
|
@ -200,9 +200,9 @@ type Config struct {
|
|||
// "rbf".
|
||||
Storage *storage.Config `toml:"storage"`
|
||||
|
||||
// RowcacheOn, if true, turns on the row cache for all storage backends.
|
||||
// The default is now off because it makes rbf queries faster and uses
|
||||
// much less memory.
|
||||
// RowcacheOn permanently disabled. No longer useful w/ RBF. Left
|
||||
// for backward compatibility but will be removed in a future
|
||||
// version.
|
||||
RowcacheOn bool `toml:"rowcache-on"`
|
||||
|
||||
// RBFConfig defines all externally configurable RBF flags.
|
||||
|
|
|
|||
|
|
@ -482,7 +482,7 @@ func (m *Command) SetupServer() error {
|
|||
pilosa.OptServerClusterName(m.Config.Cluster.Name),
|
||||
pilosa.OptServerSerializer(proto.Serializer{}),
|
||||
pilosa.OptServerStorageConfig(m.Config.Storage),
|
||||
pilosa.OptServerRowcacheOn(m.Config.RowcacheOn),
|
||||
pilosa.OptServerRowcacheOn(false),
|
||||
pilosa.OptServerRBFConfig(m.Config.RBFConfig),
|
||||
pilosa.OptServerMaxQueryMemory(m.Config.MaxQueryMemory),
|
||||
pilosa.OptServerQueryHistoryLength(m.Config.QueryHistoryLength),
|
||||
|
|
|
|||
15
txfactory.go
15
txfactory.go
|
|
@ -242,10 +242,15 @@ func (qcx *Qcx) GetTx(o Txo) (tx Tx, finisher func(perr *error), err error) {
|
|||
}
|
||||
|
||||
// qcx.write reflects the top executor determination
|
||||
// if a write will be done at the end, so we upgrade
|
||||
// the "local" read Tx to be writes, so that they
|
||||
// don't deadlock against themselves.
|
||||
o.Write = o.Write || qcx.write
|
||||
// if a write will be happen at some point, in which case, to avoid
|
||||
// locking problems with multi-shard things, we (probably incorrectly)
|
||||
// treat every Tx as its own individual separate Tx.
|
||||
//
|
||||
// But we still want to open non-write transactions individually, we
|
||||
// just can't recycle them (because write operations will come in and
|
||||
// we want them to work and commit right away so we're not holding a write
|
||||
// lock for long).
|
||||
writeLogic := o.Write || qcx.write
|
||||
|
||||
// In general, we make ALL write transactions local, and never reuse them
|
||||
// below. Previously this was to help lmdb.
|
||||
|
|
@ -273,7 +278,7 @@ func (qcx *Qcx) GetTx(o Txo) (tx Tx, finisher func(perr *error), err error) {
|
|||
return *qcx.RequiredForAtomicWriteTx, NoopFinisher, nil
|
||||
}
|
||||
|
||||
if !o.Write && qcx.Grp != nil {
|
||||
if !writeLogic && qcx.Grp != nil {
|
||||
// read, with a group in place.
|
||||
finisher = func(perr *error) {} // finisher is a returned value
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue