From 3c9a45a0a7978b5491a4d5d664c9be8e78bca15f Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 19 Jun 2014 16:36:39 -0500 Subject: [PATCH] adjusted cassandra config; locked down fragment setup and renable autoaloc --- core/etcd.go | 38 +++++++++++----------- core/service.go | 2 ++ db/constants.go | 63 +++++++++++++++++++++++++++++++++++++ index/fragment_container.go | 8 +---- index/storage_cass.go | 26 +++++++-------- 5 files changed, 97 insertions(+), 40 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index ece14633c..eb808580f 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -21,7 +21,7 @@ type TopologyMapper struct { namespace string } -func (self *TopologyMapper) Run() { +func (self *TopologyMapper) Setup() { log.Println(self.namespace + "/db") db_path := self.namespace + "/db" resp, err := self.service.Etcd.Get(db_path, false, true) @@ -36,12 +36,17 @@ func (self *TopologyMapper) Run() { log.Fatal(err) } } + //need to lock the world for _, node := range flatten(resp.Node) { err := self.handlenode(node) if err != nil { log.Println(err) } } + +} +func (self *TopologyMapper) Run() { + receiver := make(chan *etcd.Response) stop := make(chan bool) go func() { @@ -50,12 +55,12 @@ func (self *TopologyMapper) Run() { for { ns := self.namespace + "/db" log.Println(" ETCD watcher:", ns) - resp, err = self.service.Etcd.Watch(ns, 0, true, receiver, stop) + resp, err := self.service.Etcd.Watch(ns, 0, true, receiver, stop) log.Println("TopologyMapper ETCD watcher", resp, err) } }() go func() { - for resp = range receiver { + for resp := range receiver { switch resp.Action { case "set": self.handlenode(resp.Node) @@ -100,32 +105,31 @@ func getLightestProcess(m map[string]int) (Pair, error) { } func (self *TopologyMapper) MakeFragments(db string, slice_int int) error { - /* - frames_to_create := config.GetStringArrayDefault("supported_frames", []string{"b.n", "l.n", "t.t", "d"}) - for _, frame := range frames_to_create { - err := self.AllocateFragment(db, frame, slice_int) - if err != nil { - log.Println(err) - } + frames_to_create := config.GetStringArrayDefault("supported_frames", []string{"b.n", "l.n", "t.t", "d"}) + for _, frame := range frames_to_create { + err := self.AllocateFragment(db, frame, slice_int) + if err != nil { + log.Println(err) } - */ + } return nil } func (self *TopologyMapper) AllocateFragment(db, frame string, slice_int int) error { //get Lock to create the fragment - ttl := uint64(300) //secs to hold the lock + ttl := uint64(config.GetIntDefault("fragment_alloc_lock_time_secs", 600)) lock_key := fmt.Sprintf("%s/lock/%s-%s-%d", self.namespace, db, frame, slice_int) response, err := self.service.Etcd.RawCreate(lock_key, "0", ttl) if err == nil { - switch response.StatusCode { - case 201: //key created + if response.StatusCode == 201 { //key created //figure out least loaded process..possibly check max process //to create the node, just write off the items to etcd and the watch should spawn //be nice if something would notify perhaps queue m := make(map[string]int) + id_string := self.service.Id.String() + m[id_string] = 0 //at least have one process if none created for _, dbs := range self.service.Cluster.GetDatabases() { for _, fsi := range dbs.GetFramSliceIntersects() { for _, fragment := range fsi.GetFragments() { @@ -152,12 +156,6 @@ func (self *TopologyMapper) AllocateFragment(db, frame string, slice_int int) er // need to check value to see how many we have left _, err = self.service.Etcd.Set(fragment_key, process_guid, 0) log.Printf("Fragment sent to etcd: %s(%s)", fragment_key, process_guid) - case 400: //key already present - return errors.New("Fragment creation already in process:" + lock_key) - default: - log.Println(spew.Sdump(response)) - return errors.New(fmt.Sprintf("Unknown Etcd status:%d", response.StatusCode)) - } } diff --git a/core/service.go b/core/service.go index 317908426..1efd5a15c 100644 --- a/core/service.go +++ b/core/service.go @@ -101,6 +101,8 @@ func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) { } func (service *Service) Run() { + log.Println("Setup service...", service.version) + service.TopologyMapper.Setup() log.Println("Running service...", service.version) go service.TopologyMapper.Run() go service.ProcessMapper.Run() diff --git a/db/constants.go b/db/constants.go index d2fbdd0eb..656be0f29 100644 --- a/db/constants.go +++ b/db/constants.go @@ -1,3 +1,66 @@ package db const SLICE_WIDTH = 65536 + +/* +DEMOGRAPHIC_TILES = { + 'gender': { + 'male': 1342, + 'female': 1343, + }, + 'marital_status': { + 'married': 1352, + 'single': 1353, + }, + 'age_range': { + '18-21': 2936, + '21-24': 1344, + '25-34': 1345, + '35-44': 1346, + '45-54': 1347, + '55-64': 1348, + '65+': 2937, + }, + 'children': { + 'Yes': 1349, + }, + 'education': { + 'Completed High School': 15723, + 'Attended College': 15726, + 'Completed College': 15724, + 'Completed Graduate School': 15725, + }, + 'home_owner_status': { + 'Own': 1350, + 'Rent': 1351, + }, + 'home_market_value': { + '1k-25k': 5700, + '25k-50k': 2954, + '50k-75k': 2955, + '75k-100k': 2956, + '100k-150k': 1360, + '150k-200k': 1364, + '200k-250k': 2957, + '250k-300k': 2958, + '300k-350k': 2959, + '350k-500k': 16048, + '500k-1mm': 1361, + '1mm+': 1359, + }, + 'household_income_range': { + '25k-35k': 2960, + '125k-150k': 16049, + '100k-125k': 16050, + '150k-175k': 1362, + '35k-50k': 16051, + '75k-100k': 1363, + '50k-75k': 16052, + '15k-25k': 2961, + '250k+': 2962, + '200k-250k': 16053, + '0-15k': 16054, + '175k-200k': 16055, + }, +} +*/ diff --git a/index/fragment_container.go b/index/fragment_container.go index af36a0306..3c7503adb 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -244,13 +244,7 @@ func getStorage(db string, slice int, frame string, fid SUUID) Storage { full_dir := fmt.Sprintf("%s/%s/%d/%s/%s", base_path, db, slice, frame, SUUID_to_Hex(fid)) return NewLevelDBStorage(full_dir) case "cassandra": - // host := config.GetString("cassandra_host") - hosts := config.GetStringArrayDefault("cassandra_hosts", []string{"localhost"}) - keyspace := config.GetString("cassandra_keyspace") - if keyspace == "" { - keyspace = "hotbox" - } - return NewCassStorage(hosts, keyspace) + return NewCassStorage() } } diff --git a/index/storage_cass.go b/index/storage_cass.go index 05d051504..4b29490ec 100644 --- a/index/storage_cass.go +++ b/index/storage_cass.go @@ -3,8 +3,8 @@ package index // #cgo CFLAGS:-mpopcnt import ( - "fmt" "log" + "pilosa/config" "pilosa/util" "time" @@ -20,6 +20,17 @@ type CassandraStorage struct { batch_counter int } +var cluster *gocql.ClusterConfig + +func init() { + hosts := config.GetStringArrayDefault("cassandra_hosts", []string{"localhost"}) + keyspace := config.GetStringDefault("cassandra_keyspace", "hotbox") + cluster = gocql.NewCluster(hosts...) + cluster.Keyspace = keyspace + cluster.Consistency = gocql.One + cluster.Timeout = 3 * time.Second +} + func BuildSchema() { /* "CREATE KEYSPACE IF NOT EXISTS hotbox WITH strategy_class = SimpleStrategy AND strategy_options:replication_factor = 1" @@ -29,25 +40,14 @@ func BuildSchema() { */ } -func NewCassStorage(hosts []string, keyspace string) Storage { +func NewCassStorage() Storage { obj := new(CassandraStorage) - //cluster := gocql.NewCluster("127.0.0.1") - cluster := gocql.NewCluster(hosts...) - cluster.Keyspace = keyspace - //cluster.Consistency = gocql.Quorum - cluster.Consistency = gocql.One - cluster.Consistency = gocql.One - cluster.Timeout = 3 * time.Second - //cluster.ProtoVersion = 1 // cluster.CQLVersion = "3.0.0" session, err := cluster.CreateSession() if err != nil { log.Fatal(err) } - err = session.Query(fmt.Sprintf("USE %s", keyspace)).Exec() - if err != nil { - } obj.db = session obj.stmt = `INSERT INTO bitmap ( bitmap_id, db, frame, slice , filter, ChunkKey, BlockIndex, block) VALUES (?,?,?,?,?,?,?,?);` obj.batch = nil