adjusted cassandra config; locked down fragment setup and renable autoaloc

This commit is contained in:
Todd Gruben 2014-06-19 16:36:39 -05:00
parent 574b700634
commit 3c9a45a0a7
5 changed files with 97 additions and 40 deletions

View file

@ -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))
}
}

View file

@ -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()

View file

@ -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,
},
}
*/

View file

@ -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()
}
}

View file

@ -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