mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-09 06:31:02 +00:00
standardizes on MaxSlice (instead of SliceN)
creates DB locally even if node doesn't have data for that DB
This commit is contained in:
parent
4d1c12b4d7
commit
05f65ef09c
8 changed files with 44 additions and 40 deletions
|
|
@ -45,8 +45,8 @@ func NewClient(host string) (*Client, error) {
|
|||
// Host returns the host the client was initialized with.
|
||||
func (c *Client) Host() string { return c.host }
|
||||
|
||||
// SliceNs returns the number of slices on a server by database.
|
||||
func (c *Client) SliceNs(ctx context.Context) (MaxSlices, error) {
|
||||
// MaxSliceByDatabase returns the number of slices on a server by database.
|
||||
func (c *Client) MaxSliceByDatabase(ctx context.Context) (map[string]uint64, error) {
|
||||
// Execute request against the host.
|
||||
u := url.URL{
|
||||
Scheme: "http",
|
||||
|
|
@ -362,13 +362,13 @@ func (c *Client) BackupTo(ctx context.Context, w io.Writer, db, frame string) er
|
|||
tw := tar.NewWriter(w)
|
||||
|
||||
// Find the maximum number of slices.
|
||||
sliceNs, err := c.SliceNs(ctx)
|
||||
maxSlices, err := c.MaxSliceByDatabase(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("slice n: %s", err)
|
||||
}
|
||||
|
||||
// Backup every slice to the tar file.
|
||||
for i := uint64(0); i <= sliceNs[db]; i++ {
|
||||
for i := uint64(0); i <= maxSlices[db]; i++ {
|
||||
if err := c.backupSliceTo(ctx, tw, db, frame, i); err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
|
|||
|
|
@ -480,13 +480,13 @@ func (cmd *ExportCommand) Run(ctx context.Context) error {
|
|||
}
|
||||
|
||||
// Determine slice count.
|
||||
sliceNs, err := client.SliceNs(ctx)
|
||||
maxSlices, err := client.MaxSliceByDatabase(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Export each slice.
|
||||
for slice := uint64(0); slice <= sliceNs[cmd.Database]; slice++ {
|
||||
for slice := uint64(0); slice <= maxSlices[cmd.Database]; slice++ {
|
||||
logger.Printf("exporting slice: %d", slice)
|
||||
if err := client.ExportCSV(ctx, cmd.Database, cmd.Frame, slice, w); err != nil {
|
||||
return err
|
||||
|
|
|
|||
6
db.go
6
db.go
|
|
@ -116,8 +116,8 @@ func (db *DB) Close() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// SliceN returns the max slice in the database according to this node.
|
||||
func (db *DB) SliceN() uint64 {
|
||||
// MaxSlice returns the max slice in the database according to this node.
|
||||
func (db *DB) MaxSlice() uint64 {
|
||||
if db == nil {
|
||||
return 0
|
||||
}
|
||||
|
|
@ -126,7 +126,7 @@ func (db *DB) SliceN() uint64 {
|
|||
|
||||
max := db.remoteMaxSlice
|
||||
for _, f := range db.frames {
|
||||
if slice := f.SliceN(); slice > max {
|
||||
if slice := f.MaxSlice(); slice > max {
|
||||
max = slice
|
||||
}
|
||||
}
|
||||
|
|
|
|||
16
executor.go
16
executor.go
|
|
@ -53,10 +53,10 @@ func (e *Executor) Execute(ctx context.Context, db string, q *pql.Query, slices
|
|||
// If slices aren't specified, then include all of them.
|
||||
if len(slices) == 0 {
|
||||
// Round up the number of slices.
|
||||
sliceN := e.Index.DB(db).SliceN()
|
||||
maxSlice := e.Index.DB(db).MaxSlice()
|
||||
|
||||
// Generate a slices of all slices.
|
||||
slices = make([]uint64, sliceN+1)
|
||||
slices = make([]uint64, maxSlice+1)
|
||||
for i := range slices {
|
||||
slices[i] = uint64(i)
|
||||
}
|
||||
|
|
@ -713,7 +713,7 @@ func (e *Executor) mapReduce(ctx context.Context, db string, slices []uint64, c
|
|||
|
||||
// Iterate over all map responses and reduce.
|
||||
var result interface{}
|
||||
var sliceN int
|
||||
var maxSlice int
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
|
|
@ -738,8 +738,8 @@ func (e *Executor) mapReduce(ctx context.Context, db string, slices []uint64, c
|
|||
result = reduceFn(result, resp.result)
|
||||
|
||||
// If all slices have been processed then return.
|
||||
sliceN += len(resp.slices)
|
||||
if sliceN >= len(slices) {
|
||||
maxSlice += len(resp.slices)
|
||||
if maxSlice >= len(slices) {
|
||||
return result, nil
|
||||
}
|
||||
}
|
||||
|
|
@ -797,7 +797,7 @@ func (e *Executor) mapperLocal(ctx context.Context, slices []uint64, mapFn mapFu
|
|||
}
|
||||
|
||||
// Reduce results
|
||||
var sliceN int
|
||||
var maxSlice int
|
||||
var result interface{}
|
||||
for {
|
||||
select {
|
||||
|
|
@ -808,11 +808,11 @@ func (e *Executor) mapperLocal(ctx context.Context, slices []uint64, mapFn mapFu
|
|||
return nil, resp.err
|
||||
}
|
||||
result = reduceFn(result, resp.result)
|
||||
sliceN++
|
||||
maxSlice++
|
||||
}
|
||||
|
||||
// Exit once all slices are processed.
|
||||
if sliceN == len(slices) {
|
||||
if maxSlice == len(slices) {
|
||||
return result, nil
|
||||
}
|
||||
}
|
||||
|
|
|
|||
8
frame.go
8
frame.go
|
|
@ -58,8 +58,8 @@ func (f *Frame) Path() string { return f.path }
|
|||
// BitmapAttrStore returns the attribute storage.
|
||||
func (f *Frame) BitmapAttrStore() *AttrStore { return f.bitmapAttrStore }
|
||||
|
||||
// SliceN returns the max slice in the frame.
|
||||
func (f *Frame) SliceN() uint64 {
|
||||
// MaxSlice returns the max slice in the frame.
|
||||
func (f *Frame) MaxSlice() uint64 {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
|
|
@ -128,7 +128,7 @@ func (f *Frame) openFragments() error {
|
|||
frag.BitmapAttrStore = f.bitmapAttrStore
|
||||
f.fragments[frag.Slice()] = frag
|
||||
|
||||
f.stats.Count("sliceN", 1)
|
||||
f.stats.Count("maxSlice", 1)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
@ -202,7 +202,7 @@ func (f *Frame) createFragmentIfNotExists(slice uint64) (*Fragment, error) {
|
|||
// Save to lookup.
|
||||
f.fragments[slice] = frag
|
||||
|
||||
f.stats.Count("sliceN", 1)
|
||||
f.stats.Count("maxSlice", 1)
|
||||
|
||||
return frag, nil
|
||||
}
|
||||
|
|
|
|||
12
handler.go
12
handler.go
|
|
@ -103,7 +103,7 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
|||
case "/slices/max":
|
||||
switch r.Method {
|
||||
case "GET":
|
||||
h.handleGetMaxSlices(w, r)
|
||||
h.handleGetSliceMax(w, r)
|
||||
default:
|
||||
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
|
||||
}
|
||||
|
|
@ -252,8 +252,8 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
}
|
||||
|
||||
func (h *Handler) handleGetMaxSlices(w http.ResponseWriter, r *http.Request) error {
|
||||
ms := h.Index.SliceNs()
|
||||
func (h *Handler) handleGetSliceMax(w http.ResponseWriter, r *http.Request) error {
|
||||
ms := h.Index.MaxSlices()
|
||||
if strings.Contains(r.Header.Get("Accept"), "application/x-protobuf") {
|
||||
pb := &internal.MaxSlicesResponse{
|
||||
MaxSlices: ms,
|
||||
|
|
@ -271,7 +271,7 @@ func (h *Handler) handleGetMaxSlices(w http.ResponseWriter, r *http.Request) err
|
|||
}
|
||||
|
||||
type sliceMaxResponse struct {
|
||||
MaxSlices MaxSlices `json:"MaxSlices"`
|
||||
MaxSlices map[string]uint64 `json:"MaxSlices"`
|
||||
}
|
||||
|
||||
// handleDeleteDB handles DELETE /db request.
|
||||
|
|
@ -818,7 +818,7 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request)
|
|||
}
|
||||
|
||||
// Determine the maximum number of slices.
|
||||
sliceNs, err := client.SliceNs(r.Context())
|
||||
maxSlices, err := client.MaxSliceByDatabase(r.Context())
|
||||
if err != nil {
|
||||
http.Error(w, "cannot determine remote slice count: "+err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
|
|
@ -826,7 +826,7 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request)
|
|||
|
||||
// Loop over each slice and import it if this node owns it.
|
||||
//travis
|
||||
for slice := uint64(0); slice <= sliceNs[db]; slice++ {
|
||||
for slice := uint64(0); slice <= maxSlices[db]; slice++ {
|
||||
// Ignore this slice if we don't own it.
|
||||
if !h.Cluster.OwnsFragment(h.Host, db, slice) {
|
||||
continue
|
||||
|
|
|
|||
13
index.go
13
index.go
|
|
@ -106,14 +106,11 @@ func (i *Index) Close() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// MaxSlices contains the max known slice -by this node- for each DB
|
||||
type MaxSlices map[string]uint64
|
||||
|
||||
// SliceNs returns MaxSlice map for all databases.
|
||||
func (i *Index) SliceNs() MaxSlices {
|
||||
a := make(MaxSlices)
|
||||
// MaxSlices returns MaxSlice map for all databases.
|
||||
func (i *Index) MaxSlices() map[string]uint64 {
|
||||
a := make(map[string]uint64)
|
||||
for _, db := range i.DBs() {
|
||||
a[db.Name()] = db.SliceN()
|
||||
a[db.Name()] = db.MaxSlice()
|
||||
}
|
||||
return a
|
||||
}
|
||||
|
|
@ -345,7 +342,7 @@ func (s *IndexSyncer) SyncIndex() error {
|
|||
return fmt.Errorf("frame sync error: db=%s, frame=%s, err=%s", di.Name, fi.Name, err)
|
||||
}
|
||||
|
||||
for slice := uint64(0); slice <= s.Index.DB(di.Name).SliceN(); slice++ {
|
||||
for slice := uint64(0); slice <= s.Index.DB(di.Name).MaxSlice(); slice++ {
|
||||
// Ignore slices that this host doesn't own.
|
||||
if !s.Cluster.OwnsFragment(s.Host, di.Name, slice) {
|
||||
continue
|
||||
|
|
|
|||
17
server.go
17
server.go
|
|
@ -197,19 +197,26 @@ func (s *Server) monitorMaxSlices() {
|
|||
case <-ticker.C:
|
||||
}
|
||||
|
||||
oldmaxslices := s.Index.SliceNs()
|
||||
oldmaxslices := s.Index.MaxSlices()
|
||||
for _, node := range s.Cluster.Nodes {
|
||||
if s.Host != node.Host {
|
||||
maxSlices, _ := checkMaxSlices(node.Host)
|
||||
for db, newmax := range maxSlices {
|
||||
// we're not going to create a db locally if we don't know about it already
|
||||
// TODO: consider changing this so we DO create a db locally
|
||||
// do we want/need nodes to have empty files structures?
|
||||
// if we don't know about a db locally, create it
|
||||
// so that the /schema endpoint can report it
|
||||
if localdb := s.Index.DB(db); localdb != nil {
|
||||
if newmax > oldmaxslices[db] {
|
||||
oldmaxslices[db] = newmax
|
||||
localdb.SetRemoteMaxSlice(newmax)
|
||||
}
|
||||
} else {
|
||||
d, err := s.Index.CreateDBIfNotExists(db)
|
||||
if err != nil {
|
||||
s.logger().Printf("Failed to create DB locally: %s", db)
|
||||
return
|
||||
}
|
||||
oldmaxslices[db] = newmax
|
||||
d.SetRemoteMaxSlice(newmax)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -217,7 +224,7 @@ func (s *Server) monitorMaxSlices() {
|
|||
}
|
||||
}
|
||||
|
||||
func checkMaxSlices(hostport string) (MaxSlices, error) {
|
||||
func checkMaxSlices(hostport string) (map[string]uint64, error) {
|
||||
// Create HTTP request.
|
||||
req, err := http.NewRequest("GET", (&url.URL{
|
||||
Scheme: "http",
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue