batch size 65536

This commit is contained in:
Todd Gruben 2022-03-08 13:00:08 -06:00
parent ac450a0c95
commit d322cc7aef

View file

@ -8309,33 +8309,23 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i
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 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)
}
commitor, err = deleteKeyTranslation(ctx, idx, shard, columns)
if err != nil {
return false, err
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")
}
src := NewRowFromBitmap(columns)
deletedRowID, err = transactExistRow(ctx, idx, shard, existenceFragment, src)
}
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 {
@ -8420,6 +8410,114 @@ 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) (bool, error) {
var existenceFragment *fragment
var deletedRowID uint64
var commitor Commitor = &NopCommitor{}
var err error
writeTx := idx.Txf().NewTx(Txo{Write: writable, Index: idx, Shard: shard})
if err != nil {
return false, err
}
defer writeTx.Rollback()
changed := false
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
}
}()
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() {
frag, ok := view.fragments[shard]
if !ok {
continue
}
c, err := clearFragment(frag)
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
if normalFlow {
existenceFragment.clearRow(writeTx, deletedRowID)
} else {
// 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) (bool, 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
}
var change bool
var err error
limit := 65536
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