From c66daabc81b17ce8193d5f828590d9928ffad9d8 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 10 Dec 2018 20:37:26 -0600 Subject: [PATCH] convert the anti-entropy logic to use `ImportRoaring` instead of `QueryNode` --- api.go | 2 - fragment.go | 102 +++++++++++++++++++++++++++++++------------------ holder_test.go | 11 +++--- 3 files changed, 69 insertions(+), 46 deletions(-) diff --git a/api.go b/api.go index cf27522e6..ce4fbd0ca 100644 --- a/api.go +++ b/api.go @@ -333,8 +333,6 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, } return err }) - go func(node *Node) { - }(node) } else if !remote { // if remote == true we don't forward to other nodes // forward it on eg.Go(func() error { diff --git a/fragment.go b/fragment.go index 028d105b4..4f71809c9 100644 --- a/fragment.go +++ b/fragment.go @@ -28,6 +28,7 @@ import ( "math" "os" "sort" + "strings" "sync" "syscall" "time" @@ -2315,46 +2316,38 @@ func (s *fragmentSyncer) syncBlock(id int) error { // Write updates to remote blocks. for i := 0; i < len(uris); i++ { set, clear := sets[i], clears[i] - count := 0 - // Ignore if there are no differences. - if len(set.columnIDs) == 0 && len(clear.columnIDs) == 0 { - continue - } - - // Generate query with sets & clears, and group the requests to not exceed MaxWritesPerRequest. - total := len(set.columnIDs) + len(clear.columnIDs) - maxWrites := s.Cluster.maxWritesPerRequest - if maxWrites <= 0 { - maxWrites = 5000 - } - buffers := make([]bytes.Buffer, int(math.Ceil(float64(total)/float64(maxWrites)))) - - // Only sync the standard block. - for j := 0; j < len(set.columnIDs); j++ { - fmt.Fprintf(&(buffers[count/maxWrites]), "Set(%d, %s=%d)\n", (f.shard*ShardWidth)+set.columnIDs[j], f.field, set.rowIDs[j]) - count++ - } - for j := 0; j < len(clear.columnIDs); j++ { - fmt.Fprintf(&(buffers[count/maxWrites]), "Clear(%d, %s=%d)\n", (f.shard*ShardWidth)+clear.columnIDs[j], f.field, clear.rowIDs[j]) - count++ - } - - // Iterate over the buffers. - for k := 0; k < len(buffers); k++ { - // Verify sync is not prematurely closing. - if s.isClosing() { - return nil - } - - // Execute query. - queryRequest := &QueryRequest{ - Query: buffers[k].String(), - Remote: true, - } - _, err := s.Cluster.InternalClient.QueryNode(ctx, uris[i], f.index, queryRequest) + // Handle Sets. + if len(set.columnIDs) > 0 { + setData, err := bitsToRoaringData(set) if err != nil { - return errors.Wrap(err, "executing") + return errors.Wrap(err, "converting bits to roaring data (set)") + } + + setReq := &ImportRoaringRequest{ + Clear: false, + Views: map[string][]byte{cleanViewName(f.view): setData}, + } + + if err := s.Cluster.InternalClient.ImportRoaring(ctx, uris[i], f.index, f.field, f.shard, true, setReq); err != nil { + return errors.Wrap(err, "sending roaring data (set)") + } + } + + // Handle Clears. + if len(clear.columnIDs) > 0 { + clearData, err := bitsToRoaringData(clear) + if err != nil { + return errors.Wrap(err, "converting bits to roaring data (clear)") + } + + clearReq := &ImportRoaringRequest{ + Clear: true, + Views: map[string][]byte{"": clearData}, + } + + if err := s.Cluster.InternalClient.ImportRoaring(ctx, uris[i], f.index, f.field, f.shard, true, clearReq); err != nil { + return errors.Wrap(err, "sending roaring data (clear)") } } } @@ -2362,6 +2355,39 @@ func (s *fragmentSyncer) syncBlock(id int) error { return nil } +// cleanViewName converts a viewname into the equivalent +// string required by the external api. Because views are +// not exposed externally, the conversion looks like this: +// "standard" -> "" +// "standard_YYYYMMDD" -> "YYYYMMDD" +// "other" -> "other" (there is currently not a use for this) +func cleanViewName(v string) string { + viewPrefix := viewStandard + "_" + if strings.HasPrefix(v, viewPrefix) { + return v[len(viewPrefix):] + } else if v == viewStandard { + return "" + } + return v +} + +// bitsToRoaringData converts a pairSet into a roaring.Bitmap +// which represents the data within a single shard. +func bitsToRoaringData(ps pairSet) ([]byte, error) { + bmp := roaring.NewBitmap() + for j := 0; j < len(ps.columnIDs); j++ { + bmp.DirectAdd(ps.rowIDs[j]*ShardWidth + (ps.columnIDs[j] % ShardWidth)) + } + + var buf bytes.Buffer + _, err := bmp.WriteTo(&buf) + if err != nil { + return nil, errors.Wrap(err, "writing to buffer") + } + + return buf.Bytes(), nil +} + func madvise(b []byte, advice int) error { // nolint: unparam _, _, err := syscall.Syscall(syscall.SYS_MADVISE, uintptr(unsafe.Pointer(&b[0])), uintptr(len(b)), uintptr(advice)) if err != 0 { diff --git a/holder_test.go b/holder_test.go index 3920e4cc7..2835e5c4f 100644 --- a/holder_test.go +++ b/holder_test.go @@ -391,27 +391,26 @@ func TestHolderSyncer_TimeQuantum(t *testing.T) { hldr0 := &test.Holder{Holder: c[0].Server.Holder()} hldr1 := &test.Holder{Holder: c[1].Server.Holder()} - // Set data on the local holder. + // Set data on the local holder for node0. t1 := time.Date(2018, 8, 1, 12, 30, 0, 0, time.UTC) t2 := time.Date(2018, 8, 2, 12, 30, 0, 0, time.UTC) hldr0.SetBitTime("i", "f", 0, 1, &t1) hldr0.SetBitTime("i", "f", 0, 2, &t2) + // Set data on node1 + hldr0.SetBitTime("i", "f", 0, 22, &t2) + err = c[0].Server.SyncData() if err != nil { t.Fatalf("syncing node 0: %v", err) } - err = c[1].Server.SyncData() - if err != nil { - t.Fatalf("syncing node 1: %v", err) - } // Verify data is the same on both nodes. for i, hldr := range []*test.Holder{hldr0, hldr1} { if a := hldr.RowTime("i", "f", 0, t1, quantum).Columns(); !reflect.DeepEqual(a, []uint64{1}) { t.Errorf("unexpected columns(%d/0): %+v", i, a) } - if a := hldr.RowTime("i", "f", 0, t2, quantum).Columns(); !reflect.DeepEqual(a, []uint64{2}) { + if a := hldr.RowTime("i", "f", 0, t2, quantum).Columns(); !reflect.DeepEqual(a, []uint64{2, 22}) { t.Errorf("unexpected columns(%d/0): %+v", i, a) } }