mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 04:17:51 +00:00
Merge pull request #1970 from molecula/FB-1185
[FB-1185] Batch delete transactions
This commit is contained in:
commit
3135c40f17
5 changed files with 145 additions and 58 deletions
191
executor.go
191
executor.go
|
|
@ -8239,10 +8239,12 @@ func (e *executor) executeDeleteRecords(ctx context.Context, qcx *Qcx, index str
|
|||
} else if len(c.Children) > 1 {
|
||||
return false, errors.New("Delete() only accepts a single bitmap input")
|
||||
}
|
||||
qcx.Abort()
|
||||
qcx.Reset() //release the qcx to allow for rbf checkpoint
|
||||
|
||||
// Execute calls in bulk on each remote node and merge.
|
||||
mapFn := func(ctx context.Context, shard uint64, mopt *mapOptions) (_ interface{}, err error) {
|
||||
return e.executeDeleteRecordFromShard(ctx, qcx, index, c, shard)
|
||||
return e.executeDeleteRecordFromShard(ctx, index, c.Children[0], shard)
|
||||
}
|
||||
|
||||
// Merge returned results at coordinating node.
|
||||
|
|
@ -8276,64 +8278,74 @@ func transactExistRow(ctx context.Context, idx *Index, shard uint64, frag *fragm
|
|||
}
|
||||
return rowID, tx.Commit()
|
||||
}
|
||||
func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (changed bool, err error) {
|
||||
func (e *executor) executeDeleteRecordFromShard(ctx context.Context, index string, bmCall *pql.Call, shard uint64) (changed bool, err error) {
|
||||
span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeDeleteRecordFromShard")
|
||||
defer span.Finish()
|
||||
//need to build the bitmap in the call
|
||||
child := c.Children[0]
|
||||
src, er := e.executeBitmapCallShard(ctx, qcx, index, child, shard)
|
||||
if er != nil {
|
||||
err = er
|
||||
return
|
||||
}
|
||||
if len(src.segments) == 0 { //nothing to remove
|
||||
return
|
||||
}
|
||||
columns := src.segments[0].data //should only be one segment
|
||||
if columns.Count() == 0 {
|
||||
return
|
||||
}
|
||||
// Fetch index.
|
||||
idx := e.Holder.Index(index)
|
||||
if idx == nil {
|
||||
err = newNotFoundError(ErrIndexNotFound, index)
|
||||
return
|
||||
}
|
||||
|
||||
qcx := idx.Txf().NewQcx()
|
||||
//bmCall is a bitmap
|
||||
row, err := e.executeBitmapCallShard(ctx, qcx, index, bmCall, shard)
|
||||
qcx.Abort()
|
||||
if len(row.segments) == 0 {
|
||||
return
|
||||
}
|
||||
columns := row.segments[0].data
|
||||
if columns.Count() == 0 {
|
||||
return
|
||||
}
|
||||
src := NewRowFromBitmap(columns)
|
||||
return DeleteRowsWithFlow(ctx, src, idx, shard, true)
|
||||
}
|
||||
|
||||
func DeleteRows(ctx context.Context, src *Row, idx *Index, shard uint64) (bool, error) {
|
||||
return DeleteRowsWithFlow(ctx, src, idx, shard, false)
|
||||
}
|
||||
func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, normalFlow bool) (bool, error) {
|
||||
func clearFragment(writeTx Tx, columns *roaring.Bitmap, frag *fragment, resChan chan countResults) (changed bool, err error) {
|
||||
posChan := make(chan uint64, 8192)
|
||||
findExisting := roaring.NewBitmapBitmapFilter(columns, func(pos uint64) error {
|
||||
posChan <- pos
|
||||
return nil
|
||||
})
|
||||
go writeTx.RemoveChannel(frag.index(), frag.field(), frag.view(), frag.shard, posChan, resChan)
|
||||
|
||||
err = writeTx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting)
|
||||
close(posChan)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
r := <-resChan
|
||||
|
||||
changed = r.changeCount > 0
|
||||
err = r.err
|
||||
return
|
||||
}
|
||||
|
||||
func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (bool, error) {
|
||||
var existenceFragment *fragment
|
||||
var deletedRowID uint64
|
||||
var commitor Commitor = &NopCommitor{}
|
||||
var err error
|
||||
if len(src.segments) == 0 { //nothing to remove
|
||||
return false, nil
|
||||
}
|
||||
columns := src.segments[0].data //should only be one segment
|
||||
if columns.Count() == 0 {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
if idx.Keys() {
|
||||
//store columns in exits field ToBeDelete row commited
|
||||
if normalFlow { // normalFlow is the standard path, "not normal" is recoverory
|
||||
existenceFragment = idx.Holder().fragment(idx.Name(), existenceFieldName, viewStandard, shard)
|
||||
if existenceFragment == nil {
|
||||
//no exists field
|
||||
return false, errors.New("can't bulk delete without existence field")
|
||||
}
|
||||
deletedRowID, err = transactExistRow(ctx, idx, shard, existenceFragment, src)
|
||||
var err error //store columns in exits field ToBeDelete row commited
|
||||
if normalFlow { // normalFlow is the standard path, "not normal" is recoverory
|
||||
existenceFragment = idx.Holder().fragment(idx.Name(), existenceFieldName, viewStandard, shard)
|
||||
if existenceFragment == nil {
|
||||
//no exists field
|
||||
return false, errors.New("can't bulk delete without existence field")
|
||||
}
|
||||
commitor, err = deleteKeyTranslation(ctx, idx, shard, columns)
|
||||
src := NewRowFromBitmap(columns)
|
||||
deletedRowID, err = transactExistRow(ctx, idx, shard, existenceFragment, src)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
commitor, err = deleteKeyTranslation(ctx, idx, shard, columns)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
writeTx := idx.Txf().NewTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
return false, err
|
||||
|
|
@ -8363,24 +8375,7 @@ func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64,
|
|||
}
|
||||
|
||||
}()
|
||||
findExisting := roaring.NewBitmapBitmapFilter(columns, func(p uint64) error { return nil })
|
||||
resChan := make(chan countResults)
|
||||
clearFragment := func(frag *fragment) (bool, error) {
|
||||
posChan := make(chan uint64, 8192)
|
||||
findExisting.SetCallback(func(pos uint64) error {
|
||||
posChan <- pos
|
||||
return nil
|
||||
})
|
||||
go writeTx.RemoveChannel(frag.index(), frag.field(), frag.view(), frag.shard, posChan, resChan)
|
||||
|
||||
err = writeTx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting)
|
||||
close(posChan)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
r := <-resChan
|
||||
return r.changeCount > 0, r.err
|
||||
}
|
||||
|
||||
for _, field := range idx.Fields() {
|
||||
for _, view := range field.views() {
|
||||
|
|
@ -8389,7 +8384,7 @@ func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64,
|
|||
if !ok {
|
||||
continue
|
||||
}
|
||||
c, err := clearFragment(frag)
|
||||
c, err := clearFragment(writeTx, columns, frag, resChan)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
|
@ -8417,6 +8412,92 @@ func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64,
|
|||
return changed, nil
|
||||
}
|
||||
|
||||
func DeleteRowsWithOutKeysFlow(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (changed bool, err error) {
|
||||
var existenceFragment *fragment
|
||||
var deletedRowID uint64
|
||||
var commitor Commitor = &NopCommitor{}
|
||||
writeTx := idx.Txf().NewTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
defer writeTx.Rollback()
|
||||
defer func() {
|
||||
// if there is an error in the key commit, then rollback the delete
|
||||
// write records before keys to remove possiblity of unmatch keys=records
|
||||
err := writeTx.Commit()
|
||||
if err != nil {
|
||||
changed = false
|
||||
commitor.Rollback()
|
||||
return
|
||||
}
|
||||
}()
|
||||
resChan := make(chan countResults)
|
||||
for _, field := range idx.Fields() {
|
||||
for _, view := range field.views() {
|
||||
|
||||
frag, ok := view.fragments[shard]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
c, err := clearFragment(writeTx, columns, frag, resChan)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if c {
|
||||
changed = true
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
close(resChan)
|
||||
if existenceFragment == nil { //a string keys have been deleted and the deleteRow was created
|
||||
return changed, nil
|
||||
}
|
||||
|
||||
if normalFlow {
|
||||
existenceFragment.clearRow(writeTx, deletedRowID)
|
||||
return changed, nil
|
||||
}
|
||||
|
||||
// this is if we are recovering from failure and cleaning up
|
||||
rows, err := existenceFragment.rows(ctx, writeTx, 1)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
for _, rowId := range rows {
|
||||
existenceFragment.clearRow(writeTx, rowId)
|
||||
}
|
||||
return changed, nil
|
||||
}
|
||||
|
||||
func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, normalFlow bool) (change bool, err error) {
|
||||
if len(src.segments) == 0 { //nothing to remove
|
||||
return false, nil
|
||||
}
|
||||
columns := src.segments[0].data //should only be one segment
|
||||
if columns.Count() == 0 {
|
||||
return false, nil
|
||||
}
|
||||
bits := src.segments[0].data.Slice()
|
||||
min := func(a, b int) int {
|
||||
if a <= b {
|
||||
return a
|
||||
}
|
||||
return b
|
||||
}
|
||||
limit := idx.holder.cfg.RBFConfig.MaxDelete
|
||||
for i := 0; i < len(bits); i += limit {
|
||||
batch := roaring.NewBitmap(bits[i:min(i+limit, len(bits))]...)
|
||||
if idx.Keys() {
|
||||
change, err = DeleteRowsWithFlowWithKeys(ctx, batch, idx, shard, normalFlow)
|
||||
} else {
|
||||
change, err = DeleteRowsWithOutKeysFlow(ctx, batch, idx, shard, normalFlow)
|
||||
}
|
||||
if err != nil {
|
||||
return change, err
|
||||
}
|
||||
}
|
||||
return change, err
|
||||
|
||||
}
|
||||
|
||||
type Commitor interface {
|
||||
Rollback()
|
||||
Commit() error
|
||||
|
|
|
|||
2
rbf.go
2
rbf.go
|
|
@ -253,6 +253,8 @@ func (tx *RBFTx) RemoveChannel(index, field, view string, shard uint64, a chan u
|
|||
return
|
||||
}
|
||||
} else {
|
||||
rc = roaring.Optimize(rc)
|
||||
|
||||
err = tx.tx.PutContainer(name, lastHi, rc)
|
||||
if err != nil {
|
||||
resChan <- countResults{0, errors.Wrap(err, "failed to put container")}
|
||||
|
|
|
|||
|
|
@ -41,6 +41,9 @@ type Config struct {
|
|||
// background checkpoints. It cannot be set from toml. The default is
|
||||
// to use stderr.
|
||||
Logger logger.Logger `toml:"-"`
|
||||
|
||||
// The maximum number of bits to be deleted in a single transaction default(65536)
|
||||
MaxDelete int `toml:"max-delete"`
|
||||
}
|
||||
|
||||
func NewDefaultConfig() *Config {
|
||||
|
|
@ -50,6 +53,7 @@ func NewDefaultConfig() *Config {
|
|||
MinWALCheckpointSize: DefaultMinWALCheckpointSize,
|
||||
MaxWALCheckpointSize: DefaultMaxWALCheckpointSize,
|
||||
FsyncEnabled: true,
|
||||
MaxDelete: DefaultMaxDelete,
|
||||
|
||||
// CI passed with 20. 50 was too big for CI, even on X-large instances.
|
||||
// For now we default to 0, which means use sync.Pool.
|
||||
|
|
|
|||
|
|
@ -13,3 +13,6 @@ const DefaultMaxSize = 4 * (1 << 30)
|
|||
// size of the WAL. The size can be increased by updating the DB.MaxWALSize
|
||||
// and reopening the database. This setting mainly affects virtual space usage.
|
||||
const DefaultMaxWALSize = 4 * (1 << 30)
|
||||
|
||||
// DefaultMaxDelete is the maximum number of bits that will be deleted in a single batch
|
||||
const DefaultMaxDelete = 65536
|
||||
|
|
|
|||
|
|
@ -158,9 +158,6 @@ func (q *Qcx) Abort() {
|
|||
func (q *Qcx) Reset() {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if !q.done {
|
||||
vprint.PanicOn("must call Qcx.Abort() or Qcx.Finish() before calling Reset().")
|
||||
}
|
||||
q.unprotected_reset()
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue