Add id alloc, col attrs, & row attrs backup

This commit is contained in:
Ben Johnson 2021-05-06 11:17:38 -06:00
parent 963bb3ec59
commit 7284c4dd10
8 changed files with 266 additions and 1 deletions

23
api.go
View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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", ...)

View file

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

View file

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