From 5ab025638bcd6104d25cd3181ac846a03c86dba7 Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 15 Apr 2016 12:50:47 -0500 Subject: [PATCH 1/2] fixed typos. changed -config flag to represent prefered Go format --- README.md | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index 333060c54..dfc7e5427 100644 --- a/README.md +++ b/README.md @@ -27,11 +27,11 @@ pilosa ## Configuration -You can specify a configuration by setting the `--config` flag when running `pilosa`. +You can specify a configuration by setting the `-config` flag when running `pilosa`. ```sh -pilosa --config=custom-config-file.cfg +pilosa -config custom-config-file.cfg ``` The config file uses the [TOML](https://github.com/toml-lang/toml) configuration file format, @@ -141,7 +141,7 @@ attributes set using `SetBitmapAttrs()` and `bits` are the bits set using `SetBi ``` Union(Bitmap(id=10, frame="foo"), Bitmap(id=20, frame="foo"))) ``` -Returns a result set similar to that of a `Bitmap()` query, only the `attrs` dictionaly will be empty: `{"results":[{"attrs":{},"bits":[1,2]}]}`. +Returns a result set similar to that of a `Bitmap()` query, only the `attrs` dictionary will be empty: `{"results":[{"attrs":{},"bits":[1,2]}]}`. Note that a `Union()` query can be nested within other queries anywhere that you would otherwise provide a `Bitmap()`. --- @@ -149,7 +149,7 @@ Note that a `Union()` query can be nested within other queries anywhere that you ``` Intersect(Bitmap(id=10, frame="foo"), Bitmap(id=20, frame="foo"))) ``` -Returns a result set similar to that of a `Bitmap()` query, only the `attrs` dictionaly will be empty: `{"results":[{"attrs":{},"bits":[1]}]}`. +Returns a result set similar to that of a `Bitmap()` query, only the `attrs` dictionary will be empty: `{"results":[{"attrs":{},"bits":[1]}]}`. Note that an `Intersect()` query can be nested within other queries anywhere that you would otherwise provide a `Bitmap()`. --- @@ -157,7 +157,7 @@ Note that an `Intersect()` query can be nested within other queries anywhere tha ``` Difference(Bitmap(id=10, frame="foo"), Bitmap(id=20, frame="foo"))) ``` -`Difference()` represents all of the bits that are set in the first `Bitmap()` but are not set in the second `Bitmap()`. It returns a result set similar to that of a `Bitmap()` query, only the `attrs` dictionaly will be empty: `{"results":[{"attrs":{},"bits":[2]}]}`. +`Difference()` represents all of the bits that are set in the first `Bitmap()` but are not set in the second `Bitmap()`. It returns a result set similar to that of a `Bitmap()` query, only the `attrs` dictionary will be empty: `{"results":[{"attrs":{},"bits":[2]}]}`. Note that a `Difference()` query can be nested within other queries anywhere that you would otherwise provide a `Bitmap()`. --- From f29ad7dfa8985bea86bb53e134f827927cfa3696 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Mon, 18 Apr 2016 16:32:27 -0600 Subject: [PATCH 2/2] add anti-entropy monitor --- cmd/pilosa/main.go | 80 +++++++++++++++++++++++++++++++++++++++++++++- executor.go | 5 +++ 2 files changed, 84 insertions(+), 1 deletion(-) diff --git a/cmd/pilosa/main.go b/cmd/pilosa/main.go index fe614c636..5c7971e6a 100644 --- a/cmd/pilosa/main.go +++ b/cmd/pilosa/main.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "io/ioutil" + "log" "math/rand" "net" "net/http" @@ -17,6 +18,7 @@ import ( "runtime/pprof" "strconv" "strings" + "sync" "time" "github.com/BurntSushi/toml" @@ -42,6 +44,9 @@ const ( // DefaultHost is the default hostname and port to use. DefaultHost = "localhost:15000" + + // DefaultAntiEntropyInterval is the default interval to run AAE. + DefaultAntiEntropyInterval = 10 * time.Minute ) func main() { @@ -81,6 +86,10 @@ type Main struct { ticker *time.Ticker pollingSecs int + // Close management. + wg sync.WaitGroup + closing chan struct{} + // Path to the configuration file. ConfigPath string @@ -103,6 +112,8 @@ type Main struct { // NewMain returns a new instance of Main. func NewMain() *Main { return &Main{ + closing: make(chan struct{}), + Config: NewConfig(), Stdin: os.Stdin, Stdout: os.Stdout, @@ -185,6 +196,9 @@ func (m *Main) Run(args ...string) error { // Serve HTTP. go func() { http.Serve(ln, h) }() + // Start anti-entropy background workers. + m.startAntiEntropyMonitors() + // Sync up max slice if more than one node if len(m.Cluster.Nodes) > 1 { m.ticker = time.NewTicker(time.Second * time.Duration(m.pollingSecs)) @@ -212,8 +226,61 @@ func (m *Main) Run(args ...string) error { return nil } -func checkMaxSlice(hostport string) (uint64, error) { +func (m *Main) startAntiEntropyMonitors() { + for _, node := range m.Cluster.Nodes { + // Skip this node. + if node.Host == m.Host { + continue + } + m.wg.Add(1) + go func(node *pilosa.Node) { + defer m.wg.Done() + m.monitorAntiEntropy(node) + }(node) + } +} + +func (m *Main) monitorAntiEntropy(node *pilosa.Node) { + ticker := time.NewTicker(time.Duration(m.Config.AntiEntropy.Interval)) + defer ticker.Stop() + + m.logger().Printf("index sync monitor initializing: host=%s", node.Host) + + for { + // Wait for tick or a close. + select { + case <-m.closing: + return + case <-ticker.C: + } + + m.logger().Printf("index sync beginning: host=%s", node.Host) + + // Set up remote client. + client, err := pilosa.NewClient(node.Host) + if err != nil { + m.logger().Printf("anti-entropy client error: host=%s", node.Host) + continue + } + + // Initialize syncer with local index and remote client. + var syncer pilosa.IndexSyncer + syncer.Index = m.index + syncer.Client = client + + // Sync indexes. + if err := syncer.SyncIndex(); err != nil { + m.logger().Printf("index sync error: host=%s, err=%s", node.Host, err) + continue + } + + // Record successful sync in log. + m.logger().Printf("index sync complete: host=%s", node.Host) + } +} + +func checkMaxSlice(hostport string) (uint64, error) { // Create HTTP request. req, err := http.NewRequest("GET", (&url.URL{ Scheme: "http", @@ -260,6 +327,10 @@ func checkMaxSlice(hostport string) (uint64, error) { // Close shuts down the process. func (m *Main) Close() error { + // Notify goroutines to stop. + close(m.closing) + m.wg.Wait() + if m.ticker != nil { m.ticker.Stop() } @@ -312,6 +383,8 @@ func (m *Main) ParseFlags(args []string) error { return nil } +func (m *Main) logger() *log.Logger { return log.New(m.Stderr, "", log.LstdFlags) } + // Config represents the configuration for the command. type Config struct { DataDir string `toml:"data-dir"` @@ -325,6 +398,10 @@ type Config struct { Plugins struct { Path string `toml:"path"` } `toml:"plugins"` + + AntiEntropy struct { + Interval Duration `toml:"interval"` + } `toml:"anti-entropy"` } type ConfigNode struct { @@ -337,6 +414,7 @@ func NewConfig() *Config { Host: DefaultHost, } c.Cluster.ReplicaN = pilosa.DefaultReplicaN + c.AntiEntropy.Interval = Duration(DefaultAntiEntropyInterval) return c } diff --git a/executor.go b/executor.go index 90f8995ad..e46ed30a9 100644 --- a/executor.go +++ b/executor.go @@ -420,6 +420,11 @@ func (e *Executor) executeSetBit(db string, c *pql.SetBit, opt *ExecOptions) (bo continue } + // Do not forward call if this is already being forwarded. + if opt.Remote { + continue + } + // Forward call to remote node otherwise. if res, err := e.exec(node, db, &pql.Query{Calls: pql.Calls{c}}, nil, opt); err != nil { return false, err