diff --git a/cmd/restore_tar.go b/cmd/restore_tar.go index d626a488c..c00b9f8da 100644 --- a/cmd/restore_tar.go +++ b/cmd/restore_tar.go @@ -23,6 +23,8 @@ The Restore command will take a tar-formatted backup archive and restore it to a flags.DurationVar(&cmd.RetryPeriod, "retry-period", cmd.RetryPeriod, "Length of time after HTTP request failure to continue retrying request.") flags.StringVar(&cmd.Pprof, "pprof", cmd.Pprof, "host:port to listen for profiling requests at /debug/pprof and /debug/fgprof.") flags.StringVar(&cmd.AuthToken, "auth-token", "", "Authentication token") + flags.StringVar(&cmd.TempDir, "temp-dir", cmd.TempDir, "Location of tempory spillover files the default is the system's default(usually /tmp)") + ctl.SetTLSConfig( flags, "", &cmd.TLS.CertificatePath, diff --git a/ctl/backup_tar.go b/ctl/backup_tar.go index f2244ff17..6dcf95c7e 100644 --- a/ctl/backup_tar.go +++ b/ctl/backup_tar.go @@ -3,6 +3,7 @@ package ctl import ( "archive/tar" + "bufio" "bytes" "context" "crypto/tls" @@ -492,3 +493,10 @@ func (fb *FileBuffer) Reset() error { fb.buf.Reset() return fb.Close() } + +func (fb *FileBuffer) NewReader() io.Reader { + if fb.file != nil { + return bytes.NewReader(fb.buf.Bytes()) + } + return bufio.NewReader(fb.file) +} diff --git a/ctl/restore_tar.go b/ctl/restore_tar.go index 7432a3b75..bcfee44bc 100644 --- a/ctl/restore_tar.go +++ b/ctl/restore_tar.go @@ -3,7 +3,6 @@ package ctl import ( "archive/tar" - "bytes" "compress/gzip" "context" "crypto/tls" @@ -39,6 +38,8 @@ type RestoreTarCommand struct { // Host:port on which to listen for pprof. Pprof string `json:"pprof"` + TempDir string + // Reusable client. client *pilosa.InternalClient @@ -134,7 +135,10 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) { return errors.New("no primary") } c := &gohttp.Client{} - buf := new(bytes.Buffer) + mb512 := 2 << 29 + buf := NewFileBuffer(mb512, cmd.TempDir) + defer buf.Close() + for { buf.Reset() header, err := tarReader.Next() @@ -188,7 +192,8 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) { node := node g.Go(func() error { client := &gohttp.Client{} - rd := bytes.NewReader(buf.Bytes()) + // rd := bytes.NewReader(buf.Bytes()) + rd := buf.NewReader() logger.Printf("shard %v %v", shard, indexName) url := node.URI.Path(fmt.Sprintf("/internal/restore/%v/%v", indexName, shard)) _, err = client.Post(url, "application/octet-stream", rd) @@ -220,7 +225,8 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) { node := node g.Go(func() error { client := &gohttp.Client{} - rd := bytes.NewReader(buf.Bytes()) + // rd := bytes.NewReader(buf.Bytes()) + rd := buf.NewReader() logger.Printf("dataframe shard %v %v", shard, indexName) url := node.URI.Path(fmt.Sprintf("/internal/dataframe/restore/%v/%v", indexName, shard)) _, err = client.Post(url, "application/octet-stream", rd) @@ -252,7 +258,8 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) { g.Go(func() error { // rd := bytes.NewReader(shardBytes) rd := func() (io.Reader, error) { - return bytes.NewReader(buf.Bytes()), nil + return buf.NewReader(), nil + // return bytes.NewReader(buf.Bytes()), nil } return cmd.client.ImportIndexKeys(ctx, &node.URI, indexName, partitionID, false, rd) @@ -269,20 +276,18 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) { switch action := record[4]; action { case "translate": logger.Printf("field keys %v %v", indexName, fieldName) - bc := NewBroadcaster(tarReader, len(nodes)) + _, err = io.Copy(buf, tarReader) g, _ := errgroup.WithContext(ctx) - for i, node := range nodes { - i := i + for _, node := range nodes { node := node g.Go(func() error { rd := func() (io.Reader, error) { - return bc.Readers[i], nil + return buf.NewReader(), nil } return cmd.client.ImportFieldKeys(ctx, &node.URI, indexName, fieldName, false, rd) }) } - bc.Consume() if err := g.Wait(); err != nil { return err }