mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-10 21:07:53 +00:00
Initial work on query planner, parser, and connection logic.
This commit is contained in:
commit
3f2ff1a0e3
10 changed files with 1011 additions and 0 deletions
25
commands/pilosa-cruncher/cruncher.go
Normal file
25
commands/pilosa-cruncher/cruncher.go
Normal file
|
|
@ -0,0 +1,25 @@
|
|||
package main
|
||||
|
||||
import (
|
||||
"pilosa/cruncher"
|
||||
"pilosa/core"
|
||||
"flag"
|
||||
"log"
|
||||
)
|
||||
|
||||
var locationString string
|
||||
|
||||
func init() {
|
||||
flag.StringVar(&locationString, "l", "127.0.0.1:1300", "ip:port to listen on")
|
||||
flag.Parse()
|
||||
}
|
||||
|
||||
func main() {
|
||||
location, err := core.NewLocation(locationString)
|
||||
if err != nil {
|
||||
log.Fatal("Location not valid:", locationString)
|
||||
}
|
||||
|
||||
cruncher := cruncher.NewCruncher(location)
|
||||
cruncher.Run()
|
||||
}
|
||||
25
commands/pilosa-router/router.go
Normal file
25
commands/pilosa-router/router.go
Normal file
|
|
@ -0,0 +1,25 @@
|
|||
package main
|
||||
|
||||
import (
|
||||
"pilosa/router"
|
||||
"pilosa/core"
|
||||
"flag"
|
||||
"log"
|
||||
)
|
||||
|
||||
var locationString string
|
||||
|
||||
func init() {
|
||||
flag.StringVar(&locationString, "l", "127.0.0.1:1200", "ip:port to listen on")
|
||||
flag.Parse()
|
||||
}
|
||||
|
||||
func main() {
|
||||
location, err := core.NewLocation(locationString)
|
||||
if err != nil {
|
||||
log.Fatal("Location not valid:", locationString)
|
||||
}
|
||||
|
||||
router := router.NewRouter(location)
|
||||
router.Run()
|
||||
}
|
||||
413
core/service.go
Normal file
413
core/service.go
Normal file
|
|
@ -0,0 +1,413 @@
|
|||
package core
|
||||
|
||||
import (
|
||||
"github.com/coreos/go-etcd/etcd"
|
||||
"github.com/coreos/etcd/store"
|
||||
"log"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
"strings"
|
||||
"sync"
|
||||
"net"
|
||||
"encoding/gob"
|
||||
//"net"
|
||||
//"flag"
|
||||
//"encoding/gob"
|
||||
//"net/http"
|
||||
//"encoding/json"
|
||||
//"io"
|
||||
)
|
||||
type Message struct {
|
||||
Key string `json:key`
|
||||
Data interface{} `json:data`
|
||||
Destination Location
|
||||
}
|
||||
|
||||
type Envelope struct {
|
||||
Message *Message
|
||||
Location *Location
|
||||
}
|
||||
|
||||
type Connection struct {
|
||||
Conn net.Conn
|
||||
Encoder *gob.Encoder
|
||||
Decoder *gob.Decoder
|
||||
}
|
||||
|
||||
type Stopper struct {
|
||||
TermChans []chan int
|
||||
DoneChans []chan int
|
||||
Mutex sync.RWMutex
|
||||
}
|
||||
|
||||
func (stopper *Stopper) Stop() {
|
||||
var i chan int
|
||||
var o chan int
|
||||
stopper.Mutex.RLock()
|
||||
for _, i = range stopper.TermChans {
|
||||
go func() {
|
||||
i <- 1
|
||||
}()
|
||||
}
|
||||
for _, o = range stopper.DoneChans {
|
||||
<-o
|
||||
}
|
||||
stopper.Mutex.RUnlock()
|
||||
return
|
||||
}
|
||||
|
||||
func (stopper *Stopper) GetExitChannels() (chan int, chan int) {
|
||||
termchan := make(chan int, 1)
|
||||
donechan := make(chan int, 1)
|
||||
stopper.Mutex.Lock()
|
||||
stopper.TermChans = append(stopper.TermChans, termchan)
|
||||
stopper.DoneChans = append(stopper.DoneChans, donechan)
|
||||
stopper.Mutex.Unlock()
|
||||
return termchan, donechan
|
||||
}
|
||||
|
||||
type Service struct {
|
||||
Stopper
|
||||
Handler func(Message)
|
||||
Mailbox chan Message
|
||||
Inbox chan *Message
|
||||
Port string
|
||||
PortHttp string
|
||||
Etcd *etcd.Client
|
||||
Location *Location
|
||||
NodeMap NodeMap
|
||||
NodeMapMutex sync.RWMutex
|
||||
Outbox chan *Envelope
|
||||
ConnectionMap map[Location]*Connection
|
||||
ConnectionMapMutex sync.RWMutex
|
||||
Listener net.Listener
|
||||
ConnectionRegisterChannel chan *PersistentConnection
|
||||
//Cluster query.Cluster
|
||||
}
|
||||
|
||||
func NewService(location *Location) *Service {
|
||||
service := new(Service)
|
||||
service.Location = location
|
||||
service.Outbox = make(chan *Envelope)
|
||||
service.Inbox = make(chan *Message)
|
||||
return service
|
||||
}
|
||||
|
||||
func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) {
|
||||
hupChan := make(chan os.Signal, 1)
|
||||
termChan := make(chan os.Signal, 1)
|
||||
signal.Notify(hupChan, syscall.SIGHUP)
|
||||
signal.Notify(termChan, syscall.SIGINT, syscall.SIGTERM)
|
||||
return termChan, hupChan
|
||||
}
|
||||
|
||||
func (service *Service) SendMessage(message *Message) error {
|
||||
if message.Destination == *service.Location {
|
||||
return service.DeliverMessage(message)
|
||||
}
|
||||
router := service.GetRouterLocation(message.Destination)
|
||||
var dest *Location
|
||||
if router == *service.Location {
|
||||
dest = &message.Destination
|
||||
} else {
|
||||
dest = &router
|
||||
}
|
||||
return service.DoSendMessage(message, dest)
|
||||
}
|
||||
|
||||
func (service *Service) DeliverMessage(message *Message) error {
|
||||
log.Println("do handle of", message)
|
||||
switch message.Key {
|
||||
case "ping":
|
||||
location, ok := message.Data.(Location)
|
||||
if !ok {
|
||||
log.Println("Invalid ping request!", message)
|
||||
return nil
|
||||
}
|
||||
service.SendMessage(&Message{"pong", nil, location})
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (service *Service) DoSendMessage(message *Message, destination *Location) error {
|
||||
log.Println("send message", message, "to", destination)
|
||||
service.Outbox <- &Envelope{message, destination}
|
||||
return nil
|
||||
}
|
||||
|
||||
type PersistentConnection struct {
|
||||
Outbox chan *Message
|
||||
Inbox chan *Message
|
||||
Encoder *gob.Encoder
|
||||
Decoder *gob.Decoder
|
||||
Location Location
|
||||
Service *Service
|
||||
Connection *net.Conn
|
||||
Identified bool
|
||||
}
|
||||
|
||||
func NewPersistentConnection(service *Service, location *Location) *PersistentConnection {
|
||||
conn := &PersistentConnection{}
|
||||
conn.Service = service
|
||||
conn.Identified = true
|
||||
conn.Outbox = make(chan *Message)
|
||||
conn.Inbox = make(chan *Message)
|
||||
if location != nil {
|
||||
conn.Location = *location
|
||||
}
|
||||
return conn
|
||||
}
|
||||
|
||||
func (conn *PersistentConnection) TryConnect() error {
|
||||
locationString := conn.Location.ToString()
|
||||
log.Println("Dialing", locationString)
|
||||
connection, err := net.Dial("tcp", locationString)
|
||||
if err != nil {
|
||||
log.Println("Error in dial!")
|
||||
return err
|
||||
}
|
||||
conn.Connection = &connection
|
||||
conn.Encoder = gob.NewEncoder(connection)
|
||||
conn.Decoder = gob.NewDecoder(connection)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (conn *PersistentConnection) Manage() {
|
||||
log.Println("Starting manager")
|
||||
var connectionError chan error
|
||||
var err error
|
||||
|
||||
go func() {
|
||||
var message *Message
|
||||
var err error
|
||||
for {
|
||||
err = conn.Decoder.Decode(&message)
|
||||
if err != nil {
|
||||
connectionError <- err
|
||||
return
|
||||
}
|
||||
conn.Inbox <- message
|
||||
}
|
||||
}()
|
||||
|
||||
var message *Message
|
||||
for {
|
||||
select {
|
||||
case message = <-conn.Outbox:
|
||||
log.Println("sending", message)
|
||||
err = conn.Encoder.Encode(message)
|
||||
if err != nil {
|
||||
// if e, ok := err.(*net.OpError); ok {
|
||||
// if e.Err == syscall.EPIPE {
|
||||
// // Client disconnected
|
||||
|
||||
log.Println("error sending", err)
|
||||
go func() { conn.Outbox <- message }() // Resend failed message
|
||||
if !conn.Identified {
|
||||
log.Println("Stopping because connection not identified")
|
||||
return
|
||||
}
|
||||
conn.Connect()
|
||||
}
|
||||
case message = <-conn.Inbox:
|
||||
log.Println("receiving", message)
|
||||
if message.Key == "identify" {
|
||||
log.Println("Registering connection")
|
||||
conn.Location = message.Data.(Location)
|
||||
go conn.Service.RegisterConnection(conn)
|
||||
} else {
|
||||
conn.Service.Inbox <- message
|
||||
}
|
||||
case err = <-connectionError:
|
||||
log.Println("error receiving", err)
|
||||
if !conn.Identified {
|
||||
log.Println("Stopping because connection not identified")
|
||||
return
|
||||
}
|
||||
conn.Connect()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (conn *PersistentConnection) Connect() {
|
||||
log.Println("Connnecting to", conn.Location)
|
||||
err := conn.TryConnect()
|
||||
if err != nil {
|
||||
log.Println(err)
|
||||
log.Println("Connection failed! Waiting 1 second...")
|
||||
time.Sleep(time.Second)
|
||||
// Infinite recursion if node never comes up. Not sure if this is a problem.
|
||||
conn.Connect()
|
||||
} else {
|
||||
log.Println("Send identify message")
|
||||
conn.Encoder.Encode(Message{"identify", *conn.Service.Location, conn.Location})
|
||||
//go func() { conn.Outbox <- &Message{"identify", *conn.Service.Location, conn.Location} }()
|
||||
}
|
||||
log.Println("Connect ending")
|
||||
}
|
||||
|
||||
type ConnectionMapping map[Location]*PersistentConnection
|
||||
|
||||
func (service *Service) RegisterConnection(conn *PersistentConnection) {
|
||||
service.ConnectionRegisterChannel <- conn
|
||||
}
|
||||
func (service *Service) HandleConnections() {
|
||||
connections := ConnectionMapping{}
|
||||
|
||||
log.Println("Handling connections...")
|
||||
for {
|
||||
select {
|
||||
case envelope := <-service.Outbox:
|
||||
conn, ok := connections[*envelope.Location]
|
||||
if !ok {
|
||||
conn = NewPersistentConnection(service, envelope.Location)
|
||||
connections[*envelope.Location] = conn
|
||||
conn.Identified = true
|
||||
go func () {
|
||||
conn.Connect()
|
||||
conn.Manage()
|
||||
}()
|
||||
}
|
||||
// Spawning new goroutine so it doesn't block the main event loop while it's sending.
|
||||
// This emulates an infinitely buffered channel.
|
||||
go func() { conn.Outbox <- envelope.Message }()
|
||||
case conn := <-service.ConnectionRegisterChannel:
|
||||
connections[conn.Location] = conn
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (service *Service) SetupEtcd() {
|
||||
gob.Register(Location{})
|
||||
service.Etcd = etcd.NewClient()
|
||||
service.NodeMapMutex.Lock()
|
||||
defer service.NodeMapMutex.Unlock()
|
||||
service.NodeMap = NodeMap{}
|
||||
|
||||
nodes, err := service.Etcd.Get("nodes")
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
for _, node := range nodes {
|
||||
nodestring := strings.Split(node.Key, "/")[2]
|
||||
location, err := NewLocation(nodestring)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
routerlocation, err := NewLocation(node.Value)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
service.NodeMap[*location] = *routerlocation
|
||||
}
|
||||
log.Println(service.NodeMap)
|
||||
}
|
||||
|
||||
func (service *Service) WatchEtcd() {
|
||||
var receiver = make(chan *store.Response)
|
||||
var stop chan bool
|
||||
go func () {
|
||||
_, err := service.Etcd.Watch("nodes/", 0, receiver, stop)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
}()
|
||||
|
||||
exit, done := service.GetExitChannels()
|
||||
|
||||
for {
|
||||
select {
|
||||
case response := <-receiver:
|
||||
switch response.Action {
|
||||
case "SET":
|
||||
nodestring := strings.Split(response.Key, "/")[2]
|
||||
node, err := NewLocation(nodestring)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
router, err := NewLocation(response.Value)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
service.NodeMapMutex.Lock()
|
||||
service.NodeMap[*node] = *router
|
||||
service.NodeMapMutex.Unlock()
|
||||
case "DELETE":
|
||||
nodestring := strings.Split(response.Key, "/")[2]
|
||||
node, err := NewLocation(nodestring)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
service.NodeMapMutex.Lock()
|
||||
delete(service.NodeMap, *node)
|
||||
service.NodeMapMutex.Unlock()
|
||||
default:
|
||||
log.Println("unhandled etcd message", response)
|
||||
}
|
||||
//log.Println(response.Action, response.Key, response.Value)
|
||||
log.Println(service.NodeMap)
|
||||
case <-exit:
|
||||
log.Println("cleaning up watchetcd service thing.")
|
||||
time.Sleep(2*time.Second)
|
||||
log.Println("done!")
|
||||
done <- 1
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (service *Service) GetRouterLocation(node Location) Location {
|
||||
service.NodeMapMutex.RLock()
|
||||
defer service.NodeMapMutex.RUnlock()
|
||||
location, ok := service.NodeMap[node]
|
||||
if !ok {
|
||||
location, ok = service.NodeMap[*service.Location]
|
||||
if !ok {
|
||||
log.Fatal("Cannot find router!!!")
|
||||
}
|
||||
}
|
||||
return location
|
||||
}
|
||||
|
||||
func (service *Service) SetupNetwork() {
|
||||
log.Println("Setup network")
|
||||
var err error
|
||||
locationString := service.Location.ToString()
|
||||
service.Listener, err = net.Listen("tcp", locationString)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
}
|
||||
func (service *Service) Serve() {
|
||||
for {
|
||||
conn, err := service.Listener.Accept()
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
con := NewPersistentConnection(service, nil)
|
||||
con.Connection = &conn
|
||||
con.Encoder = gob.NewEncoder(*con.Connection)
|
||||
con.Decoder = gob.NewDecoder(*con.Connection)
|
||||
go con.Manage()
|
||||
}
|
||||
}
|
||||
|
||||
//func (app *Application) serveHTTP() {
|
||||
// http.HandleFunc("/message", func(w http.ResponseWriter, r *http.Request) {
|
||||
// if r.Method != "POST" {
|
||||
// http.Error(w, "Only POST allowed", http.StatusMethodNotAllowed)
|
||||
// return
|
||||
// }
|
||||
// var message Message
|
||||
// decoder := json.NewDecoder(r.Body)
|
||||
// if decoder.Decode(&message) != nil {
|
||||
// http.Error(w, "Invalid JSON", http.StatusBadRequest)
|
||||
// return
|
||||
// }
|
||||
// app.Mailbox <- message
|
||||
// })
|
||||
// http.ListenAndServe(app.PortHttp, nil)
|
||||
//}
|
||||
//
|
||||
|
||||
134
core/topology.go
Normal file
134
core/topology.go
Normal file
|
|
@ -0,0 +1,134 @@
|
|||
package core
|
||||
|
||||
import (
|
||||
"github.com/stathat/consistent"
|
||||
"log"
|
||||
"fmt"
|
||||
"errors"
|
||||
"strings"
|
||||
"strconv"
|
||||
)
|
||||
|
||||
var FrameDoesNotExistError = errors.New("Frame does not exist.")
|
||||
|
||||
type Location struct {
|
||||
Ip string
|
||||
Port int
|
||||
}
|
||||
|
||||
// Create a Location struct given a string in form "0.0.0.0:0"
|
||||
func NewLocation(location_string string) (*Location, error) {
|
||||
splitstring := strings.Split(location_string, ":")
|
||||
if len(splitstring) != 2 {
|
||||
return nil, errors.New("Location string must be in form 0.0.0.0:0")
|
||||
}
|
||||
ip := splitstring[0]
|
||||
port, err := strconv.Atoi(splitstring[1])
|
||||
if err != nil{
|
||||
return nil, errors.New("Port is not a number!")
|
||||
}
|
||||
return &Location{ip, port}, nil
|
||||
}
|
||||
|
||||
func (location *Location) ToString() string {
|
||||
return fmt.Sprintf("%s:%d", location.Ip, location.Port)
|
||||
}
|
||||
|
||||
// Map of node location to their router
|
||||
type NodeMap map[Location]Location
|
||||
|
||||
// A fragment is a collection of bitmaps within a slice. The fragment contains a reference to the responsible node for that fragment. The node is in the form ip:port
|
||||
type Fragment struct {
|
||||
Node string
|
||||
}
|
||||
|
||||
// A slice is the vertical combination of every fragment. It contains the hashring used to delegate bitmaps to fragments
|
||||
type Slice struct {
|
||||
Fragments []Fragment
|
||||
Hashring *consistent.Consistent
|
||||
}
|
||||
|
||||
// A frame is a collection of slices in a given category (brands, demographics, etc), specific to a database
|
||||
type Frame struct {
|
||||
Name string
|
||||
Slices []*Slice
|
||||
}
|
||||
|
||||
// Add a slice to a frame with given Node addresses
|
||||
func (f *Frame) AddSlice(addrs ...string) *Slice {
|
||||
slice := Slice{}
|
||||
slice.Hashring = consistent.New()
|
||||
slice.Hashring.NumberOfReplicas = 200
|
||||
sliceIndex := len(f.Slices)
|
||||
for index, addr := range addrs {
|
||||
slice.Fragments = append(slice.Fragments, Fragment{addr})
|
||||
slice.Hashring.Add(fmt.Sprintf("%d %d", sliceIndex, index))
|
||||
}
|
||||
f.Slices = append(f.Slices, &slice)
|
||||
return &slice
|
||||
}
|
||||
|
||||
// Represents the entire cluster, and a reference to the Node this instance is running on
|
||||
type Cluster struct {
|
||||
Databases map[string]*Database
|
||||
Self string
|
||||
}
|
||||
|
||||
// Add a database to a cluster
|
||||
func (c *Cluster) AddDatabase(name string) *Database {
|
||||
database := Database{Name: name}
|
||||
if c.Databases == nil {
|
||||
c.Databases = make(map[string]*Database)
|
||||
}
|
||||
c.Databases[name] = &database
|
||||
return &database
|
||||
}
|
||||
|
||||
// A database is a collection of all the frames within a given profile space
|
||||
type Database struct {
|
||||
Name string
|
||||
Frames []*Frame
|
||||
}
|
||||
|
||||
// Count the number of slices in a database
|
||||
func (d *Database) NumSlices() (int, error) {
|
||||
if len(d.Frames) < 1 {
|
||||
return 0, errors.New("Database is empty")
|
||||
}
|
||||
return len(d.Frames[0].Slices), nil
|
||||
}
|
||||
|
||||
// Add a frame to a database
|
||||
func (d *Database) AddFrame(name string) *Frame {
|
||||
frame := Frame{Name: name}
|
||||
d.Frames = append(d.Frames, &frame)
|
||||
return &frame
|
||||
}
|
||||
|
||||
// Get a frame from a database
|
||||
func (d *Database) GetFrame(name string) (*Frame, error) {
|
||||
for _, frame := range d.Frames {
|
||||
if frame.Name == name {
|
||||
return frame, nil
|
||||
}
|
||||
}
|
||||
return nil, FrameDoesNotExistError
|
||||
}
|
||||
|
||||
// For debugging, prints cluster information
|
||||
func (c *Cluster) Describe() {
|
||||
for _, database := range c.Databases {
|
||||
log.Println("frames", database.Frames)
|
||||
for _, frame := range database.Frames {
|
||||
log.Println(frame.Name, database.Name)
|
||||
for _, slice := range frame.Slices {
|
||||
log.Println(" ", slice)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type Bitmap struct {
|
||||
FrameType string
|
||||
Id int
|
||||
}
|
||||
53
cruncher/cruncher.go
Normal file
53
cruncher/cruncher.go
Normal file
|
|
@ -0,0 +1,53 @@
|
|||
package cruncher
|
||||
|
||||
import (
|
||||
"pilosa/core"
|
||||
"log"
|
||||
)
|
||||
|
||||
type Cruncher struct {
|
||||
core.Service
|
||||
}
|
||||
|
||||
func (c *Cruncher) Init() {
|
||||
log.Println("Initializing cruncher...")
|
||||
}
|
||||
|
||||
func (c *Cruncher) Run() {
|
||||
log.Println("Running cruncher...")
|
||||
c.SetupEtcd()
|
||||
//go r.SyncEtcd()
|
||||
go c.WatchEtcd()
|
||||
go c.HandleConnections()
|
||||
c.SetupNetwork()
|
||||
go c.Serve()
|
||||
go c.HandleInbox()
|
||||
|
||||
sigterm, sighup := c.GetSignals()
|
||||
for {
|
||||
select {
|
||||
case <- sighup:
|
||||
log.Println("SIGHUP! Reloading configuration...")
|
||||
// TODO: reload configuration
|
||||
case <- sigterm:
|
||||
log.Println("SIGTERM! Cleaning up...")
|
||||
c.Stop()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func NewCruncher(location *core.Location) *Cruncher {
|
||||
service := core.NewService(location)
|
||||
cruncher := Cruncher{*service}
|
||||
return &cruncher
|
||||
}
|
||||
|
||||
func (c *Cruncher) HandleInbox() {
|
||||
for {
|
||||
select {
|
||||
case message := <-c.Inbox:
|
||||
log.Println("process", message)
|
||||
}
|
||||
}
|
||||
}
|
||||
61
query/parser.go
Normal file
61
query/parser.go
Normal file
|
|
@ -0,0 +1,61 @@
|
|||
package query
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"pilosa/core"
|
||||
)
|
||||
|
||||
var InvalidQueryError = errors.New("Invalid query format.")
|
||||
|
||||
type QueryParser struct {
|
||||
QueryString string
|
||||
}
|
||||
|
||||
func (q *QueryParser) Walk(data interface{}) (*Query, error) {
|
||||
query := new(Query)
|
||||
|
||||
slice, ok := data.([]interface{})
|
||||
if !ok {
|
||||
return nil, InvalidQueryError
|
||||
}
|
||||
operation, ok := slice[0].(string)
|
||||
|
||||
if !ok {
|
||||
return nil, InvalidQueryError
|
||||
}
|
||||
if operation == "union" || operation == "intersect" {
|
||||
query.Operation = operation
|
||||
inputs := slice[1:]
|
||||
query.Inputs = make([]QueryInput, len(inputs))
|
||||
for idx, input := range inputs {
|
||||
subquery, err := q.Walk(input)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
query.Inputs[idx] = subquery
|
||||
}
|
||||
} else if operation == "bitmap" {
|
||||
query.Operation = "get"
|
||||
frame, ok := slice[1].(string)
|
||||
if !ok {
|
||||
return nil, InvalidQueryError
|
||||
}
|
||||
id, ok := slice[2].(float64)
|
||||
if !ok {
|
||||
return nil, InvalidQueryError
|
||||
}
|
||||
id_int := int(id)
|
||||
query.Inputs = []QueryInput{core.Bitmap{frame, id_int}}
|
||||
}
|
||||
|
||||
return query, nil
|
||||
}
|
||||
|
||||
func (q *QueryParser) Parse() (*Query, error) {
|
||||
var data interface{}
|
||||
if err := json.Unmarshal([]byte(q.QueryString), &data); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return q.Walk(data)
|
||||
}
|
||||
141
query/planner.go
Normal file
141
query/planner.go
Normal file
|
|
@ -0,0 +1,141 @@
|
|||
package query
|
||||
|
||||
import (
|
||||
"github.com/nu7hatch/gouuid"
|
||||
"strconv"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"pilosa/core"
|
||||
)
|
||||
|
||||
// A single step in the query plan.
|
||||
type QueryStep struct {
|
||||
id uuid.UUID
|
||||
operation string
|
||||
inputs []QueryInput
|
||||
location string
|
||||
destination string
|
||||
}
|
||||
|
||||
func (q QueryStep) String() string {
|
||||
return fmt.Sprintf("%s %s %s, LOC: %s, DEST: %s", q.operation, q.id.String(), q.inputs, q.location, q.destination)
|
||||
}
|
||||
|
||||
type QueryInput interface {}
|
||||
|
||||
// Represents a parsed query. Inputs can be Query or Bitmap objects
|
||||
type Query struct{
|
||||
Operation string
|
||||
Inputs []QueryInput // Maybe Bitmap and Query objects should have different fields to avoid using interface{}
|
||||
}
|
||||
|
||||
// This is the output of the query planner. Contains a list of steps which can be performed in parallel
|
||||
type QueryPlan []QueryStep
|
||||
|
||||
type QueryPlanner struct {
|
||||
Cluster *core.Cluster
|
||||
Database *core.Database
|
||||
}
|
||||
|
||||
type QueryTree interface {
|
||||
getLocation(d *core.Database) string
|
||||
}
|
||||
|
||||
// QueryTree for UNION, INTER, and CAT queries
|
||||
type CompositeQueryTree struct {
|
||||
operation string
|
||||
subqueries []QueryTree
|
||||
location string
|
||||
}
|
||||
|
||||
// Randomly select location from subqueries (so subqueries roll up into composite queries while minimizing inter-node data traffic)
|
||||
func (qt *CompositeQueryTree) getLocation(d *core.Database) string {
|
||||
if qt.location == "" {
|
||||
subqueryLength := len(qt.subqueries)
|
||||
if subqueryLength > 1 {
|
||||
locationIndex := rand.Intn(subqueryLength)
|
||||
subquery := qt.subqueries[locationIndex]
|
||||
qt.location = subquery.getLocation(d)
|
||||
}
|
||||
}
|
||||
return qt.location
|
||||
}
|
||||
|
||||
// QueryTree for GET queries
|
||||
type GetQueryTree struct {
|
||||
bitmap core.Bitmap
|
||||
slice int
|
||||
}
|
||||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *GetQueryTree) getLocation(d *core.Database) string {
|
||||
frame, err := d.GetFrame(qt.bitmap.FrameType)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
slice := frame.Slices[qt.slice]
|
||||
hashString, err := slice.Hashring.Get(strconv.Itoa(qt.bitmap.Id))
|
||||
|
||||
var sliceIndex int
|
||||
var fragIndex int
|
||||
fmt.Sscan(hashString, &fragIndex, &sliceIndex)
|
||||
fragment := slice.Fragments[fragIndex]
|
||||
|
||||
return fmt.Sprintf(fragment.Node)
|
||||
}
|
||||
|
||||
// Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices
|
||||
func (qp *QueryPlanner) buildTree(query *Query, slice int) QueryTree {
|
||||
var tree QueryTree
|
||||
if slice == -1 {
|
||||
tree = &CompositeQueryTree{operation:"cat"}
|
||||
numSlices, err := qp.Database.NumSlices()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
for slice := 0; slice < numSlices; slice++ {
|
||||
subtree := qp.buildTree(query, slice)
|
||||
composite := tree.(*CompositeQueryTree)
|
||||
composite.subqueries = append(composite.subqueries, subtree)
|
||||
}
|
||||
} else {
|
||||
if query.Operation == "get" {
|
||||
tree = &GetQueryTree{query.Inputs[0].(core.Bitmap), slice}
|
||||
return tree
|
||||
} else {
|
||||
subqueries := make([]QueryTree, len(query.Inputs))
|
||||
for i, input := range query.Inputs {
|
||||
subqueries[i] = qp.buildTree(input.(*Query), slice)
|
||||
}
|
||||
tree = &CompositeQueryTree{operation:query.Operation, subqueries:subqueries}
|
||||
}
|
||||
}
|
||||
return tree
|
||||
}
|
||||
|
||||
// Produces flattened QueryPlan from QueryTree input
|
||||
func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location string) *QueryPlan {
|
||||
plan := QueryPlan{}
|
||||
if composite, ok := qt.(*CompositeQueryTree); ok {
|
||||
inputs := make([]QueryInput, len(composite.subqueries))
|
||||
step := QueryStep{*id, composite.operation, inputs, composite.getLocation(qp.Database), location}
|
||||
for index, subq := range composite.subqueries {
|
||||
sub_id, _ := uuid.NewV4()
|
||||
step.inputs[index] = []QueryInput{"wait", sub_id}
|
||||
subq_steps := qp.flatten(subq, sub_id, composite.getLocation(qp.Database))
|
||||
plan = append(plan, *subq_steps...)
|
||||
}
|
||||
plan = append(plan, step)
|
||||
} else if get, ok := qt.(*GetQueryTree); ok {
|
||||
step := QueryStep{*id, "get", []QueryInput{get.bitmap, get.slice}, get.getLocation(qp.Database), location}
|
||||
plan := QueryPlan{step}
|
||||
return &plan
|
||||
}
|
||||
return &plan
|
||||
}
|
||||
|
||||
// Transforms Query into QueryTree and flattens to QueryPlan object
|
||||
func (qp *QueryPlanner) Plan(query *Query, id *uuid.UUID, destination string, slice int) *QueryPlan {
|
||||
queryTree := qp.buildTree(query, -1)
|
||||
return qp.flatten(queryTree, id, destination)
|
||||
}
|
||||
67
query/planner_example.go
Normal file
67
query/planner_example.go
Normal file
|
|
@ -0,0 +1,67 @@
|
|||
package query
|
||||
//package main
|
||||
//
|
||||
//import (
|
||||
// "pilosa/query"
|
||||
// "pilosa/core"
|
||||
// "math/rand"
|
||||
// "github.com/nu7hatch/gouuid"
|
||||
// "time"
|
||||
// "log"
|
||||
//)
|
||||
//
|
||||
//func main() {
|
||||
// rand.Seed(time.Now().UnixNano())
|
||||
// cluster := core.Cluster{Self:"192.168.1.100:1201"}
|
||||
// database := cluster.AddDatabase("property49")
|
||||
// frame := database.AddFrame("general")
|
||||
// frame.AddSlice("192.168.1.100:1201", "192.168.1.100:1202", "192.168.1.100:1203")
|
||||
// frame.AddSlice("192.168.1.101:1201", "192.168.1.101:1202", "192.168.1.101:1203")
|
||||
// frame.AddSlice("192.168.1.102:1201", "192.168.1.102:1202", "192.168.1.102:1203")
|
||||
// frame2 := database.AddFrame("brands")
|
||||
// frame2.AddSlice("192.168.1.200:1201", "192.168.1.200:1202", "192.168.1.200:1203")
|
||||
// frame2.AddSlice("192.168.1.201:1201", "192.168.1.201:1202", "192.168.1.201:1203")
|
||||
// frame2.AddSlice("192.168.1.202:1201", "192.168.1.202:1202", "192.168.1.202:1203")
|
||||
//
|
||||
// //cluster.describe()
|
||||
//
|
||||
// //query := Query{"get", []QueryInput{Bitmap{"general", 10}}}
|
||||
//
|
||||
// //query := Query{"union", []QueryInput{
|
||||
// // &Query{"get", []QueryInput{Bitmap{"general", 20}}},
|
||||
// // &Query{"get", []QueryInput{Bitmap{"brands", 30}}},
|
||||
// //}}
|
||||
//
|
||||
// //queryString = "union(bitmap(general, 33), bitmap(brands, 44))"
|
||||
// //queryString := `["union", ["intersect", ["bitmap", "general", 33], ["bitmap", "brands", 44]], ["bitmap", "general", 55]]`
|
||||
// //query := `count(union(bitmap(""), bitmap(brands, 55)))`
|
||||
// //query := `setbit(bitmap(cats, 33), 75000)`
|
||||
// queryString := `["union", ["intersect", ["bitmap", "general", 33], ["bitmap", "brands", 44]], ["bitmap", "general", 55]]`
|
||||
// queryParser := query.QueryParser{queryString}
|
||||
// queryParsed, err := queryParser.Parse()
|
||||
// log.Println(queryParsed)
|
||||
// if err != nil {
|
||||
// log.Println("ERROR!", err)
|
||||
// }
|
||||
//
|
||||
//// query := qp.Query{"inter", []qp.QueryInput{
|
||||
//// &qp.Query{"union", []qp.QueryInput{
|
||||
//// &qp.Query{"get", []qp.QueryInput{core.Bitmap{"general", 20}}},
|
||||
//// &qp.Query{"get", []qp.QueryInput{core.Bitmap{"brands", 30}}},
|
||||
//// }},
|
||||
//// &qp.Query{"get", []qp.QueryInput{core.Bitmap{"general", 10}}},
|
||||
//// }}
|
||||
////
|
||||
// planner := query.QueryPlanner{&cluster, database}
|
||||
//
|
||||
// id, _ := uuid.NewV4()
|
||||
// dest := "1.2.3.4:1234"
|
||||
// log.Println("Query id", id, "Dest", dest)
|
||||
//
|
||||
// queryplan := planner.Plan(queryParsed, id, dest, -1)
|
||||
//
|
||||
// for _, i := range *queryplan {
|
||||
// log.Println(i)
|
||||
// }
|
||||
// //fmt.Println(queryplan)
|
||||
//}
|
||||
10
query/planner_test.go
Normal file
10
query/planner_test.go
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
package query
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"log"
|
||||
)
|
||||
|
||||
func TestQueryPlanner(t *testing.T) {
|
||||
log.Println("query planner test")
|
||||
}
|
||||
82
router/router.go
Normal file
82
router/router.go
Normal file
|
|
@ -0,0 +1,82 @@
|
|||
package router
|
||||
|
||||
import (
|
||||
"log"
|
||||
"time"
|
||||
//"strings"
|
||||
"pilosa/core"
|
||||
)
|
||||
|
||||
type Router struct {
|
||||
core.Service
|
||||
}
|
||||
|
||||
func (r *Router) Init() {
|
||||
log.Println("Initializing router...")
|
||||
}
|
||||
|
||||
func (r *Router) Run() {
|
||||
log.Println("Running router...")
|
||||
r.SetupEtcd()
|
||||
//go r.SyncEtcd()
|
||||
go r.WatchEtcd()
|
||||
go r.HandleConnections()
|
||||
|
||||
go func() {
|
||||
for {
|
||||
r.SendMessage(&core.Message{"ping", core.Location{"127.0.0.1", 1200}, core.Location{"127.0.0.1", 1300}})
|
||||
time.Sleep(2*time.Second)
|
||||
//log.Println(r.GetRouterLocation(core.Location{"127.0.0.1", 1200}))
|
||||
}
|
||||
}()
|
||||
|
||||
sigterm, sighup := r.GetSignals()
|
||||
for {
|
||||
select {
|
||||
case <- sighup:
|
||||
log.Println("SIGHUP! Reloading configuration...")
|
||||
// TODO: reload configuration
|
||||
case <- sigterm:
|
||||
log.Println("SIGTERM! Cleaning up...")
|
||||
r.Stop()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Router) SetupEtcd() {
|
||||
r.Service.SetupEtcd()
|
||||
//routers, err := r.Etcd.Get("topology")
|
||||
//if err != nil {
|
||||
// log.Fatal(err)
|
||||
//}
|
||||
//for _, router := range routers {
|
||||
// routerstring := strings.Split(router.Key, "/")[2]
|
||||
// routerlocation, err := core.NewLocation(routerstring)
|
||||
// if err != nil {
|
||||
// log.Fatal(err)
|
||||
// }
|
||||
// nodes, err := r.Etcd.Get("topology/" + routerstring)
|
||||
// if err != nil {
|
||||
// log.Fatal(err)
|
||||
// }
|
||||
// for _, node := range nodes {
|
||||
// nodestring := strings.Split(node.Key, "/")[3]
|
||||
// nodelocation, err := core.NewLocation(nodestring)
|
||||
// if err != nil {
|
||||
// log.Fatal(err)
|
||||
// }
|
||||
// r.NodeMap[*nodelocation] = *routerlocation
|
||||
// }
|
||||
//}
|
||||
}
|
||||
|
||||
func (r *Router) HandleMessage(m *core.Message) {
|
||||
log.Println(m)
|
||||
}
|
||||
|
||||
func NewRouter(location *core.Location) *Router {
|
||||
service := core.NewService(location)
|
||||
router := Router{*service}
|
||||
return &router
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue