From 7284c4dd10084d84f18f05999dfbefd3cd405b14 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 6 May 2021 11:17:38 -0600 Subject: [PATCH] Add id alloc, col attrs, & row attrs backup --- api.go | 23 +++++++++++ attr.go | 6 +++ boltdb/attrstore.go | 12 +++++- client.go | 16 ++++++++ ctl/backup.go | 98 +++++++++++++++++++++++++++++++++++++++++++++ http/client.go | 68 +++++++++++++++++++++++++++++++ http/handler.go | 29 ++++++++++++++ idalloc.go | 15 +++++++ 8 files changed, 266 insertions(+), 1 deletion(-) diff --git a/api.go b/api.go index 6516f56de..c40cc1a93 100644 --- a/api.go +++ b/api.go @@ -280,6 +280,15 @@ func (api *API) DeleteIndex(ctx context.Context, indexName string) error { return nil } +func (api *API) WriteColumnAttrDataTo(ctx context.Context, w io.Writer, indexName string) error { + index := api.holder.Index(indexName) + if index == nil { + return newNotFoundError(ErrIndexNotFound, indexName) + } + _, err := index.ColumnAttrStore().WriteTo(w) + return err +} + // CreateField makes the named field in the named index with the given options. // This method currently only takes a single functional option, but that may be // changed in the future to support multiple options. @@ -337,6 +346,15 @@ func (api *API) Field(ctx context.Context, indexName, fieldName string) (*Field, return field, nil } +func (api *API) WriteRowAttrDataTo(ctx context.Context, w io.Writer, indexName, fieldName string) error { + field := api.holder.Field(indexName, fieldName) + if field == nil { + return newNotFoundError(ErrFieldNotFound, fieldName) + } + _, err := field.RowAttrStore().WriteTo(w) + return err +} + func setUpImportOptions(opts ...ImportOption) (*ImportOptions, error) { options := &ImportOptions{} for _, opt := range opts { @@ -2297,6 +2315,11 @@ func (api *API) ResetIDAlloc(index string) error { return api.holder.ida.reset(index) } +func (api *API) WriteIDAllocDataTo(w io.Writer) error { + _, err := api.holder.ida.WriteTo(w) + return err +} + // TranslateIndexDB is an internal function to load the index keys database // rd is a boltdb file. func (api *API) TranslateIndexDB(ctx context.Context, indexName string, partitionID int, rd io.Reader) error { diff --git a/attr.go b/attr.go index 66667f806..fa5d11eca 100644 --- a/attr.go +++ b/attr.go @@ -16,6 +16,7 @@ package pilosa import ( "bytes" + "io" "sort" "github.com/gogo/protobuf/proto" @@ -32,6 +33,8 @@ const ( // AttrStore represents an interface for handling row/column attributes. type AttrStore interface { + io.WriterTo + Path() string Open() error Close() error @@ -76,6 +79,9 @@ func (s nopAttrStore) Blocks() ([]AttrBlock, error) { return nil, nil } // BlockData is a no-op implementation of AttrStore BlockData method. func (s nopAttrStore) BlockData(i uint64) (map[uint64]map[string]interface{}, error) { return nil, nil } +// WriteTo is a no-op implementation of AttrStore WriteTo method. +func (s nopAttrStore) WriteTo(w io.Writer) (int64, error) { return 0, nil } + // AttrBlock represents a checksummed block of the attribute store. type AttrBlock struct { ID uint64 `json:"id"` diff --git a/boltdb/attrstore.go b/boltdb/attrstore.go index cf5c62473..4c74a70a2 100644 --- a/boltdb/attrstore.go +++ b/boltdb/attrstore.go @@ -16,9 +16,9 @@ package boltdb import ( "bytes" - "encoding/binary" "fmt" + "io" "sort" "sync" "time" @@ -276,6 +276,16 @@ func (s *attrStore) BlockData(i uint64) (m map[uint64]map[string]interface{}, er return m, nil } +// WriteTo writes the underlying database to w. +func (s *attrStore) WriteTo(w io.Writer) (int64, error) { + tx, err := s.db.Begin(false) + if err != nil { + return 0, err + } + defer tx.Rollback() + return tx.WriteTo(w) +} + // txAttrs returns a map of attributes for an id. func txAttrs(tx *bolt.Tx, id uint64) (map[string]interface{}, error) { v := tx.Bucket([]byte("attrs")).Get(u64tob(id)) diff --git a/client.go b/client.go index 70736bf60..9985e1342 100644 --- a/client.go +++ b/client.go @@ -82,8 +82,12 @@ type InternalClient interface { ImportRoaring(ctx context.Context, uri *pnet.URI, index, field string, shard uint64, remote bool, req *ImportRoaringRequest) error ImportColumnAttrs(ctx context.Context, uri *pnet.URI, index string, req *ImportColumnAttrsRequest) error ShardReader(ctx context.Context, index string, shard uint64) (io.ReadCloser, error) + + IDAllocDataReader(ctx context.Context) (io.ReadCloser, error) IndexTranslateDataReader(ctx context.Context, index string, partitionID int) (io.ReadCloser, error) + IndexAttrDataReader(ctx context.Context, index string) (io.ReadCloser, error) FieldTranslateDataReader(ctx context.Context, index, field string) (io.ReadCloser, error) + FieldAttrDataReader(ctx context.Context, index, field string) (io.ReadCloser, error) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool) (*Transaction, error) FinishTransaction(ctx context.Context, id string) (*Transaction, error) @@ -209,14 +213,26 @@ func (n nopInternalClient) ShardReader(ctx context.Context, index string, shard return nil, nil } +func (n nopInternalClient) IDAllocDataReader(ctx context.Context) (io.ReadCloser, error) { + return nil, nil +} + func (n nopInternalClient) IndexTranslateDataReader(ctx context.Context, index string, partitionID int) (io.ReadCloser, error) { return nil, nil } +func (n nopInternalClient) IndexAttrDataReader(ctx context.Context, index string) (io.ReadCloser, error) { + return nil, nil +} + func (n nopInternalClient) FieldTranslateDataReader(ctx context.Context, index, field string) (io.ReadCloser, error) { return nil, nil } +func (n nopInternalClient) FieldAttrDataReader(ctx context.Context, index, field string) (io.ReadCloser, error) { + return nil, nil +} + func (n nopInternalClient) EnsureIndex(ctx context.Context, name string, options IndexOptions) error { return nil } diff --git a/ctl/backup.go b/ctl/backup.go index ffab74fb8..403860656 100644 --- a/ctl/backup.go +++ b/ctl/backup.go @@ -144,6 +144,37 @@ func (cmd *BackupCommand) backupSchema(ctx context.Context, tw *tar.Writer, sche return nil } +func (cmd *BackupCommand) backupIDAllocData(ctx context.Context, tw *tar.Writer) error { + logger := cmd.Logger() + logger.Printf("backing up id alloc data") + + rc, err := cmd.client.IDAllocDataReader(ctx) + if err != nil { + return fmt.Errorf("fetching id alloc data reader: %w", err) + } + defer rc.Close() + + // Read to buffer to determine size. + var buf bytes.Buffer + if _, err := buf.ReadFrom(rc); err != nil { + return fmt.Errorf("copying id alloc data to memory: %w", err) + } + + // Build header & copy data to archive. + if err = tw.WriteHeader(&tar.Header{ + Name: "idalloc", + Mode: 0666, + Size: int64(buf.Len()), + ModTime: time.Now(), + }); err != nil { + return err + } else if _, err := io.Copy(tw, &buf); err != nil { + return fmt.Errorf("copying id alloc data to archive: %w", err) + } + + return nil +} + // backupIndex backs up all shards for a given index. func (cmd *BackupCommand) backupIndex(ctx context.Context, tw *tar.Writer, ii *pilosa.IndexInfo) error { logger := cmd.Logger() @@ -165,10 +196,18 @@ func (cmd *BackupCommand) backupIndex(ctx context.Context, tw *tar.Writer, ii *p if err := cmd.backupIndexTranslateData(ctx, tw, ii.Name); err != nil { return err } + if err := cmd.backupIndexAttrData(ctx, tw, ii.Name); err != nil { + return err + } + + // Back up field translation & attribute data. for _, fi := range ii.Fields { if err := cmd.backupFieldTranslateData(ctx, tw, ii.Name, fi.Name); err != nil { return fmt.Errorf("cannot backup field translation data for field %q on index %q: %w", fi.Name, ii.Name, err) } + if err := cmd.backupFieldAttrData(ctx, tw, ii.Name, fi.Name); err != nil { + return fmt.Errorf("cannot backup field attr data for field %q on index %q: %w", fi.Name, ii.Name, err) + } } return nil @@ -253,6 +292,36 @@ func (cmd *BackupCommand) backupIndexPartitionTranslateData(ctx context.Context, return nil } +func (cmd *BackupCommand) backupIndexAttrData(ctx context.Context, tw *tar.Writer, name string) error { + logger := cmd.Logger() + logger.Printf("backing up index attr data: %s", name) + + rc, err := cmd.client.IndexAttrDataReader(ctx, name) + if err != nil { + return fmt.Errorf("fetching index attr data reader: %w", err) + } + defer rc.Close() + + // Read to buffer to determine size. + var buf bytes.Buffer + if _, err := buf.ReadFrom(rc); err != nil { + return fmt.Errorf("copying index attr data to memory: %w", err) + } + + // Build header & copy data to archive. + if err = tw.WriteHeader(&tar.Header{ + Name: path.Join("indexes", name, "attributes"), + Mode: 0666, + Size: int64(buf.Len()), + ModTime: time.Now(), + }); err != nil { + return err + } else if _, err := io.Copy(tw, &buf); err != nil { + return fmt.Errorf("copying index attr data to archive: %w", err) + } + return nil +} + func (cmd *BackupCommand) backupFieldTranslateData(ctx context.Context, tw *tar.Writer, indexName, fieldName string) error { logger := cmd.Logger() logger.Printf("backing up field translation data: %s/%s", indexName, fieldName) @@ -282,7 +351,36 @@ func (cmd *BackupCommand) backupFieldTranslateData(ctx context.Context, tw *tar. } else if _, err := io.Copy(tw, &buf); err != nil { return fmt.Errorf("copying translate data to archive: %w", err) } + return nil +} +func (cmd *BackupCommand) backupFieldAttrData(ctx context.Context, tw *tar.Writer, indexName, fieldName string) error { + logger := cmd.Logger() + logger.Printf("backing up field attr data: %s/%s", indexName, fieldName) + + rc, err := cmd.client.FieldAttrDataReader(ctx, indexName, fieldName) + if err != nil { + return fmt.Errorf("fetching field attr data reader: %w", err) + } + defer rc.Close() + + // Read to buffer to determine size. + var buf bytes.Buffer + if _, err := buf.ReadFrom(rc); err != nil { + return fmt.Errorf("copying field attr data to memory: %w", err) + } + + // Build header & copy data to archive. + if err = tw.WriteHeader(&tar.Header{ + Name: path.Join("indexes", indexName, "fields", fieldName, "attributes"), + Mode: 0666, + Size: int64(buf.Len()), + ModTime: time.Now(), + }); err != nil { + return err + } else if _, err := io.Copy(tw, &buf); err != nil { + return fmt.Errorf("copying field attr data to archive: %w", err) + } return nil } diff --git a/http/client.go b/http/client.go index dbca3c501..fdf2cf7d7 100644 --- a/http/client.go +++ b/http/client.go @@ -2136,6 +2136,28 @@ func (c *InternalClient) ShardReader(ctx context.Context, index string, shard ui return resp.Body, nil } +// IDAllocDataReader returns a reader that provides a snapshot of ID allocation data. +func (c *InternalClient) IDAllocDataReader(ctx context.Context) (io.ReadCloser, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.IDAllocDataReader") + defer span.Finish() + + // Build request. + req, err := http.NewRequest("GET", c.defaultURI.String()+"/internal/idalloc/data", nil) + if err != nil { + return nil, errors.Wrap(err, "creating request") + } + + req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) + req.Header.Set("Accept", "application/octet-stream") + + // Execute request. + resp, err := c.executeRequest(req.WithContext(ctx)) + if err != nil { + return nil, err + } + return resp.Body, nil +} + // IndexTranslateDataReader returns a reader that provides a snapshot of // translation data for a partition in an index. func (c *InternalClient) IndexTranslateDataReader(ctx context.Context, index string, partitionID int) (io.ReadCloser, error) { @@ -2165,6 +2187,29 @@ func (c *InternalClient) IndexTranslateDataReader(ctx context.Context, index str return resp.Body, nil } +// IndexAttrDataReader returns a reader that provides a snapshot of column attributes data. +func (c *InternalClient) IndexAttrDataReader(ctx context.Context, index string) (io.ReadCloser, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.IndexAttrDataReader") + defer span.Finish() + + // Build request. + u := fmt.Sprintf("%s/internal/index/%s/attr/data", c.defaultURI.String(), url.QueryEscape(index)) + req, err := http.NewRequest("GET", u, nil) + if err != nil { + return nil, errors.Wrap(err, "creating request") + } + + req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) + req.Header.Set("Accept", "application/octet-stream") + + // Execute request. + resp, err := c.executeRequest(req.WithContext(ctx)) + if err != nil { + return nil, err + } + return resp.Body, nil +} + // FieldTranslateDataReader returns a reader that provides a snapshot of // translation data for a field. func (c *InternalClient) FieldTranslateDataReader(ctx context.Context, index, field string) (io.ReadCloser, error) { @@ -2194,6 +2239,29 @@ func (c *InternalClient) FieldTranslateDataReader(ctx context.Context, index, fi return resp.Body, nil } +// FieldAttrDataReader returns a reader that provides a snapshot of row attributes data. +func (c *InternalClient) FieldAttrDataReader(ctx context.Context, index, field string) (io.ReadCloser, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.FieldAttrDataReader") + defer span.Finish() + + // Build request. + u := fmt.Sprintf("%s/internal/index/%s/field/%s/attr/data", c.defaultURI.String(), url.QueryEscape(index), url.QueryEscape(field)) + req, err := http.NewRequest("GET", u, nil) + if err != nil { + return nil, errors.Wrap(err, "creating request") + } + + req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) + req.Header.Set("Accept", "application/octet-stream") + + // Execute request. + resp, err := c.executeRequest(req.WithContext(ctx)) + if err != nil { + return nil, err + } + return resp.Body, nil +} + // Status function is just a public function for this particular implementation of InternalClient. // It's not require by pilosa.InternalClient interface. // The function returns pilosa cluster state as a string ("NORMAL", "DEGRADED", "DOWN", "RESIZING", ...) diff --git a/http/handler.go b/http/handler.go index ae1a43706..94f14a011 100644 --- a/http/handler.go +++ b/http/handler.go @@ -250,7 +250,9 @@ func (h *Handler) populateValidators() { h.validators["GetFragmentData"] = queryValidationSpecRequired("index", "field", "view", "shard") h.validators["GetFragmentNodes"] = queryValidationSpecRequired("shard", "index") h.validators["PostIndexAttrDiff"] = queryValidationSpecRequired() + h.validators["GetIndexAttrData"] = queryValidationSpecRequired() h.validators["PostFieldAttrDiff"] = queryValidationSpecRequired() + h.validators["GetFieldAttrData"] = queryValidationSpecRequired() h.validators["GetNodes"] = queryValidationSpecRequired() h.validators["GetShardMax"] = queryValidationSpecRequired() h.validators["GetTransactionList"] = queryValidationSpecRequired() @@ -425,12 +427,14 @@ func newRouter(handler *Handler) http.Handler { router.HandleFunc("/internal/fragment/data", handler.handleGetFragmentData).Methods("GET").Name("GetFragmentData") router.HandleFunc("/internal/fragment/nodes", handler.handleGetFragmentNodes).Methods("GET").Name("GetFragmentNodes") router.HandleFunc("/internal/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST").Name("PostIndexAttrDiff") + router.HandleFunc("/internal/index/{index}/attr/data", handler.handleGetIndexAttrData).Methods("GET").Name("GetIndexAttrData") router.HandleFunc("/internal/translate/data", handler.handleGetTranslateData).Methods("GET").Name("GetTranslateData") router.HandleFunc("/internal/translate/data", handler.handlePostTranslateData).Methods("POST").Name("PostTranslateData") router.HandleFunc("/internal/translate/keys", handler.handlePostTranslateKeys).Methods("POST").Name("PostTranslateKeys") router.HandleFunc("/internal/translate/ids", handler.handlePostTranslateIDs).Methods("POST").Name("PostTranslateIDs") router.HandleFunc("/internal/index/{index}/field/{field}/attr/diff", handler.handlePostFieldAttrDiff).Methods("POST").Name("PostFieldAttrDiff") router.HandleFunc("/internal/index/{index}/field/{field}/remote-available-shards/{shardID}", handler.handleDeleteRemoteAvailableShard).Methods("DELETE") + router.HandleFunc("/internal/index/{index}/field/{field}/attr/data", handler.handleGetFieldAttrData).Methods("GET").Name("GetFieldAttrData") router.HandleFunc("/internal/index/{index}/shard/{shard}/snapshot", handler.handleGetIndexShardSnapshot).Methods("GET").Name("GetIndexShardSnapshot") router.HandleFunc("/internal/index/{index}/shards", handler.handleGetIndexAvailableShards).Methods("GET").Name("GetIndexAvailableShards") router.HandleFunc("/internal/nodes", handler.handleGetNodes).Methods("GET").Name("GetNodes") @@ -446,6 +450,7 @@ func newRouter(handler *Handler) http.Handler { router.HandleFunc("/internal/idalloc/reserve", handler.handleReserveIDs).Methods("POST").Name("ReserveIDs") router.HandleFunc("/internal/idalloc/commit", handler.handleCommitIDs).Methods("POST").Name("CommitIDs") router.HandleFunc("/internal/idalloc/reset/{index}", handler.handleResetIDAlloc).Methods("POST").Name("ResetIDAlloc") + router.HandleFunc("/internal/idalloc/data", handler.handleIDAllocData).Methods("GET").Name("IDAllocData") // endpoints for collecting cpu profiles from a chosen begin point to // when the client wants to stop. Used for profiling imports that @@ -1217,6 +1222,14 @@ func (h *Handler) handlePostIndexAttrDiff(w http.ResponseWriter, r *http.Request } } +// handleGetIndexAttrData handles GET /internal/index/{index}/attr/data requests. +func (h *Handler) handleGetIndexAttrData(w http.ResponseWriter, r *http.Request) { + if err := h.api.WriteColumnAttrDataTo(r.Context(), w, mux.Vars(r)["index"]); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } +} + func (h *Handler) handleGetActiveQueries(w http.ResponseWriter, r *http.Request) { var rtype string switch { @@ -1741,6 +1754,14 @@ type postFieldAttrDiffResponse struct { Attrs map[uint64]map[string]interface{} `json:"attrs"` } +// handleGetFieldAttrData handles GET /internal/index/{index}/field/{field}/attr/data requests. +func (h *Handler) handleGetFieldAttrData(w http.ResponseWriter, r *http.Request) { + if err := h.api.WriteRowAttrDataTo(r.Context(), w, mux.Vars(r)["index"], mux.Vars(r)["field"]); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } +} + // handleGetIndexShardSnapshot handles GET /internal/index/{index}/shard/{shard}/snapshot requests. func (h *Handler) handleGetIndexShardSnapshot(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] @@ -3071,3 +3092,11 @@ func (h *Handler) handleResetIDAlloc(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) w.Write([]byte("OK")) //nolint:errcheck } + +func (h *Handler) handleIDAllocData(w http.ResponseWriter, r *http.Request) { + w.Header().Add("Content-Type", "application/octet-stream") + if err := h.api.WriteIDAllocDataTo(w); err != nil { + http.Error(w, fmt.Sprintf("writeing id allocation data: %v", err.Error()), http.StatusInternalServerError) + return + } +} diff --git a/idalloc.go b/idalloc.go index 97657f6e4..d73a1a103 100644 --- a/idalloc.go +++ b/idalloc.go @@ -17,6 +17,7 @@ package pilosa import ( "encoding/binary" "fmt" + "io" "math/bits" "sort" "time" @@ -72,6 +73,20 @@ func (ida *idAllocator) Close() error { return ida.db.Close() } +func (ida *idAllocator) WriteTo(w io.Writer) (int64, error) { + if ida == nil || ida.db == nil { + return 0, fmt.Errorf("idAllocator closed") + } + + tx, err := ida.db.Begin(false) + if err != nil { + return 0, err + } + defer tx.Rollback() + + return tx.WriteTo(w) +} + // ErrIDOffsetDesync is an error generated when attempting to reserve IDs at a committed offset. // This will typically happen when kafka partitions are moved between kafka ingesters - there may be a brief period in which 2 ingesters are processing the same messages at the same time. // The ingester can resolve this by ignoring messages under base.