Merge pull request #1780 from travisturner/anti-entropy-roaring

convert the anti-entropy logic to use `ImportRoaring` instead of `QueryNode`
This commit is contained in:
Travis Turner 2018-12-10 22:01:26 -06:00 committed by GitHub
commit 01e520d9ea
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
3 changed files with 69 additions and 46 deletions

2
api.go
View file

@ -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 {

View file

@ -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 {

View file

@ -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)
}
}