diff --git a/cluster.go b/cluster.go index 820349ded..cdcf50b26 100644 --- a/cluster.go +++ b/cluster.go @@ -144,14 +144,18 @@ type Cluster struct { // Threshold for logging long-running queries LongQueryTime time.Duration + + // Maximum number of SetBit() or ClearBit() commands per request. + MaxWritesPerRequest int } // NewCluster returns a new instance of Cluster with defaults. func NewCluster() *Cluster { return &Cluster{ - Hasher: &jmphasher{}, - PartitionN: DefaultPartitionN, - ReplicaN: DefaultReplicaN, + Hasher: &jmphasher{}, + PartitionN: DefaultPartitionN, + ReplicaN: DefaultReplicaN, + MaxWritesPerRequest: DefaultMaxWritesPerRequest, } } diff --git a/fragment.go b/fragment.go index 073bdfa84..930c8e02b 100644 --- a/fragment.go +++ b/fragment.go @@ -493,7 +493,7 @@ func (f *Fragment) FieldValue(columnID uint64, bitDepth uint) (value uint64, exi f.mu.Lock() defer f.mu.Unlock() - // If existance bit is unset then ignore remaining bits. + // If existence bit is unset then ignore remaining bits. if v, err := f.bit(uint64(bitDepth), columnID); err != nil { return 0, false, err } else if !v { @@ -587,7 +587,7 @@ func (f *Fragment) importSetFieldValue(columnID uint64, bitDepth uint, value uin // FieldSum returns the sum of a given field as well as the number of columns involved. // A bitmap can be passed in to optionally filter the computed columns. func (f *Fragment) FieldSum(filter *Bitmap, bitDepth uint) (sum, count uint64, err error) { - // Compute count based on the existance bit. + // Compute count based on the existence bit. row := f.Row(uint64(bitDepth)) if filter != nil { count = row.IntersectionCount(filter) @@ -616,6 +616,7 @@ func (f *Fragment) FieldSum(filter *Bitmap, bitDepth uint) (sum, count uint64, e return sum, count, nil } +// FieldRange returns bitmaps with a field value encoding matching the predicate. func (f *Fragment) FieldRange(op pql.Token, bitDepth uint, predicate uint64) (*Bitmap, error) { switch op { case pql.EQ: @@ -754,6 +755,7 @@ func (f *Fragment) FieldNotNull(bitDepth uint) (*Bitmap, error) { return f.Row(uint64(bitDepth)), nil } +// FieldRangeBetween returns bitmaps with a field value encoding matching any value between predicateMin and predicateMax. func (f *Fragment) FieldRangeBetween(bitDepth uint, predicateMin, predicateMax uint64) (*Bitmap, error) { b := f.Row(uint64(bitDepth)) keep1 := NewBitmap() // GTE @@ -1826,36 +1828,43 @@ func (s *FragmentSyncer) syncBlock(id int) error { // Write updates to remote blocks. for i := 0; i < len(clients); 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. - var buf bytes.Buffer + // Generate query with sets & clears, and group the requests to not exceed MaxWritesPerRequest. + total := len(set.ColumnIDs) + len(clear.ColumnIDs) + buffers := make([]bytes.Buffer, int(math.Ceil(float64(total)/float64(s.Cluster.MaxWritesPerRequest)))) // Only sync the standard block. for j := 0; j < len(set.ColumnIDs); j++ { - fmt.Fprintf(&buf, "SetBit(frame=%q, rowID=%d, columnID=%d)\n", f.Frame(), set.RowIDs[j], (f.Slice()*SliceWidth)+set.ColumnIDs[j]) + fmt.Fprintf(&(buffers[count/s.Cluster.MaxWritesPerRequest]), "SetBit(frame=%q, rowID=%d, columnID=%d)\n", f.Frame(), set.RowIDs[j], (f.Slice()*SliceWidth)+set.ColumnIDs[j]) + count++ } for j := 0; j < len(clear.ColumnIDs); j++ { - fmt.Fprintf(&buf, "ClearBit(frame=%q, rowID=%d, columnID=%d)\n", f.Frame(), clear.RowIDs[j], (f.Slice()*SliceWidth)+clear.ColumnIDs[j]) + fmt.Fprintf(&(buffers[count/s.Cluster.MaxWritesPerRequest]), "ClearBit(frame=%q, rowID=%d, columnID=%d)\n", f.Frame(), clear.RowIDs[j], (f.Slice()*SliceWidth)+clear.ColumnIDs[j]) + count++ } - // Verify sync is not prematurely closing. - if s.isClosing() { - return nil - } + // 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 := &internal.QueryRequest{ - Query: buf.String(), - Remote: true, - } - _, err := clients[i].ExecuteQuery(context.Background(), f.Index(), queryRequest) - if err != nil { - return err + // Execute query. + queryRequest := &internal.QueryRequest{ + Query: buffers[k].String(), + Remote: true, + } + _, err := clients[i].ExecuteQuery(context.Background(), f.Index(), queryRequest) + if err != nil { + return err + } } } diff --git a/server.go b/server.go index ffaad0b71..208d0cd6f 100644 --- a/server.go +++ b/server.go @@ -177,6 +177,7 @@ func (s *Server) Open() error { e.Host = s.URI.HostPort() e.Cluster = s.Cluster e.MaxWritesPerRequest = s.MaxWritesPerRequest + s.Cluster.MaxWritesPerRequest = s.MaxWritesPerRequest // Initialize HTTP handler. s.Handler.Broadcaster = s.Broadcaster