diff --git a/cluster.go b/cluster.go index aae5a344a..5929f6308 100644 --- a/cluster.go +++ b/cluster.go @@ -21,6 +21,16 @@ type Node struct { // Nodes represents a list of nodes. type Nodes []*Node +// Contains returns true if a node exists in the list. +func (a Nodes) Contains(n *Node) bool { + for i := range a { + if a[i] == n { + return true + } + } + return false +} + // ContainsHost returns true if host matches on of the node's host. func (a Nodes) ContainsHost(host string) bool { for _, n := range a { @@ -31,6 +41,17 @@ func (a Nodes) ContainsHost(host string) bool { return false } +// Filter returns a new list of nodes with node removed. +func (a Nodes) Filter(n *Node) []*Node { + other := make([]*Node, 0, len(a)) + for i := range a { + if a[i] != n { + other = append(other, a[i]) + } + } + return other +} + // FilterHost returns a new list of nodes with host removed. func (a Nodes) FilterHost(host string) []*Node { other := make([]*Node, 0, len(a)) @@ -51,6 +72,13 @@ func (a Nodes) Hosts() []string { return hosts } +// Clone returns a shallow copy of nodes. +func (a Nodes) Clone() []*Node { + other := make([]*Node, len(a)) + copy(other, a) + return other +} + // Cluster represents a collection of nodes. type Cluster struct { Nodes []*Node @@ -74,6 +102,16 @@ func NewCluster() *Cluster { } } +// NodeByHost returns a node reference by host. +func (c *Cluster) NodeByHost(host string) *Node { + for _, n := range c.Nodes { + if n.Host == host { + return n + } + } + return nil +} + // Partition returns the partition that a slice belongs to. func (c *Cluster) Partition(db string, slice uint64) int { var buf [8]byte diff --git a/executor.go b/executor.go index ac4da5d90..a2fb40fea 100644 --- a/executor.go +++ b/executor.go @@ -106,29 +106,28 @@ func (e *Executor) executeCall(ctx context.Context, db string, c pql.Call, slice // executeBitmapCall executes a call that returns a bitmap. func (e *Executor) executeBitmapCall(ctx context.Context, db string, c pql.BitmapCall, slices []uint64, opt *ExecOptions) (*Bitmap, error) { - other := NewBitmap() - for node, nodeSlices := range e.slicesByNode(db, slices) { - // Execute locally if the hostname matches. - if node.Host == e.Host { - for _, slice := range nodeSlices { - bm, err := e.executeBitmapCallSlice(ctx, db, c, slice) - if err != nil { - return nil, err - } - other.Merge(bm) - } - continue - } + // Execute calls in bulk on each remote node and merge. + mapFn := func(slice uint64) (interface{}, error) { + return e.executeBitmapCallSlice(ctx, db, c, slice) + } - // Otherwise execute remotely. - res, err := e.exec(ctx, node, db, &pql.Query{Calls: pql.Calls{c}}, nodeSlices, opt) - if err != nil { - return nil, err + // Merge returned results at coordinating node. + reduceFn := func(prev, v interface{}) interface{} { + other, _ := prev.(*Bitmap) + if other == nil { + other = NewBitmap() } - other.Merge(res[0].(*Bitmap)) + other.Merge(v.(*Bitmap)) + return other + } + + other, err := e.mapReduce(ctx, db, slices, c, opt, mapFn, reduceFn) + if err != nil { + return nil, err } // Attach bitmap attributes for Bitmap() calls. + bm, _ := other.(*Bitmap) if c, ok := c.(*pql.Bitmap); ok { fr := e.Index.Frame(db, c.Frame) if fr != nil { @@ -136,11 +135,11 @@ func (e *Executor) executeBitmapCall(ctx context.Context, db string, c pql.Bitma if err != nil { return nil, err } - other.Attrs = attrs + bm.Attrs = attrs } } - return other, nil + return bm, nil } // executeBitmapCallSlice executes a bitmap call for a single slice. @@ -187,43 +186,23 @@ func (e *Executor) executeTopN(ctx context.Context, db string, c *pql.TopN, slic } func (e *Executor) executeTopNSlices(ctx context.Context, db string, c *pql.TopN, slices []uint64, opt *ExecOptions) ([]Pair, error) { - slicesByNode := e.slicesByNode(db, slices) - - type resp struct { - pairs []Pair - err error - } - ch := make(chan resp, len(slicesByNode)) - - for node, nodeSlices := range slicesByNode { - go func(node *Node, nodeSlices []uint64) { - // Execute locally if the hostname matches. - if node.Host == e.Host { - pairs, err := e.executeTopNSlicesLocal(ctx, db, c, nodeSlices) - ch <- resp{pairs: pairs, err: err} - return - } - - // Otherwise execute remotely. - res, err := e.exec(ctx, node, db, &pql.Query{Calls: pql.Calls{c}}, nodeSlices, opt) - if err != nil { - ch <- resp{err: err} - return - } - ch <- resp{pairs: res[0].([]Pair)} - }(node, nodeSlices) + // Execute calls in bulk on each remote node and merge. + mapFn := func(slice uint64) (interface{}, error) { + return e.executeTopNSlice(ctx, db, c, slice) } - // Collect results. - var results []Pair - for range slicesByNode { - r := <-ch - if r.err != nil { - return nil, r.err - } - results = Pairs(results).Add(r.pairs) + // Merge returned results at coordinating node. + reduceFn := func(prev, v interface{}) interface{} { + other, _ := prev.([]Pair) + return Pairs(other).Add(v.([]Pair)) } + other, err := e.mapReduce(ctx, db, slices, c, opt, mapFn, reduceFn) + if err != nil { + return nil, err + } + results, _ := other.([]Pair) + // Sort final merged results. sort.Sort(Pairs(results)) @@ -235,34 +214,6 @@ func (e *Executor) executeTopNSlices(ctx context.Context, db string, c *pql.TopN return results, nil } -func (e *Executor) executeTopNSlicesLocal(ctx context.Context, db string, c *pql.TopN, slices []uint64) ([]Pair, error) { - type resp struct { - pairs []Pair - err error - } - ch := make(chan resp, len(slices)) - - // Execute TopN() in parallel across slices. - for _, slice := range slices { - go func(slice uint64) { - pairs, err := e.executeTopNSlice(ctx, db, c, slice) - ch <- resp{pairs: pairs, err: err} - }(slice) - } - - // Collect results. - var results []Pair - for range slices { - r := <-ch - if r.err != nil { - return nil, r.err - } - results = Pairs(results).Add(r.pairs) - } - - return results, nil -} - // executeTopNSlice executes a TopN call for a single slice. func (e *Executor) executeTopNSlice(ctx context.Context, db string, c *pql.TopN, slice uint64) ([]Pair, error) { // Retrieve bitmap used to intersect. @@ -381,27 +332,27 @@ func (e *Executor) executeUnionSlice(ctx context.Context, db string, c *pql.Unio // executeCount executes a count() call. func (e *Executor) executeCount(ctx context.Context, db string, c *pql.Count, slices []uint64, opt *ExecOptions) (uint64, error) { - var n uint64 - for node, nodeSlices := range e.slicesByNode(db, slices) { - // Execute locally if the hostname matches. - if node.Host == e.Host { - for _, slice := range nodeSlices { - bm, err := e.executeBitmapCallSlice(ctx, db, c.Input, slice) - if err != nil { - return 0, err - } - n += bm.Count() - } - continue - } - - // Otherwise execute remotely. - res, err := e.exec(ctx, node, db, &pql.Query{Calls: pql.Calls{c}}, nodeSlices, opt) + // Execute calls in bulk on each remote node and merge. + mapFn := func(slice uint64) (interface{}, error) { + bm, err := e.executeBitmapCallSlice(ctx, db, c.Input, slice) if err != nil { return 0, err } - n += res[0].(uint64) + return bm.Count(), nil } + + // Merge returned results at coordinating node. + reduceFn := func(prev, v interface{}) interface{} { + other, _ := prev.(uint64) + return other + v.(uint64) + } + + result, err := e.mapReduce(ctx, db, slices, c, opt, mapFn, reduceFn) + if err != nil { + return 0, err + } + n, _ := result.(uint64) + return n, nil } @@ -717,17 +668,170 @@ func (e *Executor) exec(ctx context.Context, node *Node, db string, q *pql.Query } // slicesByNode returns a mapping of nodes to slices. -// -// NOTE: Currently the only primary node is used. -func (e *Executor) slicesByNode(db string, slices []uint64) map[*Node][]uint64 { +// Returns errSliceUnavailable if a slice cannot be allocated to a node. +func (e *Executor) slicesByNode(nodes []*Node, db string, slices []uint64) (map[*Node][]uint64, error) { m := make(map[*Node][]uint64) - for _, slice := range slices { - nodes := e.Cluster.FragmentNodes(db, slice) - node := nodes[0] - m[node] = append(m[node], slice) +loop: + for _, slice := range slices { + for _, node := range e.Cluster.FragmentNodes(db, slice) { + if Nodes(nodes).Contains(node) { + m[node] = append(m[node], slice) + continue loop + } + } + return nil, errSliceUnavailable } - return m + return m, nil +} + +// mapReduce maps and reduces data across the cluster. +// +// If a mapping of slices to a node fails then the slices are resplit across +// secondary nodes and retried. This continues to occur until all nodes are exhausted. +func (e *Executor) mapReduce(ctx context.Context, db string, slices []uint64, c pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc) (interface{}, error) { + ch := make(chan mapResponse, 0) + + // Wrap context with a cancel to kill goroutines on exit. + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + // If this is the coordinating node then start with all nodes in the cluster. + // + // However, if this request is being sent from the coordinator then all + // processing should be done locally so we start with just the local node. + var nodes []*Node + if !opt.Remote { + nodes = Nodes(e.Cluster.Nodes).Clone() + } else { + nodes = []*Node{e.Cluster.NodeByHost(e.Host)} + } + + // Start mapping across all primary owners. + if err := e.mapper(ctx, ch, nodes, db, slices, c, opt, mapFn, reduceFn); err != nil { + return nil, err + } + + // Iterate over all map responses and reduce. + var result interface{} + var sliceN int + for { + select { + case <-ctx.Done(): + return nil, ctx.Err() + case resp := <-ch: + // On error retry against remaining nodes. If an error returns then + // the context will cancel and cause all open goroutines to return. + if resp.err != nil { + // Filter out unavailable nodes. + nodes = Nodes(nodes).Filter(resp.node) + + // Begin mapper against secondary nodes. + if err := e.mapper(ctx, ch, nodes, db, resp.slices, c, opt, mapFn, reduceFn); err == errSliceUnavailable { + return nil, resp.err + } else if err != nil { + return nil, err + } + continue + } + + // Reduce value. + result = reduceFn(result, resp.result) + + // If all slices have been processed then return. + sliceN += len(resp.slices) + if sliceN >= len(slices) { + return result, nil + } + } + } +} + +func (e *Executor) mapper(ctx context.Context, ch chan mapResponse, nodes []*Node, db string, slices []uint64, c pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc) error { + // Group slices together by nodes. + m, err := e.slicesByNode(nodes, db, slices) + if err != nil { + return err + } + + // Execute each node in a separate goroutine. + for n, nodeSlices := range m { + go func(n *Node, nodeSlices []uint64) { + resp := mapResponse{node: n, slices: nodeSlices} + + // Send local slices to mapper, otherwise remote exec. + if n.Host == e.Host { + resp.result, resp.err = e.mapperLocal(ctx, nodeSlices, mapFn, reduceFn) + } else if !opt.Remote { + results, err := e.exec(ctx, n, db, &pql.Query{Calls: pql.Calls{c}}, nodeSlices, opt) + if len(results) > 0 { + resp.result = results[0] + } + resp.err = err + } + + // Return response to the channel. + select { + case <-ctx.Done(): + case ch <- resp: + } + }(n, nodeSlices) + } + + return nil +} + +// mapperLocal performs map & reduce entirely on the local node. +func (e *Executor) mapperLocal(ctx context.Context, slices []uint64, mapFn mapFunc, reduceFn reduceFunc) (interface{}, error) { + ch := make(chan mapResponse, len(slices)) + + for _, slice := range slices { + go func(slice uint64) { + result, err := mapFn(slice) + + // Return response to the channel. + select { + case <-ctx.Done(): + case ch <- mapResponse{result: result, err: err}: + } + }(slice) + } + + // Reduce results + var sliceN int + var result interface{} + for { + select { + case <-ctx.Done(): + return nil, ctx.Err() + case resp := <-ch: + if resp.err != nil { + return nil, resp.err + } + result = reduceFn(result, resp.result) + sliceN++ + } + + // Exit once all slices are processed. + if sliceN == len(slices) { + return result, nil + } + } +} + +// errSliceUnavailable is a marker error if no nodes are available. +var errSliceUnavailable = errors.New("slice unavailable") + +type mapFunc func(slice uint64) (interface{}, error) + +type reduceFunc func(prev, v interface{}) interface{} + +type mapResponse struct { + node *Node + slices []uint64 + + result interface{} + err error } // ExecOptions represents an execution context for a single Execute() call.