Add Cluster.UnavailableNodes support to skip reads on certain nodes.

This commit is contained in:
Travis 2017-05-08 10:49:41 -05:00
parent 0ce5245e5a
commit 65760496fe
No known key found for this signature in database
GPG key ID: 7F08008DFD9314C9
3 changed files with 21 additions and 3 deletions

View file

@ -83,6 +83,9 @@ func (a Nodes) Clone() []*Node {
type Cluster struct {
Nodes []*Node
// UnavailableNodes do not receive read queries.
UnavailableNodes []*Node
// Hashing algorithm used to assign partitions to nodes.
Hasher Hasher

View file

@ -13,9 +13,10 @@ type Config struct {
Host string `toml:"host"`
Cluster struct {
ReplicaN int `toml:"replicas"`
Nodes []*ConfigNode `toml:"node"`
PollingInterval Duration `toml:"polling-interval"`
ReplicaN int `toml:"replicas"`
Nodes []*ConfigNode `toml:"node"`
UnavailableNodes []*ConfigNode `toml:"unavailable-node"`
PollingInterval Duration `toml:"polling-interval"`
} `toml:"cluster"`
Plugins struct {
@ -59,6 +60,10 @@ func (c *Config) PilosaCluster() *Cluster {
cluster.Nodes = append(cluster.Nodes, &Node{Host: n.Host})
}
for _, u := range c.Cluster.UnavailableNodes {
cluster.UnavailableNodes = append(cluster.UnavailableNodes, &Node{Host: u.Host})
}
return cluster
}

View file

@ -718,6 +718,16 @@ func (e *Executor) mapReduce(ctx context.Context, db string, slices []uint64, c
var nodes []*Node
if !opt.Remote {
nodes = Nodes(e.Cluster.Nodes).Clone()
// If this is a read operation, don't send the query to nodes
// marked as unavailable in the config file.
switch c.(type) {
case pql.BitmapCall, *pql.TopN, *pql.Count:
for _, u := range e.Cluster.UnavailableNodes {
nodes = Nodes(nodes).FilterHost(u.Host)
}
}
} else {
nodes = []*Node{e.Cluster.NodeByHost(e.Host)}
}