From 37f43d7d6f658890c5358434c44f59c489db9fc3 Mon Sep 17 00:00:00 2001 From: Souhaila Noor Date: Wed, 9 Mar 2022 09:40:08 -0600 Subject: [PATCH 1/3] added concurrency for delete during start-up --- holder.go | 85 +++++++++++++++++++---------- holder_internal_test.go | 118 ++++++++++++++++++++++++---------------- 2 files changed, 126 insertions(+), 77 deletions(-) diff --git a/holder.go b/holder.go index 4e9030f40..ef313136d 100644 --- a/holder.go +++ b/holder.go @@ -287,6 +287,44 @@ func (h *Holder) IndexesPath() string { return filepath.Join(h.path, IndexesDir) } +func (h *Holder) deletePerShard(index *Index, shard uint64) error { + inprocessRecords := NewRow() + + frag := h.fragment(index.name, existenceFieldName, viewStandard, shard) + if frag == nil { + return nil + } + + tx := index.Txf().NewTx(Txo{Write: !writable, Index: index, Shard: shard}) + defer tx.Rollback() + + // filter rows based on having _exists>=1, which is used to flag delete in-flight + rows, err := frag.rows(context.Background(), tx, 1) + if err != nil { + return err + } + + // check if any rows are found + if len(rows) == 0 { + return nil + } + + for _, record := range rows { + row, err2 := frag.row(tx, record) + if err2 != nil { + return fmt.Errorf("getting row IDs: %v", err2) + } + inprocessRecords = inprocessRecords.Union(row) + } + h.Logger.Printf("retrying delete: index=%v shard=%v record count=%v", index.name, shard, inprocessRecords.Count()) + + _, err = DeleteRows(context.Background(), inprocessRecords, index, shard) + if err != nil { + return fmt.Errorf("deleting rows: %v", err) + } + return nil +} + // processDeleteInflight checks if deletion was in progress when server shutdown // the _exists field is set to row+1 when delete is started. Upon completion, the row is deleted. // if _exists>=1, we finish deleting the rows @@ -294,38 +332,25 @@ func (h *Holder) processDeleteInflight() error { for _, index := range h.Indexes() { if index.trackExistence { shards := index.AvailableShards(includeRemote).Slice() - + index := index + ch := make(chan uint64, len(shards)) for _, shard := range shards { - inprocessRowIDs := NewRow() - - frag := h.fragment(index.name, existenceFieldName, viewStandard, shard) - if frag == nil { - continue - } - - tx := index.Txf().NewTx(Txo{Write: !writable, Index: index, Shard: shard}) - defer tx.Rollback() - - // filter rows based on having _exists>=1, which is used to flag delete in-flight - rows, err := frag.rows(context.Background(), tx, 1) - if err != nil { - return err - } - - // check if any rows are found - if len(rows) == 0 { - return nil - } - - for _, rowID := range rows { - row, err2 := frag.row(tx, rowID) - if err2 != nil { - return err2 - } - inprocessRowIDs = inprocessRowIDs.Union(row) - } - DeleteRows(context.Background(), inprocessRowIDs, index, shard) + ch <- shard } + close(ch) + + g := new(errgroup.Group) + for i := 0; i < runtime.NumCPU(); i++ { + g.Go(func() error { + for shard := range ch { + if err := h.deletePerShard(index, shard); err != nil { + return err + } + } + return nil + }) + } + g.Wait() } } return nil diff --git a/holder_internal_test.go b/holder_internal_test.go index 2fec338af..20535c9f9 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -18,6 +18,48 @@ func mustHolderConfig() *HolderConfig { return cfg } +func setupTest(t *testing.T, h *Holder, rowCol []rowCols, indexName string) (*Index, *Field) { + idx, err := h.CreateIndexIfNotExists(indexName, IndexOptions{TrackExistence: true}) + if err != nil { + t.Fatalf("failed to create index %v: %v", indexName, err) + } + f, err := idx.CreateFieldIfNotExists("f", OptFieldTypeDefault()) + if err != nil { + t.Fatalf("failed to create field in index %v: %v", indexName, err) + } + existencefield := idx.existenceFld + + shard := uint64(0) + tx := idx.Txf().NewTx(Txo{Write: true, Index: idx, Shard: shard}) + defer tx.Rollback() + for _, r := range rowCol { + _, err = f.SetBit(tx, r.row, r.col, nil) + if err != nil { + t.Fatalf("failed to set bit in index %v: %v", indexName, err) + } + + _, err = existencefield.SetBit(tx, r.row, r.col, nil) + if err != nil { + t.Fatalf("failed to set bit in index %v: %v", indexName, err) + } + } + + if err = tx.Commit(); err != nil { + t.Fatalf("failed to commit tx for index %v: %v", indexName, err) + } + + shardsFound := idx.AvailableShards(includeRemote).Slice() + if len(shardsFound) != 3 { + t.Fatalf("expected 3 shards for index %v, got %v", indexName, len(shardsFound)) + } + return idx, f +} + +type rowCols struct { + row uint64 + col uint64 +} + func TestHolder_ProcessDeleteInflight(t *testing.T) { path, _ := testhook.TempDir(t, "delete-inflight") h := NewHolder(path, mustHolderConfig()) @@ -28,63 +70,45 @@ func TestHolder_ProcessDeleteInflight(t *testing.T) { t.Fatalf("failed to open holder: %v", err) } - idx, err := h.CreateIndexIfNotExists("i", IndexOptions{TrackExistence: true}) - if err != nil { - t.Fatalf("failed to create index: %v", err) - } - f, err := idx.CreateFieldIfNotExists("f", OptFieldTypeDefault()) - if err != nil { - t.Fatalf("failed to create field: %v", err) - } - - existencefield := idx.existenceFld - shard := uint64(0) - tx := idx.Txf().NewTx(Txo{Write: true, Index: idx, Shard: shard}) - defer tx.Rollback() - - rowCol := []struct { - row uint64 - col uint64 - }{ + rowCol := []rowCols{ {1, 1}, {1, 2}, - {30, 33}, - {22, 2}, - } - for _, r := range rowCol { - _, err = f.SetBit(tx, r.row, r.col, nil) - if err != nil { - t.Fatalf("failed to set bit: %v", err) - } - - _, err = existencefield.SetBit(tx, r.row, r.col, nil) - if err != nil { - t.Fatalf("failed to set bit: %v", err) - } + {10, ShardWidth + 1}, + {1, ShardWidth * 2}, } - if err = tx.Commit(); err != nil { - t.Fatalf("failed to commit tx: %v", err) - } + idx1, f1 := setupTest(t, h, rowCol, "idxdelete1") + idx2, f2 := setupTest(t, h, rowCol, "idxdelete2") err = h.processDeleteInflight() if err != nil { t.Fatalf("failed to delete: %v", err) } - tx = idx.Txf().NewTx(Txo{Write: false, Index: idx, Shard: shard}) - defer tx.Rollback() - for _, r := range rowCol { - row, err := f.Row(tx, r.row) - if err != nil { - t.Fatalf("failed to get row: %v", err) - } - existenceRow, err := existencefield.Row(tx, r.row) - if err != nil { - t.Fatalf("failed to get row: %v", err) - } - if len(row.Columns()) != 0 || len(existenceRow.Columns()) != 0 { - t.Fatalf("expected columns for fields to be empty after delete") + tests := []struct { + idx *Index + f *Field + }{ + {idx1, f1}, + {idx2, f2}, + } + + for _, test := range tests { + idx, f := test.idx, test.f + tx := idx.Txf().NewTx(Txo{Write: false, Index: idx1, Shard: uint64(0)}) + defer tx.Rollback() + for _, r := range rowCol { + row, err := f.Row(tx, r.row) + if err != nil { + t.Fatalf("failed to get row: %v", err) + } + existenceRow, err := idx.existenceFld.Row(tx, r.row) + if err != nil { + t.Fatalf("failed to get row: %v", err) + } + if len(row.Columns()) != 0 || len(existenceRow.Columns()) != 0 { + t.Fatalf("expected columns for fields to be empty after delete") + } } } } From f653457b5a5e8930e1a4063fd7d1b4e647eb93fe Mon Sep 17 00:00:00 2001 From: Souhaila Noor Date: Thu, 10 Mar 2022 09:36:18 -0600 Subject: [PATCH 2/3] address reviewer's comments --- holder.go | 6 ++++-- holder_internal_test.go | 33 ++++++++++++++++++--------------- 2 files changed, 22 insertions(+), 17 deletions(-) diff --git a/holder.go b/holder.go index ef313136d..4d451c7ee 100644 --- a/holder.go +++ b/holder.go @@ -344,13 +344,15 @@ func (h *Holder) processDeleteInflight() error { g.Go(func() error { for shard := range ch { if err := h.deletePerShard(index, shard); err != nil { - return err + return fmt.Errorf("delete shard %d: %w", shard, err) } } return nil }) } - g.Wait() + if err := g.Wait(); err != nil { + return err + } } } return nil diff --git a/holder_internal_test.go b/holder_internal_test.go index 20535c9f9..fc8ca98ea 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -94,21 +94,24 @@ func TestHolder_ProcessDeleteInflight(t *testing.T) { } for _, test := range tests { - idx, f := test.idx, test.f - tx := idx.Txf().NewTx(Txo{Write: false, Index: idx1, Shard: uint64(0)}) - defer tx.Rollback() - for _, r := range rowCol { - row, err := f.Row(tx, r.row) - if err != nil { - t.Fatalf("failed to get row: %v", err) + func() { + idx, f := test.idx, test.f + tx := idx.Txf().NewTx(Txo{Write: false, Index: idx1, Shard: uint64(0)}) + defer tx.Rollback() + for _, r := range rowCol { + row, err := f.Row(tx, r.row) + if err != nil { + t.Fatalf("failed to get row: %v", err) + } + existenceRow, err := idx.existenceFld.Row(tx, r.row) + if err != nil { + t.Fatalf("failed to get row: %v", err) + } + if len(row.Columns()) != 0 || len(existenceRow.Columns()) != 0 { + t.Fatalf("expected columns for fields to be empty after delete") + } } - existenceRow, err := idx.existenceFld.Row(tx, r.row) - if err != nil { - t.Fatalf("failed to get row: %v", err) - } - if len(row.Columns()) != 0 || len(existenceRow.Columns()) != 0 { - t.Fatalf("expected columns for fields to be empty after delete") - } - } + }() + } } From a7b849c5fa7190d09d0b6baf83269fdd7649a0df Mon Sep 17 00:00:00 2001 From: Souhaila Noor Date: Thu, 10 Mar 2022 16:18:39 -0600 Subject: [PATCH 3/3] release read tx --- holder.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/holder.go b/holder.go index 4d451c7ee..e40edc3dc 100644 --- a/holder.go +++ b/holder.go @@ -318,6 +318,8 @@ func (h *Holder) deletePerShard(index *Index, shard uint64) error { } h.Logger.Printf("retrying delete: index=%v shard=%v record count=%v", index.name, shard, inprocessRecords.Count()) + tx.Rollback() // release the read tx in case a checksum is needed in DeleteRows + _, err = DeleteRows(context.Background(), inprocessRecords, index, shard) if err != nil { return fmt.Errorf("deleting rows: %v", err)